返回 CodeWhale
lifecycle_outbox.rs
根目录 / crates / hooks / src / lifecycle_outbox.rs
1 //! Lifecycle event outbox: a local JSONL log of session/turn/subagent
2 //! lifecycle events plus an optional webhook fan-out.
3 //!
4 //! This is the machine-readable sibling of the TUI shell-hook system. Hooks
5 //! fire shell commands per event and are TUI-only; the outbox appends one
6 //! JSON line per event to a config-gated file and needs no per-event
7 //! configuration. It is additive and opt-in: with no path configured,
8 //! [`LifecycleOutbox::emit`] is a no-op.
9 //!
10 //! # Line schema
11 //!
12 //! Every line is a `codewhale_protocol::runtime::RuntimeEventEnvelope`:
13 //!
14 //! ```json
15 //! {"schema_version": 1, "seq": 3, "event": "turn_start", "kind": "turn.started",
16 //! "thread_id": "…", "turn_id": "…", "item_id": null, "timestamp": "…",
17 //! "created_at": "…", "payload": {…}}
18 //! ```
19 //!
20 //! - `seq` is monotonic per outbox file. On the first write the writer
21 //! recovers the `seq` of the file's last complete line (bounded tail scan,
22 // so an outbox that grows unbounded is never re-read in full) and continues
23 //! from `last + 1`.
24 //! - `event` is the snake-case lifecycle name (`turn_start`, `turn_end`, …);
25 //! `kind` is the dotted kind (`turn.started`, `turn.failed`, …).
26 //! - Payloads are constructed by the emit sites from bounded, pre-redacted
27 //! fields only — never raw tool arguments, environment, or full transcript
28 //! text. [`bounded_text`] enforces the same ceilings as the desktop
29 //! notification payloads: headline ≤ 80, detail ≤ 120, preview ≤ 200
30 //! characters.
31 //!
32 //! # Delivery model
33 //!
34 //! [`LifecycleOutbox::emit`] never blocks the caller: it enqueues the event
35 //! on an internal channel and a single writer task appends lines in order.
36 //! If no tokio runtime is available the event is dropped with a warning.
37 //! Webhook POSTs (`{"at": …, "event": …}`) are attempted after the local
38 //! append; failures are logged and dropped, never retried into the agent
39 //! loop. Process owners call [`LifecycleOutbox::flush`] before runtime teardown
40 //! to give queued terminal events a bounded opportunity to finish delivery.
41
42 use std::path::{Path, PathBuf};
43 use std::sync::atomic::{AtomicBool, Ordering};
44 use std::sync::{Arc, Mutex};
45 use std::time::Duration;
46
47 use anyhow::{Context, Result};
48 use chrono::Utc;
49 use codewhale_protocol::runtime::{RUNTIME_EVENT_ENVELOPE_SCHEMA_VERSION, RuntimeEventEnvelope};
50 use serde_json::{Value, json};
51 use tokio::io::{AsyncReadExt, AsyncSeekExt, AsyncWriteExt};
52 use tokio::sync::mpsc::error::TrySendError;
53 use tokio::sync::mpsc::{Receiver, Sender};
54 use tokio::sync::oneshot;
55
56 use crate::WebhookHookSink;
57
58 /// Text-length ceilings for outbox payload fields. Mirrors the desktop
59 /// notification payload limits so the outbox never carries more than the
60 /// lock-screen-capable surface already does.
61 pub const OUTBOX_HEADLINE_MAX_CHARS: usize = 80;
62 pub const OUTBOX_DETAIL_MAX_CHARS: usize = 120;
63 pub const OUTBOX_PREVIEW_MAX_CHARS: usize = 200;
64
65 /// Suffix appended when [`bounded_text`] truncates a field.
66 pub const OUTBOX_TRUNCATION_MARKER: &str = "…";
67
68 /// How far back from EOF the seq-recovery scan reads. Outbox lines are
69 /// bounded (payload ceilings above plus envelope overhead), so a line can
70 /// never approach this window and the last complete line is always inside it.
71 const SEQ_RECOVERY_TAIL_BYTES: u64 = 64 * 1024;
72
73 /// Queue capacity from emit sites to the writer (#6212). The outbox is
74 /// observability, not control flow: when a consumer is wedged longer than
75 /// this backlog, further events are dropped with a warning instead of
76 /// retaining unbounded snapshots in memory. Flush commands are never dropped
77 /// silently — a full queue fails the flush fast rather than timing it out.
78 const OUTBOX_QUEUE_CAPACITY: usize = 1024;
79
80 /// Events gathered per writer drain before one batched append. Batching only
81 /// affects how many queued events share a single file open/write/flush; each
82 /// line still lands as its own complete JSONL record.
83 const OUTBOX_MAX_BATCH: usize = 64;
84
85 /// Concurrent webhook posts the writer allows in flight. A stalled webhook
86 /// delays only its own event (plus the flush barrier, which the caller's
87 /// timeout bounds) — never the audit-log appends of later events.
88 const WEBHOOK_MAX_INFLIGHT: usize = 8;
89
90 /// One lifecycle event destined for the outbox.
91 ///
92 /// Construct one per emit site. `payload` must only contain bounded,
93 /// pre-redacted fields; apply [`bounded_text`] to anything free-form (error
94 /// messages, previews) before inserting it.
95 #[derive(Debug, Clone)]
96 pub struct LifecycleEvent {
97 /// Snake-case event name, e.g. `"turn_start"`.
98 pub event: String,
99 /// Dotted event kind, e.g. `"turn.started"` or `"turn.failed"`.
100 pub kind: String,
101 /// Owning session/thread id. Empty when the producer has none.
102 pub thread_id: String,
103 /// Current turn id, when known.
104 pub turn_id: Option<String>,
105 /// Current item id, when known.
106 pub item_id: Option<String>,
107 /// Bounded, redacted event payload.
108 pub payload: Value,
109 }
110
111 /// The lifecycle outbox handle.
112 ///
113 /// Cheap to clone (an `Arc`). When constructed without a path the outbox is
114 /// disabled and every `emit` is a no-op.
115 #[derive(Clone)]
116 pub struct LifecycleOutbox {
117 inner: Option<Arc<OutboxInner>>,
118 }
119
120 impl Default for LifecycleOutbox {
121 fn default() -> Self {
122 Self::disabled()
123 }
124 }
125
126 impl LifecycleOutbox {
127 /// Create an outbox writing to `path` when set and non-empty.
128 ///
129 /// `webhook_url` optionally adds a webhook fan-out (POST `{"at", "event"}`,
130 /// best-effort); `webhook_token` is its optional bearer token. Webhook
131 /// delivery is configured independently of the file: it only ever runs
132 /// when `webhook_url` is set, and it never replaces the local append.
133 pub fn new(
134 path: Option<PathBuf>,
135 webhook_url: Option<String>,
136 webhook_token: Option<String>,
137 ) -> Self {
138 let path = match path {
139 Some(path) if !path.as_os_str().is_empty() => crate::expand_home(path),
140 _ => return Self::disabled(),
141 };
142 let webhook = webhook_url
143 .as_deref()
144 .map(str::trim)
145 .filter(|url| !url.is_empty())
146 .map(|url| WebhookHookSink::new_with_token(url.to_string(), webhook_token));
147 let (sender, receiver) = tokio::sync::mpsc::channel(OUTBOX_QUEUE_CAPACITY);
148 Self {
149 inner: Some(Arc::new(OutboxInner {
150 path,
151 webhook,
152 sender,
153 receiver: Mutex::new(Some(receiver)),
154 writer_spawned: AtomicBool::new(false),
155 spawn_lock: Mutex::new(()),
156 })),
157 }
158 }
159
160 /// A disabled outbox that drops every event.
161 pub fn disabled() -> Self {
162 Self { inner: None }
163 }
164
165 /// True when a path was configured and events will be written.
166 pub fn is_enabled(&self) -> bool {
167 self.inner.is_some()
168 }
169
170 /// Emit one lifecycle event.
171 ///
172 /// Never blocks: the event is queued for the outbox's writer task (spawned
173 /// lazily on the current tokio runtime on first use). Events queued with
174 /// no runtime available — or after the writer task is gone — are dropped
175 /// with a warning. Delivery failures inside the writer are logged and
176 /// dropped as well; the outbox is observability, not control flow.
177 pub fn emit(&self, event: LifecycleEvent) {
178 let Some(inner) = self.inner.clone() else {
179 return;
180 };
181 if let Err(error) = inner.enqueue(OutboxCommand::Event(event)) {
182 tracing::warn!(target: "lifecycle_outbox", %error, "lifecycle event dropped");
183 }
184 }
185
186 /// Wait for events queued before this call to finish their delivery attempt.
187 ///
188 /// Local writes are flushed before the writer acknowledges this barrier.
189 /// Delivery errors remain logged and best effort; a stalled writer/webhook
190 /// or a stopped runtime returns an error within `timeout` instead of holding
191 /// process teardown indefinitely. This does not close other cloned handles.
192 pub async fn flush(&self, timeout: Duration) -> Result<()> {
193 let Some(inner) = self.inner.as_ref() else {
194 return Ok(());
195 };
196 let (reply, completed) = oneshot::channel();
197 inner.enqueue(OutboxCommand::Flush(reply))?;
198 tokio::time::timeout(timeout, completed)
199 .await
200 .context("lifecycle outbox flush timed out; queued events may be incomplete")?
201 .context("lifecycle outbox writer stopped before flush completed")
202 }
203 }
204
205 enum OutboxCommand {
206 Event(LifecycleEvent),
207 Flush(oneshot::Sender<()>),
208 }
209
210 struct OutboxInner {
211 path: PathBuf,
212 webhook: Option<WebhookHookSink>,
213 sender: Sender<OutboxCommand>,
214 /// The writer task's receive half. Taken exactly once by the writer task.
215 receiver: Mutex<Option<Receiver<OutboxCommand>>>,
216 writer_spawned: AtomicBool,
217 /// Serializes the lazy writer-task spawn so two racing first emits cannot
218 /// start two writers.
219 spawn_lock: Mutex<()>,
220 }
221
222 impl OutboxInner {
223 /// Queue an event and make sure the writer task exists to drain it.
224 ///
225 /// Ordering: `send` happens before the spawn so events queued before the
226 /// writer starts are drained first, preserving enqueue order. The queue
227 /// is bounded: a wedged consumer drops further events with a warning
228 /// (observability, not control flow) but fails a flush fast instead of
229 /// dropping its reply channel.
230 fn enqueue(self: &Arc<Self>, command: OutboxCommand) -> Result<()> {
231 let is_flush = matches!(command, OutboxCommand::Flush(_));
232 match self.sender.try_send(command) {
233 Ok(()) => {}
234 Err(TrySendError::Full(_)) if !is_flush => {
235 tracing::warn!(
236 target: "lifecycle_outbox",
237 queue_capacity = OUTBOX_QUEUE_CAPACITY,
238 "lifecycle event dropped: outbox queue is full"
239 );
240 }
241 Err(TrySendError::Full(_)) => {
242 anyhow::bail!("lifecycle outbox queue is full; flush rejected");
243 }
244 Err(TrySendError::Closed(_)) => {
245 anyhow::bail!("lifecycle outbox writer task is gone");
246 }
247 }
248 self.ensure_writer_spawned();
249 Ok(())
250 }
251
252 fn ensure_writer_spawned(self: &Arc<Self>) {
253 if self.writer_spawned.load(Ordering::Acquire) {
254 return;
255 }
256 let _guard = self
257 .spawn_lock
258 .lock()
259 .unwrap_or_else(|poisoned| poisoned.into_inner());
260 if self.writer_spawned.load(Ordering::Acquire) {
261 return;
262 }
263 let Ok(handle) = tokio::runtime::Handle::try_current() else {
264 tracing::warn!(
265 target: "lifecycle_outbox",
266 "no tokio runtime available; lifecycle events are queued but will not be written"
267 );
268 return;
269 };
270 let receiver = self
271 .receiver
272 .lock()
273 .unwrap_or_else(|poisoned| poisoned.into_inner())
274 .take();
275 let Some(receiver) = receiver else {
276 return;
277 };
278 let mut state = WriterState {
279 path: self.path.clone(),
280 webhook: self.webhook.clone(),
281 next_seq: 0,
282 recovered: false,
283 receiver,
284 };
285 self.writer_spawned.store(true, Ordering::Release);
286 handle.spawn(async move {
287 state.run().await;
288 });
289 }
290 }
291
292 /// The outbox writer: owns the file state and the event queue drain loop.
293 struct WriterState {
294 path: PathBuf,
295 webhook: Option<WebhookHookSink>,
296 /// Next seq to assign; filled in by [`Self::recover_seq`] on first use.
297 next_seq: u64,
298 recovered: bool,
299 receiver: Receiver<OutboxCommand>,
300 }
301
302 impl WriterState {
303 /// Drain the queue until every sender is dropped, then exit.
304 ///
305 /// Events are gathered into batches ([`OUTBOX_MAX_BATCH`]) so a burst
306 /// shares one file open/write/flush. Webhook posts fan out through a
307 /// bounded [`JoinSet`](tokio::task::JoinSet) and never delay the appends
308 /// of later batches — a stalled webhook only holds back its own event
309 /// plus the flush barrier, which the caller's timeout bounds (#6212).
310 async fn run(&mut self) {
311 let mut batch: Vec<OutboxCommand> = Vec::with_capacity(OUTBOX_MAX_BATCH);
312 let mut webhooks = tokio::task::JoinSet::new();
313 while self.receiver.recv_many(&mut batch, OUTBOX_MAX_BATCH).await > 0 {
314 let mut events = Vec::new();
315 let mut replies = Vec::new();
316 for command in batch.drain(..) {
317 match command {
318 OutboxCommand::Event(event) => events.push(event),
319 OutboxCommand::Flush(reply) => replies.push(reply),
320 }
321 }
322 if let Err(error) = self.deliver_batch(&events, &mut webhooks).await {
323 tracing::warn!(
324 target: "lifecycle_outbox",
325 %error,
326 path = %self.path.display(),
327 "lifecycle outbox write failed"
328 );
329 }
330 if !replies.is_empty() {
331 // The flush barrier covers every delivery attempt queued
332 // before it — appends above plus webhook attempts already
333 // spawned (mirroring the pre-batching serial writer, where
334 // a Flush was only reached after prior webhooks finished).
335 while webhooks.join_next().await.is_some() {}
336 for reply in replies {
337 let _ = reply.send(());
338 }
339 }
340 }
341 }
342
343 /// Assign a seq per event, build the envelopes, append the whole batch in
344 /// one file open/write/flush, then fan the batch's webhook posts out
345 /// concurrently (independently of the append result).
346 async fn deliver_batch(
347 &mut self,
348 events: &[LifecycleEvent],
349 webhooks: &mut tokio::task::JoinSet<()>,
350 ) -> Result<()> {
351 if events.is_empty() {
352 return Ok(());
353 }
354 if !self.recovered {
355 self.next_seq = recover_last_seq(&self.path).await?;
356 self.recovered = true;
357 }
358
359 let mut lines = Vec::with_capacity(events.len());
360 let mut webhook_posts = Vec::new();
361 for event in events {
362 let seq = self.next_seq;
363 self.next_seq = self.next_seq.saturating_add(1);
364
365 let envelope = RuntimeEventEnvelope {
366 schema_version: RUNTIME_EVENT_ENVELOPE_SCHEMA_VERSION,
367 seq,
368 event: event.event.clone(),
369 kind: event.kind.clone(),
370 thread_id: event.thread_id.clone(),
371 turn_id: event.turn_id.clone(),
372 item_id: event.item_id.clone(),
373 timestamp: Utc::now().to_rfc3339(),
374 created_at: Some(Utc::now().to_rfc3339()),
375 payload: event.payload.clone(),
376 extra: Default::default(),
377 };
378 if self.webhook.is_some() {
379 webhook_posts.push(json!({
380 "at": envelope.timestamp,
381 "event": envelope,
382 }));
383 }
384 lines.push(serde_json::to_string(&envelope).context("failed to encode outbox event")?);
385 }
386
387 // Appends land before any webhook work so the audit log never waits
388 // on a slow remote.
389 let append_result = self.append_lines(&lines).await;
390
391 if let Some(webhook) = &self.webhook {
392 for payload in webhook_posts {
393 while webhooks.len() >= WEBHOOK_MAX_INFLIGHT {
394 let _ = webhooks.join_next().await;
395 }
396 let sink = webhook.clone();
397 webhooks.spawn(async move {
398 if let Err(error) = sink.post_payload(payload).await {
399 tracing::warn!(
400 target: "lifecycle_outbox",
401 %error,
402 "lifecycle webhook delivery failed (dropped)"
403 );
404 }
405 });
406 }
407 }
408
409 append_result
410 }
411
412 /// Append complete JSONL lines, mirroring [`crate::JsonlHookSink`]:
413 /// lazy parent directories, append mode, flush before returning. The
414 /// writer task is the only appender for this outbox, so no extra lock is
415 /// needed here; the queue already serializes. The whole batch shares one
416 /// open and one flush; each line still lands as its own complete record.
417 async fn append_lines(&mut self, lines: &[String]) -> Result<()> {
418 let mut file = crate::open_private_append(&self.path).await?;
419 // Line + newline in a single `write_all` per record: with O_APPEND
420 // each `write` lands contiguously, so even a second process appending
421 // to the same file can interleave lines but can never splice one
422 // mid-line.
423 let mut records = Vec::with_capacity(lines.len());
424 for line in lines {
425 let mut record = Vec::with_capacity(line.len() + 1);
426 record.extend_from_slice(line.as_bytes());
427 record.push(b'\n');
428 records.push(record);
429 }
430 for record in &records {
431 file.write_all(record)
432 .await
433 .context("failed to write outbox event")?;
434 }
435 file.flush().await.context("failed to flush outbox event")
436 }
437 }
438
439 /// Recover the seq to continue from: the `seq` of the outbox file's last
440 /// complete line, plus 1 — or 1 for a missing/empty file.
441 ///
442 /// Only the tail of the file is read (bounded by [`SEQ_RECOVERY_TAIL_BYTES`]);
443 /// outbox lines are bounded far below that window, so the last complete line
444 /// is always within it. A partial trailing line from a crash mid-write is
445 /// ignored (the previous newline-terminated line wins).
446 async fn recover_last_seq(path: &Path) -> Result<u64> {
447 let mut file = match tokio::fs::File::open(path).await {
448 Ok(file) => file,
449 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(1),
450 Err(error) => {
451 return Err(error).with_context(|| format!("failed to open outbox {}", path.display()));
452 }
453 };
454 let len = file
455 .metadata()
456 .await
457 .with_context(|| format!("failed to stat outbox {}", path.display()))?
458 .len();
459 if len == 0 {
460 return Ok(1);
461 }
462 let start = len.saturating_sub(SEQ_RECOVERY_TAIL_BYTES);
463 file.seek(std::io::SeekFrom::Start(start)).await?;
464 let mut tail = vec![0u8; (len - start) as usize];
465 file.read_exact(&mut tail).await?;
466
467 let line = match tail.iter().rposition(|byte| *byte == b'\n') {
468 // The bytes after the final newline are a torn trailing line from a
469 // crash mid-write; drop them. What remains ends at a newline, so the
470 // last complete line is the bytes after the previous newline.
471 Some(last_nl) => {
472 let body = &tail[..last_nl];
473 match body.iter().rposition(|byte| *byte == b'\n') {
474 Some(idx) => &body[idx + 1..],
475 None => body,
476 }
477 }
478 // No newline at all: no complete line inside this tail (a line can
479 // only exceed the tail window by violating the bounded-line
480 // invariant). Treat the file as not-yet-writable.
481 None => return Ok(1),
482 };
483 let line = std::str::from_utf8(line).context("outbox tail is not UTF-8")?;
484 if line.trim().is_empty() {
485 return Ok(1);
486 }
487 let envelope: RuntimeEventEnvelope =
488 serde_json::from_str(line).context("failed to parse last outbox line")?;
489 Ok(envelope.seq.saturating_add(1))
490 }
491
492 /// Bound free-form text to at most `max_chars` characters, stripping control
493 /// bytes and ANSI escape sequences and collapsing whitespace runs first.
494 ///
495 /// The limit counts Unicode scalar values, not bytes, so multi-byte text gets
496 /// the same ceiling as ASCII. The result is safe to embed in an outbox
497 /// payload. Callers remain responsible for only ever passing non-secret
498 /// fields (error messages, previews, model/provider labels — never raw tool
499 /// arguments, environment, or full transcript text), the same discipline the
500 /// desktop notification payloads enforce.
501 pub fn bounded_text(text: &str, max_chars: usize) -> String {
502 let cleaned: String = text
503 .chars()
504 .filter(|ch| !ch.is_control())
505 .collect::<String>()
506 .split_whitespace()
507 .collect::<Vec<_>>()
508 .join(" ");
509 let mut truncated = false;
510 let mut out = String::new();
511 let mut char_count = 0usize;
512 for ch in cleaned.chars() {
513 if char_count + 1 > max_chars {
514 truncated = true;
515 break;
516 }
517 out.push(ch);
518 char_count += 1;
519 }
520 if truncated {
521 // Make room for the marker while staying under the character ceiling.
522 let marker_chars = OUTBOX_TRUNCATION_MARKER.chars().count();
523 while char_count + marker_chars > max_chars {
524 out.pop();
525 char_count -= 1;
526 }
527 out.push_str(OUTBOX_TRUNCATION_MARKER);
528 }
529 out
530 }
531
532 #[cfg(test)]
533 mod tests {
534 use super::*;
535
536 fn temp_outbox_path(name: &str) -> (tempfile::TempDir, PathBuf) {
537 let dir = tempfile::tempdir().expect("tempdir");
538 let path = dir.path().join(name);
539 (dir, path)
540 }
541
542 fn event(name: &str, kind: &str) -> LifecycleEvent {
543 LifecycleEvent {
544 event: name.to_string(),
545 kind: kind.to_string(),
546 thread_id: "session-1".to_string(),
547 turn_id: Some("turn-1".to_string()),
548 item_id: None,
549 payload: json!({"status": "completed"}),
550 }
551 }
552
553 async fn deliver_all(state: &mut WriterState, events: Vec<LifecycleEvent>) {
554 let mut webhooks = tokio::task::JoinSet::new();
555 state
556 .deliver_batch(&events, &mut webhooks)
557 .await
558 .expect("deliver");
559 while webhooks.join_next().await.is_some() {}
560 }
561
562 async fn read_lines(path: &Path) -> Vec<Value> {
563 let text = tokio::fs::read_to_string(path).await.expect("read outbox");
564 text.lines()
565 .map(|line| serde_json::from_str::<Value>(line).expect("json line"))
566 .collect()
567 }
568
569 #[test]
570 fn flush_persists_queued_turn_boundaries_before_runtime_teardown() {
571 let (_dir, path) = temp_outbox_path("shutdown.jsonl");
572 let runtime = tokio::runtime::Builder::new_current_thread()
573 .enable_all()
574 .build()
575 .expect("runtime");
576 runtime.block_on(async {
577 let outbox = LifecycleOutbox::new(Some(path.clone()), None, None);
578 outbox.emit(event("turn_start", "turn.started"));
579 let clone = outbox.clone();
580 clone.emit(event("turn_end", "turn.completed"));
581 outbox.flush(Duration::from_secs(2)).await.expect("flush");
582 });
583 drop(runtime);
584
585 let contents = std::fs::read_to_string(path).expect("outbox persisted before exit");
586 let lines: Vec<Value> = contents
587 .lines()
588 .map(|line| serde_json::from_str(line).expect("complete JSONL record"))
589 .collect();
590 assert_eq!(lines.len(), 2);
591 assert_eq!(lines[0]["event"], "turn_start");
592 assert_eq!(lines[0]["seq"], 1);
593 assert_eq!(lines[1]["event"], "turn_end");
594 assert_eq!(lines[1]["seq"], 2);
595 }
596
597 #[tokio::test]
598 async fn flush_disabled_or_unused_outbox_creates_no_file() {
599 let (_dir, path) = temp_outbox_path("unused.jsonl");
600 LifecycleOutbox::disabled()
601 .flush(Duration::ZERO)
602 .await
603 .expect("disabled flush");
604 LifecycleOutbox::new(Some(path.clone()), None, None)
605 .flush(Duration::from_secs(2))
606 .await
607 .expect("unused flush");
608 assert!(!path.exists());
609 }
610
611 /// A stalled webhook must not head-of-line block the audit log: the
612 /// appends of later events land while the slow post is still in flight
613 /// (#6212). The pre-batching writer awaited each webhook inline, so the
614 /// second event's line could not appear until the first post resolved.
615 #[tokio::test]
616 async fn stalled_webhook_does_not_block_later_appends() {
617 let (_dir, path) = temp_outbox_path("webhook-head-of-line.jsonl");
618 let server = wiremock::MockServer::start().await;
619 wiremock::Mock::given(wiremock::matchers::method("POST"))
620 .respond_with(wiremock::ResponseTemplate::new(200).set_delay(Duration::from_secs(5)))
621 .mount(&server)
622 .await;
623 let outbox = LifecycleOutbox::new(Some(path.clone()), Some(server.uri()), None);
624
625 outbox.emit(event("turn_start", "turn.started"));
626 // Give the writer time to start the stalled webhook post for the
627 // first event before queueing the second.
628 tokio::time::sleep(Duration::from_millis(100)).await;
629 outbox.emit(event("turn_end", "turn.completed"));
630
631 for _ in 0..100 {
632 if read_lines(&path).await.len() >= 2 {
633 break;
634 }
635 tokio::time::sleep(Duration::from_millis(10)).await;
636 }
637 let lines = read_lines(&path).await;
638 assert_eq!(
639 lines.len(),
640 2,
641 "the second event must be appended while the first webhook is still stalled"
642 );
643 assert_eq!(lines[1]["event"], "turn_end");
644 }
645
646 #[tokio::test]
647 async fn flush_is_bounded_when_webhook_delivery_stalls() {
648 let (_dir, path) = temp_outbox_path("stalled.jsonl");
649 let server = wiremock::MockServer::start().await;
650 wiremock::Mock::given(wiremock::matchers::method("POST"))
651 .respond_with(wiremock::ResponseTemplate::new(200).set_delay(Duration::from_secs(5)))
652 .mount(&server)
653 .await;
654 let outbox = LifecycleOutbox::new(Some(path), Some(server.uri()), None);
655 outbox.emit(event("turn_start", "turn.started"));
656 let error = tokio::time::timeout(
657 Duration::from_secs(1),
658 outbox.flush(Duration::from_millis(20)),
659 )
660 .await
661 .expect("the flush deadline must bound a blocked writer")
662 .expect_err("stalled delivery must report incomplete drain");
663 assert!(error.to_string().contains("flush timed out"));
664 }
665
666 #[test]
667 fn flush_reports_a_writer_lost_with_its_runtime() {
668 let (_dir, path) = temp_outbox_path("stopped.jsonl");
669 let outbox = LifecycleOutbox::new(Some(path), None, None);
670 let runtime = tokio::runtime::Builder::new_current_thread()
671 .enable_all()
672 .build()
673 .expect("first runtime");
674 runtime.block_on(async {
675 outbox.emit(event("turn_start", "turn.started"));
676 outbox
677 .flush(Duration::from_secs(2))
678 .await
679 .expect("first flush");
680 });
681 drop(runtime);
682 let runtime = tokio::runtime::Builder::new_current_thread()
683 .enable_all()
684 .build()
685 .expect("second runtime");
686 let error = runtime
687 .block_on(outbox.flush(Duration::from_secs(2)))
688 .expect_err("a closed writer cannot acknowledge delivery");
689 assert!(error.to_string().contains("writer task is gone"));
690 }
691
692 #[tokio::test]
693 async fn appends_one_jsonl_line_per_event_with_envelope_schema() {
694 let (_dir, path) = temp_outbox_path("schema.jsonl");
695 let mut state = WriterState {
696 path: path.clone(),
697 webhook: None,
698 next_seq: 0,
699 recovered: false,
700 receiver: tokio::sync::mpsc::channel(OUTBOX_QUEUE_CAPACITY).1,
701 };
702 deliver_all(&mut state, vec![event("turn_start", "turn.started")]).await;
703
704 let lines = read_lines(&path).await;
705 assert_eq!(lines.len(), 1);
706 let line = &lines[0];
707 assert_eq!(line["schema_version"], 1);
708 assert_eq!(line["seq"], 1);
709 assert_eq!(line["event"], "turn_start");
710 assert_eq!(line["kind"], "turn.started");
711 assert_eq!(line["thread_id"], "session-1");
712 assert_eq!(line["turn_id"], "turn-1");
713 assert_eq!(line["item_id"], Value::Null);
714 assert!(line["timestamp"].as_str().is_some());
715 assert!(line["payload"]["status"].as_str() == Some("completed"));
716 }
717
718 /// Every emit site now carries `payload.workspace` (and subagent events
719 /// additionally `payload.subagent`) for consumer-side routing. The writer
720 /// must preserve those fields verbatim through the envelope round trip
721 /// for every event type.
722 #[tokio::test]
723 async fn payload_workspace_and_subagent_fields_survive_the_round_trip() {
724 let (_dir, path) = temp_outbox_path("routing-fields.jsonl");
725 let mut state = WriterState {
726 path: path.clone(),
727 webhook: None,
728 next_seq: 0,
729 recovered: false,
730 receiver: tokio::sync::mpsc::channel(OUTBOX_QUEUE_CAPACITY).1,
731 };
732 let workspace = "/home/cw/wt-lane";
733 let subagent = "explore-1";
734 let subagent_payload = json!({ "workspace": workspace, "subagent": subagent });
735 deliver_all(
736 &mut state,
737 vec![
738 LifecycleEvent {
739 event: "session_start".to_string(),
740 kind: "session.started".to_string(),
741 thread_id: "session-1".to_string(),
742 turn_id: None,
743 item_id: None,
744 payload: json!({ "workspace": workspace }),
745 },
746 LifecycleEvent {
747 event: "turn_start".to_string(),
748 kind: "turn.started".to_string(),
749 thread_id: "session-1".to_string(),
750 turn_id: Some("turn-1".to_string()),
751 item_id: None,
752 payload: json!({ "workspace": workspace }),
753 },
754 LifecycleEvent {
755 event: "turn_end".to_string(),
756 kind: "turn.completed".to_string(),
757 thread_id: "session-1".to_string(),
758 turn_id: Some("turn-1".to_string()),
759 item_id: None,
760 payload: json!({ "workspace": workspace }),
761 },
762 LifecycleEvent {
763 event: "turn_stalled".to_string(),
764 kind: "turn.stalled".to_string(),
765 thread_id: "session-1".to_string(),
766 turn_id: Some("turn-1".to_string()),
767 item_id: None,
768 payload: json!({ "workspace": workspace }),
769 },
770 LifecycleEvent {
771 event: "subagent_spawn".to_string(),
772 kind: "subagent.spawned".to_string(),
773 thread_id: "session-1".to_string(),
774 turn_id: Some("turn-1".to_string()),
775 item_id: None,
776 payload: subagent_payload.clone(),
777 },
778 LifecycleEvent {
779 event: "subagent_complete".to_string(),
780 kind: "subagent.completed".to_string(),
781 thread_id: "session-1".to_string(),
782 turn_id: Some("turn-1".to_string()),
783 item_id: None,
784 payload: subagent_payload.clone(),
785 },
786 LifecycleEvent {
787 event: "session_end".to_string(),
788 kind: "session.ended".to_string(),
789 thread_id: "session-1".to_string(),
790 turn_id: None,
791 item_id: None,
792 payload: json!({ "workspace": workspace }),
793 },
794 ],
795 )
796 .await;
797
798 let lines = read_lines(&path).await;
799 let events: Vec<&str> = lines
800 .iter()
801 .map(|line| line["event"].as_str().expect("event"))
802 .collect();
803 assert_eq!(
804 events,
805 vec![
806 "session_start",
807 "turn_start",
808 "turn_end",
809 "turn_stalled",
810 "subagent_spawn",
811 "subagent_complete",
812 "session_end",
813 ],
814 "the routing-field contract must cover every lifecycle event type"
815 );
816 for line in &lines {
817 assert_eq!(
818 line["payload"]["workspace"],
819 json!(workspace),
820 "workspace must survive the round trip for event {}",
821 line["event"]
822 );
823 }
824 for event in ["subagent_spawn", "subagent_complete"] {
825 let line = lines
826 .iter()
827 .find(|line| line["event"] == event)
828 .expect(event);
829 assert_eq!(
830 line["payload"]["subagent"],
831 json!(subagent),
832 "subagent must survive the round trip for event {event}"
833 );
834 }
835 }
836
837 #[tokio::test]
838 async fn seq_is_monotonic_and_recovers_across_reopen() {
839 let (_dir, path) = temp_outbox_path("seq.jsonl");
840 let mut state = WriterState {
841 path: path.clone(),
842 webhook: None,
843 next_seq: 0,
844 recovered: false,
845 receiver: tokio::sync::mpsc::channel(OUTBOX_QUEUE_CAPACITY).1,
846 };
847 deliver_all(
848 &mut state,
849 vec![
850 event("session_start", "session.started"),
851 event("turn_start", "turn.started"),
852 event("turn_end", "turn.completed"),
853 ],
854 )
855 .await;
856
857 // A fresh writer (new process, same file) continues the sequence.
858 let mut reopened = WriterState {
859 path: path.clone(),
860 webhook: None,
861 next_seq: 0,
862 recovered: false,
863 receiver: tokio::sync::mpsc::channel(OUTBOX_QUEUE_CAPACITY).1,
864 };
865 deliver_all(&mut reopened, vec![event("turn_start", "turn.started")]).await;
866
867 let lines = read_lines(&path).await;
868 let seqs: Vec<u64> = lines
869 .iter()
870 .map(|line| line["seq"].as_u64().expect("seq"))
871 .collect();
872 assert_eq!(seqs, vec![1, 2, 3, 4]);
873 }
874
875 #[tokio::test]
876 async fn missing_and_empty_files_start_at_seq_1() {
877 let (_dir, path) = temp_outbox_path("empty.jsonl");
878 assert_eq!(recover_last_seq(&path).await.expect("missing file"), 1);
879
880 tokio::fs::write(&path, "").await.expect("empty file");
881 assert_eq!(recover_last_seq(&path).await.expect("empty file"), 1);
882 }
883
884 #[tokio::test]
885 async fn partial_trailing_line_is_ignored_during_recovery() {
886 let (_dir, path) = temp_outbox_path("partial.jsonl");
887 tokio::fs::write(
888 &path,
889 format!(
890 "{}\n{}\n{{\"schema_version\":1,\"seq\":3,\"event\":\"turn_",
891 r#"{"schema_version":1,"seq":1,"event":"session_start","kind":"session.started","thread_id":"s","turn_id":null,"item_id":null,"timestamp":"t","payload":{}}"#,
892 r#"{"schema_version":1,"seq":2,"event":"turn_start","kind":"turn.started","thread_id":"s","turn_id":null,"item_id":null,"timestamp":"t","payload":{}}"#,
893 ),
894 )
895 .await
896 .expect("write partial outbox");
897 // The torn trailing line is not a complete record; recovery continues
898 // from the last complete line's seq (2) → next seq 3.
899 assert_eq!(recover_last_seq(&path).await.expect("recover"), 3);
900 }
901
902 #[tokio::test]
903 async fn emit_queues_and_writes_in_order_without_blocking() {
904 let (_dir, path) = temp_outbox_path("emit.jsonl");
905 let outbox = LifecycleOutbox::new(Some(path.clone()), None, None);
906 assert!(outbox.is_enabled());
907
908 outbox.emit(event("session_start", "session.started"));
909 outbox.emit(event("turn_start", "turn.started"));
910 outbox.emit(event("turn_end", "turn.completed"));
911
912 // The writer task drains asynchronously; wait for the lines to land.
913 for _ in 0..100 {
914 if tokio::fs::metadata(&path)
915 .await
916 .is_ok_and(|meta| meta.len() > 0)
917 && read_lines(&path).await.len() >= 3
918 {
919 break;
920 }
921 tokio::time::sleep(std::time::Duration::from_millis(10)).await;
922 }
923 let lines = read_lines(&path).await;
924 assert_eq!(lines.len(), 3, "expected all queued events to be written");
925 let events: Vec<&str> = lines
926 .iter()
927 .map(|line| line["event"].as_str().expect("event"))
928 .collect();
929 assert_eq!(events, vec!["session_start", "turn_start", "turn_end"]);
930 let seqs: Vec<u64> = lines
931 .iter()
932 .map(|line| line["seq"].as_u64().expect("seq"))
933 .collect();
934 assert_eq!(seqs, vec![1, 2, 3], "seq must be assigned in emit order");
935 }
936
937 #[test]
938 fn disabled_outbox_drops_events_and_reports_disabled() {
939 let outbox = LifecycleOutbox::new(None, None, None);
940 assert!(!outbox.is_enabled());
941 outbox.emit(event("turn_start", "turn.started")); // must not panic
942
943 let empty_path = LifecycleOutbox::new(Some(PathBuf::new()), None, None);
944 assert!(!empty_path.is_enabled());
945
946 let default = LifecycleOutbox::default();
947 assert!(!default.is_enabled());
948 }
949
950 #[test]
951 fn webhook_only_configures_without_a_file_path() {
952 // `webhook_url` without `path` is stored losslessly in config; the
953 // outbox handle itself only activates on a path.
954 let outbox = LifecycleOutbox::new(
955 None,
956 Some("https://example.com/hook".to_string()),
957 Some("token".to_string()),
958 );
959 assert!(!outbox.is_enabled());
960 }
961
962 #[test]
963 fn bounded_text_truncates_to_limit_with_marker() {
964 assert_eq!(bounded_text("short", 80), "short");
965 let long = "x".repeat(200);
966 let bounded = bounded_text(&long, OUTBOX_DETAIL_MAX_CHARS);
967 assert_eq!(bounded.chars().count(), OUTBOX_DETAIL_MAX_CHARS);
968 assert!(bounded.ends_with(OUTBOX_TRUNCATION_MARKER));
969 assert!(bounded.starts_with('x'));
970 }
971
972 #[test]
973 fn bounded_text_strips_controls_and_collapses_whitespace() {
974 assert_eq!(
975 bounded_text("line\x1b[31m one\n\n two\t", 80),
976 "line[31m one two"
977 );
978 assert_eq!(bounded_text("", 80), "");
979 assert_eq!(bounded_text(" \n\t ", 80), "");
980 }
981
982 #[test]
983 fn bounded_text_respects_utf8_boundaries() {
984 // 30 multi-byte emoji (4 bytes each) = 120 bytes but only 30 chars.
985 let emoji = "🦈".repeat(30);
986 let bounded = bounded_text(&emoji, OUTBOX_DETAIL_MAX_CHARS);
987 assert!(bounded.chars().count() <= OUTBOX_DETAIL_MAX_CHARS);
988 assert!(bounded.starts_with('🦈'));
989 }
990
991 /// The webhook transport must POST `{"at", "event"}` JSON and, when a
992 /// token is configured, send it as `Authorization: Bearer <token>`.
993 #[tokio::test]
994 async fn webhook_posts_at_event_payload_with_bearer_token() {
995 let server = wiremock::MockServer::start().await;
996 wiremock::Mock::given(wiremock::matchers::method("POST"))
997 .and(wiremock::matchers::path("/hook"))
998 .and(wiremock::matchers::header(
999 "authorization",
1000 "Bearer secret-token",
1001 ))
1002 .and(wiremock::matchers::body_partial_json(json!({
1003 "event": {"kind": "turn.started"}
1004 })))
1005 .respond_with(wiremock::ResponseTemplate::new(200))
1006 .mount(&server)
1007 .await;
1008
1009 let webhook = WebhookHookSink::new_with_token(
1010 format!("{}/hook", server.uri()),
1011 Some("secret-token".to_string()),
1012 );
1013 webhook
1014 .post_payload(json!(
1015 {"at": "2026-08-19T00:00:00Z", "event": {"kind": "turn.started"}}
1016 ))
1017 .await
1018 .expect("webhook delivery");
1019
1020 let requests = server.received_requests().await.expect("requests");
1021 assert_eq!(requests.len(), 1, "exactly one webhook POST");
1022 }
1023
1024 /// A webhook that always fails must surface its error to the caller
1025 /// (which logs and drops it) — never panic, never retry forever.
1026 #[tokio::test]
1027 async fn webhook_failure_is_an_error_not_a_panic() {
1028 let server = wiremock::MockServer::start().await;
1029 wiremock::Mock::given(wiremock::matchers::method("POST"))
1030 .respond_with(wiremock::ResponseTemplate::new(500))
1031 .mount(&server)
1032 .await;
1033
1034 let webhook = WebhookHookSink::new_with_token(format!("{}/hook", server.uri()), None);
1035 let result = webhook.post_payload(json!({})).await;
1036 assert!(result.is_err(), "expected the failure to be reported");
1037 }
1038 }
1039
1039 lines RUST