返回 CodeWhale
lib.rs
根目录 / crates / hooks / src / lib.rs
1 use std::path::{Path, PathBuf};
2 use std::sync::Arc;
3 #[cfg(unix)]
4 use std::time::Duration;
5
6 use anyhow::{Context, Result};
7 use async_trait::async_trait;
8 use chrono::Utc;
9 use codewhale_protocol::EventFrame;
10 use serde::{Deserialize, Serialize};
11 use serde_json::{Value, json};
12 use tokio::io::AsyncWriteExt;
13
14 mod lifecycle_outbox;
15
16 pub use lifecycle_outbox::{
17 LifecycleEvent, LifecycleOutbox, OUTBOX_DETAIL_MAX_CHARS, OUTBOX_HEADLINE_MAX_CHARS,
18 OUTBOX_PREVIEW_MAX_CHARS, OUTBOX_TRUNCATION_MARKER, bounded_text,
19 };
20
21 /// All events that can be emitted through the hook system.
22 ///
23 /// Each variant represents a distinct lifecycle or streaming event. The enum is
24 /// serialised with a `"type"` discriminator using `snake_case` naming (e.g.
25 /// `"response_start"`, `"tool_lifecycle"`), making it easy to consume from
26 /// JSON-based log files or webhook receivers.
27 #[allow(clippy::large_enum_variant)] // Keep the public HookEvent shape stable for 0.8.x.
28 #[derive(Debug, Clone, Serialize, Deserialize)]
29 #[serde(tag = "type", rename_all = "snake_case")]
30 pub enum HookEvent {
31 /// A new response stream has started.
32 ResponseStart {
33 /// Unique identifier for the response being streamed.
34 response_id: String,
35 },
36 /// A chunk of text has been received for an in-progress response.
37 ResponseDelta {
38 /// Unique identifier for the response being streamed.
39 response_id: String,
40 /// The incremental text content of this chunk.
41 delta: String,
42 },
43 /// A response stream has finished.
44 ResponseEnd {
45 /// Unique identifier for the response that completed.
46 response_id: String,
47 },
48 /// A tool invocation has transitioned to a new phase (e.g. start, end, error).
49 ToolLifecycle {
50 /// Identifier of the response under which the tool was invoked.
51 response_id: String,
52 /// Name of the tool (e.g. `"shell"`, `"read_file"`).
53 tool_name: String,
54 /// Current phase of the tool execution (e.g. `"start"`, `"end"`).
55 phase: String,
56 /// Arbitrary structured payload associated with this phase.
57 payload: Value,
58 },
59 /// A background job has transitioned to a new phase.
60 JobLifecycle {
61 /// Unique identifier of the job.
62 job_id: String,
63 /// Current phase of the job (e.g. `"queued"`, `"running"`, `"done"`).
64 phase: String,
65 /// Optional progress percentage (0-100).
66 progress: Option<u8>,
67 /// Optional human-readable detail about the current phase.
68 detail: Option<String>,
69 },
70 /// An approval request has transitioned to a new phase.
71 ApprovalLifecycle {
72 /// Unique identifier of the approval request.
73 approval_id: String,
74 /// Current phase (e.g. `"requested"`, `"approved"`, `"denied"`).
75 phase: String,
76 /// Optional reason explaining the current phase.
77 reason: Option<String>,
78 },
79 /// A catch-all variant that wraps an arbitrary [`EventFrame`].
80 ///
81 /// Use this when you need to forward a protocol-level event frame without
82 /// mapping it to a more specific variant.
83 GenericEventFrame {
84 /// The raw event frame to forward.
85 frame: Box<EventFrame>,
86 },
87 }
88
89 impl HookEvent {
90 /// Serialise this event into a [`serde_json::Value`].
91 ///
92 /// Returns a JSON object with the `"type"` discriminator and all variant
93 /// fields. If serialisation fails (which should be extremely rare), a
94 /// fallback `{"type":"serialization_error"}` value is returned instead of
95 /// panicking.
96 pub fn to_json(&self) -> Value {
97 serde_json::to_value(self).unwrap_or_else(|_| json!({"type":"serialization_error"}))
98 }
99 }
100
101 /// A destination that can receive [`HookEvent`]s.
102 ///
103 /// Implementors handle the transport-specific details of delivering events
104 /// (writing to stdout, appending to a file, POSTing to a webhook, etc.).
105 /// The [`HookDispatcher`] fans out every event to all registered sinks, so a
106 /// single process can log to multiple destinations simultaneously.
107 ///
108 /// Sinks are expected to be **best-effort**: implementations should avoid
109 /// panicking and should return an [`anyhow::Error`] only for truly unexpected
110 /// failures. [`HookDispatcher::emit`] discards individual sink errors so hook
111 /// delivery failures do not abort the application.
112 #[async_trait]
113 pub trait HookSink: Send + Sync {
114 /// Deliver a single event to this sink.
115 ///
116 /// Implementations should be resilient to transient failures (e.g. a
117 /// missing listener) and should not block the caller for extended periods.
118 async fn emit(&self, event: &HookEvent) -> Result<()>;
119 }
120
121 /// A [`HookSink`] that prints each event as a single JSON line to stdout.
122 ///
123 /// Useful for local development and debugging. Events are printed via
124 /// [`println!`] so they appear interleaved with other program output.
125 #[derive(Default)]
126 pub struct StdoutHookSink;
127
128 #[async_trait]
129 impl HookSink for StdoutHookSink {
130 async fn emit(&self, event: &HookEvent) -> Result<()> {
131 println!("{}", event.to_json());
132 Ok(())
133 }
134 }
135
136 /// A [`HookSink`] that appends each event as a JSON line to a file.
137 ///
138 /// The file is created (along with any missing parent directories) on the
139 /// first emitted event. Each line is a JSON object of the form
140 /// `{"at": "<ISO 8601 timestamp>", "event": {...}}`.
141 ///
142 /// Concurrent [`emit`](HookSink::emit) calls serialize on an internal mutex so
143 /// that each JSON event is written and flushed as a complete line; without that
144 /// lock, overlapping `write_all` calls can interleave partial lines and corrupt
145 /// the JSONL log (see issue #4739).
146 pub struct JsonlHookSink {
147 path: PathBuf,
148 /// Serializes open+append+flush so concurrent tool-call emits cannot
149 /// interleave bytes mid-line.
150 write_lock: tokio::sync::Mutex<()>,
151 }
152
153 impl JsonlHookSink {
154 /// Create a new sink that writes to the file at `path`.
155 ///
156 /// A leading `~` is expanded against the user home. Parent directories
157 /// are created lazily on the first [`HookSink::emit`] call.
158 pub fn new(path: PathBuf) -> Self {
159 Self {
160 path: expand_home(path),
161 write_lock: tokio::sync::Mutex::new(()),
162 }
163 }
164 }
165
166 /// Expand a leading `~` in a configured sink path, reusing the one expansion
167 /// rule global path overrides use. Documented examples write
168 /// `~/.codewhale/...`; taken literally that created a directory named `~` in
169 /// the working directory. A path without `~` (absolute or relative) is
170 /// returned unchanged, as is `~` when no home resolves.
171 pub(crate) fn expand_home(path: PathBuf) -> PathBuf {
172 codewhale_paths::validate_absolute_path("hook sink path", path.clone()).unwrap_or(path)
173 }
174
175 /// Open a JSONL sink for append, creating it privately: on Unix new parent
176 /// directories are 0700 and a new file 0600 (existing modes are left alone).
177 /// Hook and lifecycle events carry tool payloads, prompts, and paths; under
178 /// an ordinary umask the old defaults were readable by every local user.
179 pub(crate) async fn open_private_append(path: &Path) -> Result<tokio::fs::File> {
180 if let Some(parent) = path
181 .parent()
182 .filter(|parent| !parent.as_os_str().is_empty())
183 {
184 let mut dirs = tokio::fs::DirBuilder::new();
185 dirs.recursive(true);
186 #[cfg(unix)]
187 dirs.mode(0o700);
188 dirs.create(parent)
189 .await
190 .with_context(|| format!("failed to create sink directory {}", parent.display()))?;
191 }
192 let mut options = tokio::fs::OpenOptions::new();
193 options.create(true).append(true);
194 #[cfg(unix)]
195 options.mode(0o600);
196 options
197 .open(path)
198 .await
199 .with_context(|| format!("failed to open sink file {}", path.display()))
200 }
201
202 #[async_trait]
203 impl HookSink for JsonlHookSink {
204 async fn emit(&self, event: &HookEvent) -> Result<()> {
205 // Encode outside the lock so only I/O is serialized.
206 let payload = json!({
207 "at": Utc::now().to_rfc3339(),
208 "event": event
209 });
210 let encoded = serde_json::to_string(&payload).context("failed to encode hook event")?;
211
212 let _guard = self.write_lock.lock().await;
213 let mut file = open_private_append(&self.path).await?;
214 file.write_all(encoded.as_bytes())
215 .await
216 .context("failed to write hook event")?;
217 file.write_all(b"\n")
218 .await
219 .context("failed to write hook event newline")?;
220 // Flush before drop so sequential emits (and tests that read the
221 // file immediately after) observe every completed line. Holding the
222 // mutex through flush guarantees concurrent writers never observe a
223 // partial line at EOF.
224 file.flush().await.context("failed to flush hook event")?;
225 Ok(())
226 }
227 }
228
229 /// A [`HookSink`] that POSTs each event as JSON to a remote HTTP endpoint.
230 ///
231 /// The request body is `{"at": "<ISO 8601 timestamp>", "event": {...}}`.
232 /// Failed requests are retried up to 2 times with exponential back-off
233 /// (200 ms, 400 ms). After exhausting retries the error is propagated.
234 #[derive(Clone)]
235 pub struct WebhookHookSink {
236 url: String,
237 /// Optional bearer token sent as `Authorization: Bearer <token>`.
238 bearer_token: Option<String>,
239 client: reqwest::Client,
240 }
241
242 impl WebhookHookSink {
243 /// Create a new sink that sends events to the given `url`.
244 pub fn new(url: String) -> Self {
245 Self::new_with_token(url, None)
246 }
247
248 /// Create a new sink that sends events to the given `url`, attaching
249 /// `Authorization: Bearer <token>` when a token is provided.
250 pub fn new_with_token(url: String, bearer_token: Option<String>) -> Self {
251 Self {
252 url,
253 bearer_token,
254 client: codewhale_release::platform_http_client_builder()
255 .timeout(std::time::Duration::from_secs(10))
256 .build()
257 .unwrap_or_else(|_| {
258 codewhale_release::platform_http_client_builder()
259 .build()
260 .unwrap_or_else(|_| codewhale_release::tls::reqwest_client())
261 }),
262 }
263 }
264
265 /// POST an arbitrary JSON payload to the configured endpoint.
266 ///
267 /// This is the shared delivery path behind both [`HookSink::emit`] and
268 /// the lifecycle outbox fan-out. It is deliberately not part of the
269 /// [`HookSink`] trait: outbox events are runtime event envelopes, not
270 /// [`HookEvent`]s, and only the transport needs to be shared.
271 pub async fn post_payload(&self, payload: serde_json::Value) -> Result<()> {
272 let mut retries = 0usize;
273 loop {
274 let mut request = self.client.post(&self.url).json(&payload);
275 if let Some(token) = self.bearer_token.as_deref().filter(|t| !t.is_empty()) {
276 request = request.bearer_auth(token);
277 }
278 let resp = request.send().await;
279 match resp {
280 Ok(response) if response.status().is_success() => return Ok(()),
281 Ok(response) => {
282 if retries >= 2 {
283 anyhow::bail!("webhook returned non-success status {}", response.status());
284 }
285 }
286 Err(err) => {
287 if retries >= 2 {
288 return Err(err).context("webhook request failed");
289 }
290 }
291 }
292 retries += 1;
293 tokio::time::sleep(std::time::Duration::from_millis(200 * retries as u64)).await;
294 }
295 }
296 }
297
298 #[async_trait]
299 impl HookSink for WebhookHookSink {
300 async fn emit(&self, event: &HookEvent) -> Result<()> {
301 self.post_payload(json!({
302 "at": Utc::now().to_rfc3339(),
303 "event": event,
304 }))
305 .await
306 }
307 }
308
309 /// A [`HookSink`] that sends events over a Unix domain socket.
310 ///
311 /// Each event is serialized as a single JSON line (`{"at": "...", "event": {...}}\n`)
312 /// and written to the socket. If the socket is not available (listener not running),
313 /// the event is silently dropped - hook sinks are best-effort observability, not
314 /// control flow.
315 ///
316 /// On non-Unix platforms this struct exists but its [`HookSink::emit`] is a no-op.
317 #[derive(Debug, Clone)]
318 pub struct UnixSocketHookSink {
319 #[cfg(unix)]
320 path: PathBuf,
321 }
322
323 impl UnixSocketHookSink {
324 /// Create a sink that connects to the Unix domain socket at `path`.
325 pub fn new(path: PathBuf) -> Self {
326 #[cfg(unix)]
327 {
328 Self { path }
329 }
330 #[cfg(not(unix))]
331 {
332 let _ = path;
333 Self {}
334 }
335 }
336 }
337
338 /// Longest one Unix-socket delivery (connect + write) may take. The
339 /// dispatcher awaits its sinks in order, so a listener that accepts and never
340 /// reads would otherwise stall every later sink, and the emitting turn, as
341 /// soon as the socket buffer filled.
342 #[cfg(unix)]
343 const UNIX_SOCKET_SINK_TIMEOUT: Duration = Duration::from_secs(2);
344
345 #[async_trait]
346 impl HookSink for UnixSocketHookSink {
347 #[cfg(unix)]
348 async fn emit(&self, event: &HookEvent) -> Result<()> {
349 let payload = json!({
350 "at": Utc::now().to_rfc3339(),
351 "event": event
352 });
353 let mut line = serde_json::to_string(&payload).context("failed to encode hook event")?;
354 line.push('\n');
355 let deliver = async {
356 let mut stream = match tokio::net::UnixStream::connect(&self.path).await {
357 Ok(s) => s,
358 Err(_) => return Ok(()), // listener not running, skip silently
359 };
360 stream
361 .write_all(line.as_bytes())
362 .await
363 .context("failed to write to unix socket")
364 };
365 tokio::time::timeout(UNIX_SOCKET_SINK_TIMEOUT, deliver)
366 .await
367 .map_err(|_| {
368 anyhow::anyhow!(
369 "unix socket hook sink {} did not accept the event within {:?}",
370 self.path.display(),
371 UNIX_SOCKET_SINK_TIMEOUT
372 )
373 })?
374 }
375
376 #[cfg(not(unix))]
377 async fn emit(&self, _event: &HookEvent) -> Result<()> {
378 // Unix sockets are not available on this platform.
379 Ok(())
380 }
381 }
382
383 /// Fans out [`HookEvent`]s to a collection of [`HookSink`]s.
384 ///
385 /// Register one or more sinks via [`add_sink`](HookDispatcher::add_sink),
386 /// then call [`emit`](HookDispatcher::emit) to broadcast an event to all of
387 /// them. If a sink returns an error it is silently ignored so that a failing
388 /// sink does not prevent remaining sinks from receiving the event.
389 #[derive(Default, Clone)]
390 pub struct HookDispatcher {
391 sinks: Vec<Arc<dyn HookSink>>,
392 }
393
394 impl HookDispatcher {
395 /// Register a new sink that will receive all subsequently emitted events.
396 pub fn add_sink(&mut self, sink: Arc<dyn HookSink>) {
397 self.sinks.push(sink);
398 }
399
400 /// Number of registered sinks. Exposed so transport setup can assert
401 /// exactly which sinks were wired (e.g. no stdout sink in stdio mode).
402 #[must_use]
403 pub fn sink_count(&self) -> usize {
404 self.sinks.len()
405 }
406
407 /// Broadcast an event to every registered sink.
408 ///
409 /// Errors from individual sinks are silently discarded so that one failing
410 /// sink does not block the others.
411 pub async fn emit(&self, event: HookEvent) {
412 for sink in &self.sinks {
413 let _ = sink.emit(&event).await;
414 }
415 }
416 }
417
418 #[cfg(test)]
419 mod tests {
420 use super::*;
421 use std::sync::Mutex;
422 #[cfg(unix)]
423 use std::sync::atomic::{AtomicU64, Ordering};
424 #[cfg(unix)]
425 use std::time::Duration;
426 use std::time::{SystemTime, UNIX_EPOCH};
427
428 #[cfg(unix)]
429 static SOCKET_PATH_NONCE: AtomicU64 = AtomicU64::new(0);
430
431 #[test]
432 fn hook_event_serializes_with_snake_case_type_and_payload() {
433 let event = HookEvent::ToolLifecycle {
434 response_id: "resp-1".to_string(),
435 tool_name: "shell".to_string(),
436 phase: "end".to_string(),
437 payload: json!({ "exit_code": 0 }),
438 };
439
440 let encoded = event.to_json();
441
442 assert_eq!(encoded["type"], "tool_lifecycle");
443 assert_eq!(encoded["response_id"], "resp-1");
444 assert_eq!(encoded["tool_name"], "shell");
445 assert_eq!(encoded["phase"], "end");
446 assert_eq!(encoded["payload"]["exit_code"], 0);
447 }
448
449 #[test]
450 fn generic_event_frame_serialization_is_unchanged_by_boxing() {
451 let event = HookEvent::GenericEventFrame {
452 frame: Box::new(EventFrame::ResponseStart {
453 response_id: "resp-1".to_string(),
454 }),
455 };
456
457 let encoded = event.to_json();
458
459 assert_eq!(encoded["type"], "generic_event_frame");
460 assert_eq!(encoded["frame"]["event"], "response_start");
461 assert_eq!(encoded["frame"]["response_id"], "resp-1");
462 }
463
464 #[tokio::test]
465 async fn jsonl_sink_creates_parent_dir_and_appends_events() {
466 let root = unique_temp_dir("jsonl_sink");
467 let path = root.join("nested").join("hooks.jsonl");
468 let sink = JsonlHookSink::new(path.clone());
469
470 sink.emit(&HookEvent::ResponseStart {
471 response_id: "resp-1".to_string(),
472 })
473 .await
474 .unwrap();
475 sink.emit(&HookEvent::ResponseEnd {
476 response_id: "resp-1".to_string(),
477 })
478 .await
479 .unwrap();
480
481 let raw = std::fs::read_to_string(&path).unwrap();
482 let lines = raw.lines().collect::<Vec<_>>();
483 assert_eq!(lines.len(), 2);
484
485 let first: Value = serde_json::from_str(lines[0]).unwrap();
486 let second: Value = serde_json::from_str(lines[1]).unwrap();
487 assert!(first["at"].as_str().is_some());
488 assert_eq!(first["event"]["type"], "response_start");
489 assert_eq!(first["event"]["response_id"], "resp-1");
490 assert_eq!(second["event"]["type"], "response_end");
491 assert_eq!(second["event"]["response_id"], "resp-1");
492
493 let _ = std::fs::remove_dir_all(root);
494 }
495
496 /// Concurrent emits must not interleave partial lines (#4739).
497 ///
498 /// Spawns many tasks writing through one shared [`JsonlHookSink`] and
499 /// asserts every line in the resulting file is complete, parseable JSON.
500 #[tokio::test(flavor = "multi_thread", worker_threads = 8)]
501 async fn jsonl_sink_concurrent_emits_write_atomic_json_lines() {
502 let root = unique_temp_dir("jsonl_sink_concurrent");
503 let path = root.join("hooks.jsonl");
504 let sink = Arc::new(JsonlHookSink::new(path.clone()));
505
506 const TASKS: usize = 32;
507 const EVENTS_PER_TASK: usize = 20;
508 let expected = TASKS * EVENTS_PER_TASK;
509
510 let mut handles = Vec::with_capacity(TASKS);
511 for task_id in 0..TASKS {
512 let sink = Arc::clone(&sink);
513 handles.push(tokio::spawn(async move {
514 for n in 0..EVENTS_PER_TASK {
515 let event = HookEvent::ToolLifecycle {
516 response_id: format!("resp-{task_id}"),
517 tool_name: "shell".to_string(),
518 phase: "end".to_string(),
519 payload: json!({ "n": n, "task": task_id }),
520 };
521 sink.emit(&event).await.expect("concurrent emit");
522 }
523 }));
524 }
525 for handle in handles {
526 handle.await.expect("join concurrent writer");
527 }
528
529 let raw = std::fs::read_to_string(&path).expect("read concurrent jsonl");
530 // Trailing newline yields an empty final split; keep non-empty lines only.
531 let lines: Vec<&str> = raw.lines().filter(|l| !l.is_empty()).collect();
532 assert_eq!(
533 lines.len(),
534 expected,
535 "expected {expected} complete lines, got {}; sample: {:?}",
536 lines.len(),
537 lines
538 .first()
539 .map(|s| s.chars().take(80).collect::<String>())
540 );
541
542 for (idx, line) in lines.iter().enumerate() {
543 let parsed: Value = serde_json::from_str(line).unwrap_or_else(|err| {
544 panic!("line {idx} is not complete JSON ({err}): {line:?}");
545 });
546 assert!(
547 parsed.get("at").and_then(|v| v.as_str()).is_some(),
548 "line {idx} missing at: {parsed}"
549 );
550 assert_eq!(
551 parsed["event"]["type"], "tool_lifecycle",
552 "line {idx} unexpected event type"
553 );
554 }
555
556 let _ = std::fs::remove_dir_all(root);
557 }
558
559 #[tokio::test]
560 async fn dispatcher_continues_after_sink_error() {
561 let mut dispatcher = HookDispatcher::default();
562 let first = Arc::new(RecordingSink::default());
563 let second = Arc::new(RecordingSink::default());
564
565 dispatcher.add_sink(first.clone());
566 dispatcher.add_sink(Arc::new(FailingSink));
567 dispatcher.add_sink(second.clone());
568
569 dispatcher
570 .emit(HookEvent::ApprovalLifecycle {
571 approval_id: "approval-1".to_string(),
572 phase: "requested".to_string(),
573 reason: Some("needs review".to_string()),
574 })
575 .await;
576
577 assert_eq!(
578 first.events(),
579 vec![json!({
580 "type": "approval_lifecycle",
581 "approval_id": "approval-1",
582 "phase": "requested",
583 "reason": "needs review",
584 })]
585 );
586 assert_eq!(second.events(), first.events());
587 }
588
589 #[cfg(unix)]
590 #[tokio::test]
591 async fn unix_socket_sink_skips_when_listener_absent() {
592 let (_root, socket_path) = unique_short_socket_path("missing");
593 let sink = UnixSocketHookSink::new(socket_path);
594 let result = sink
595 .emit(&HookEvent::ResponseStart {
596 response_id: "resp-1".to_string(),
597 })
598 .await;
599 assert!(result.is_ok());
600 }
601
602 #[cfg(unix)]
603 #[tokio::test]
604 async fn unix_socket_sink_sends_event_to_listener() {
605 use tokio::io::AsyncBufReadExt;
606 use tokio::net::UnixListener;
607
608 let (root, socket_path) = unique_short_socket_path("send");
609 std::fs::create_dir_all(&root).expect("mkdir");
610 let _ = std::fs::remove_file(&socket_path);
611
612 let listener = UnixListener::bind(&socket_path).expect("bind");
613 let sink = UnixSocketHookSink::new(socket_path.clone());
614 let cleanup = || {
615 let _ = std::fs::remove_file(&socket_path);
616 let _ = std::fs::remove_dir_all(&root);
617 };
618
619 let mut handle = tokio::spawn(async move {
620 let (stream, _) = listener.accept().await.expect("accept");
621 let mut reader = tokio::io::BufReader::new(stream);
622 let mut line = String::new();
623 reader.read_line(&mut line).await.expect("read_line");
624 line
625 });
626
627 let event = HookEvent::ResponseStart {
628 response_id: "resp-42".to_string(),
629 };
630 let emit = sink.emit(&event);
631 match tokio::time::timeout(Duration::from_secs(5), emit).await {
632 Ok(result) => result.expect("emit"),
633 Err(_) => {
634 handle.abort();
635 let _ = (&mut handle).await;
636 cleanup();
637 panic!("unix socket emit timed out");
638 }
639 }
640
641 let received = match tokio::time::timeout(Duration::from_secs(5), &mut handle).await {
642 Ok(result) => result.expect("join"),
643 Err(_) => {
644 handle.abort();
645 let _ = (&mut handle).await;
646 cleanup();
647 panic!("unix socket exchange timed out");
648 }
649 };
650 let parsed: Value = serde_json::from_str(&received).expect("parse");
651 assert_eq!(parsed["event"]["type"], "response_start");
652 assert_eq!(parsed["event"]["response_id"], "resp-42");
653 assert!(parsed["at"].as_str().is_some());
654
655 cleanup();
656 }
657
658 /// Audit R04-11: a listener that accepts and never reads must not stall
659 /// the sequential dispatcher past the sink's own bound.
660 #[cfg(unix)]
661 #[tokio::test]
662 async fn unix_socket_sink_gives_up_on_a_listener_that_never_reads() {
663 use tokio::net::UnixListener;
664
665 let (root, socket_path) = unique_short_socket_path("stall");
666 std::fs::create_dir_all(&root).expect("mkdir");
667 let listener = UnixListener::bind(&socket_path).expect("bind");
668 let held = tokio::spawn(async move {
669 let (stream, _) = listener.accept().await.expect("accept");
670 // Hold the connection open without ever reading from it.
671 tokio::time::sleep(Duration::from_secs(60)).await;
672 drop(stream);
673 });
674 let sink = UnixSocketHookSink::new(socket_path.clone());
675 // Far larger than any socket buffer, so `write_all` must block.
676 let event = HookEvent::ToolLifecycle {
677 response_id: "resp-stall".to_string(),
678 tool_name: "shell".to_string(),
679 phase: "end".to_string(),
680 payload: json!({ "output": "x".repeat(16 * 1024 * 1024) }),
681 };
682
683 let outcome = tokio::time::timeout(Duration::from_secs(10), sink.emit(&event)).await;
684 held.abort();
685 let _ = std::fs::remove_file(&socket_path);
686 let _ = std::fs::remove_dir_all(&root);
687 let result = outcome.expect("emit must return within its own bound");
688 assert!(result.is_err(), "a stalled delivery is reported, not faked");
689 }
690
691 /// Audit R04-10: `~` expands against the home, and a new JSONL sink is
692 /// private (dir 0700, file 0600) on Unix.
693 #[cfg(unix)]
694 #[tokio::test]
695 async fn jsonl_sink_expands_home_and_creates_private_files() {
696 use std::os::unix::fs::PermissionsExt;
697
698 let home = codewhale_paths::user_home().expect("test home");
699 assert_eq!(
700 expand_home(PathBuf::from("~/.codewhale/events.jsonl")),
701 home.join(".codewhale/events.jsonl")
702 );
703 assert_eq!(
704 expand_home(PathBuf::from("relative/events.jsonl")),
705 PathBuf::from("relative/events.jsonl")
706 );
707
708 let root = unique_temp_dir("private");
709 let path = root.join("nested").join("events.jsonl");
710 let sink = JsonlHookSink::new(path.clone());
711 sink.emit(&HookEvent::ResponseStart {
712 response_id: "resp-private".to_string(),
713 })
714 .await
715 .expect("emit");
716 let file_mode = std::fs::metadata(&path).unwrap().permissions().mode() & 0o777;
717 let dir_mode = std::fs::metadata(path.parent().unwrap())
718 .unwrap()
719 .permissions()
720 .mode()
721 & 0o777;
722 let _ = std::fs::remove_dir_all(&root);
723 assert_eq!(file_mode, 0o600);
724 assert_eq!(dir_mode, 0o700);
725 }
726
727 #[derive(Default)]
728 struct RecordingSink {
729 events: Mutex<Vec<Value>>,
730 }
731
732 impl RecordingSink {
733 fn events(&self) -> Vec<Value> {
734 self.events.lock().unwrap().clone()
735 }
736 }
737
738 #[async_trait::async_trait]
739 impl HookSink for RecordingSink {
740 async fn emit(&self, event: &HookEvent) -> Result<()> {
741 self.events.lock().unwrap().push(event.to_json());
742 Ok(())
743 }
744 }
745
746 struct FailingSink;
747
748 #[async_trait::async_trait]
749 impl HookSink for FailingSink {
750 async fn emit(&self, _event: &HookEvent) -> Result<()> {
751 anyhow::bail!("sink failed")
752 }
753 }
754
755 fn unique_temp_dir(label: &str) -> PathBuf {
756 let nanos = SystemTime::now()
757 .duration_since(UNIX_EPOCH)
758 .unwrap()
759 .as_nanos();
760 std::env::temp_dir().join(format!(
761 "deepseek-hooks-{label}-{}-{nanos}",
762 std::process::id()
763 ))
764 }
765
766 #[cfg(unix)]
767 fn unique_short_socket_path(label: &str) -> (PathBuf, PathBuf) {
768 let nanos = SystemTime::now()
769 .duration_since(UNIX_EPOCH)
770 .unwrap()
771 .as_nanos();
772 let nonce = SOCKET_PATH_NONCE.fetch_add(1, Ordering::Relaxed);
773 let root = PathBuf::from("/tmp").join(format!(
774 "cw-hk-{label}-{}-{nanos}-{nonce}",
775 std::process::id()
776 ));
777 let path = root.join("hook.sock");
778 (root, path)
779 }
780
781 #[test]
782 fn webhook_sink_new_does_not_panic() {
783 // Construction must never panic: if the configured client builder
784 // fails, the code falls back to a default `reqwest::Client` instead of
785 // calling `.expect(...)`.
786 let sink = WebhookHookSink::new("https://example.invalid/webhook".to_string());
787 let _ = sink;
788 }
789 }
790
790 lines RUST