1use std::sync::{Arc, Mutex};
2use std::time::Duration;
3
4use serde::{Deserialize, Serialize};
5use tokio_util::sync::CancellationToken;
6
7use crate::{ErrorEnvelope, TurnOutcome};
8
9use super::{
10 AwaitEventKey, AwaitEventResolver, AwaitEventWaitIdentity, EffectHost, ExecutionScope,
11 Resolution, ResolveOutcome, RuntimeEffectCommand, RuntimeEffectController,
12 RuntimeEffectEnvelope, RuntimeEffectKind, RuntimeEffectLocalExecutor, RuntimeEffectOutcome,
13 RuntimeError, RuntimeInvocation, RuntimeScope,
14};
15
16#[derive(Clone, Copy, Debug, PartialEq, Eq)]
17pub(crate) enum TurnCancelPeekIdentity {
18 StartGate,
19 PostAbortGate,
20}
21
22impl TurnCancelPeekIdentity {
23 fn as_str(self) -> &'static str {
24 match self {
25 Self::StartGate => "turn_cancel.start_gate",
26 Self::PostAbortGate => "turn_cancel.post_abort_gate",
27 }
28 }
29}
30
31#[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
37pub struct TurnAddress {
38 pub session_id: String,
39 pub turn_id: String,
40}
41
42impl TurnAddress {
43 pub fn new(session_id: impl Into<String>, turn_id: impl Into<String>) -> Self {
44 Self {
45 session_id: session_id.into(),
46 turn_id: turn_id.into(),
47 }
48 }
49
50 fn scope(&self) -> ExecutionScope {
51 ExecutionScope::turn(&self.session_id, &self.turn_id)
52 }
53
54 fn validate(&self) -> Result<(), RuntimeError> {
55 self.scope().validate()
56 }
57}
58
59#[doc(hidden)]
66#[derive(Clone, Default)]
67pub struct TurnCancelOriginHint {
68 origin: Arc<Mutex<Option<Option<String>>>>,
69}
70
71impl TurnCancelOriginHint {
72 pub fn set(&self, origin: Option<String>) {
73 let mut hint = self.origin.lock().expect("turn cancel origin hint lock");
74 if hint.is_none() {
75 *hint = Some(origin);
76 }
77 }
78
79 pub(crate) fn get(&self) -> Option<String> {
80 self.origin
81 .lock()
82 .expect("turn cancel origin hint lock")
83 .clone()
84 .flatten()
85 }
86}
87
88#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
89pub struct TurnCancellationEvidence {
90 pub request_id: String,
91 #[serde(default, skip_serializing_if = "Option::is_none")]
93 pub origin: Option<String>,
94 #[serde(default, skip_serializing_if = "Option::is_none")]
95 pub reason: Option<String>,
96}
97
98#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
99pub struct TurnCancelRequest {
100 pub address: TurnAddress,
101 pub request_id: String,
102 #[serde(default, skip_serializing_if = "Option::is_none")]
104 pub origin: Option<String>,
105 #[serde(default, skip_serializing_if = "Option::is_none")]
106 pub reason: Option<String>,
107}
108
109impl TurnCancelRequest {
110 pub fn new(
111 address: TurnAddress,
112 request_id: impl Into<String>,
113 origin: Option<String>,
114 ) -> Self {
115 Self {
116 address,
117 request_id: request_id.into(),
118 origin,
119 reason: None,
120 }
121 }
122
123 pub fn with_reason(mut self, reason: impl Into<String>) -> Self {
124 self.reason = Some(reason.into());
125 self
126 }
127
128 fn validate(&self) -> Result<(), RuntimeError> {
129 self.address.validate()?;
130 if self.request_id.trim().is_empty() {
131 return Err(RuntimeError::new(
132 "invalid_turn_cancel_request",
133 "turn cancellation requires a non-empty request id",
134 ));
135 }
136 Ok(())
137 }
138
139 fn evidence(&self) -> TurnCancellationEvidence {
140 TurnCancellationEvidence {
141 request_id: self.request_id.clone(),
142 origin: self.origin.clone(),
143 reason: self.reason.clone(),
144 }
145 }
146}
147
148#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
149#[serde(tag = "outcome", content = "cancellation", rename_all = "snake_case")]
150pub enum TurnCancelOutcome {
151 Requested(TurnCancellationEvidence),
152 AlreadyRequested(TurnCancellationEvidence),
153 CompletionWonRace,
154 UnknownOrRevoked,
155}
156
157#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
168pub struct TurnCancelReceipt {
169 pub durability_tier: crate::DurabilityTier,
170 pub outcome: TurnCancelOutcome,
171}
172
173#[derive(Clone, Debug, Serialize, Deserialize)]
174#[serde(tag = "status", rename_all = "snake_case")]
175pub enum TurnTerminal {
176 Committed {
177 outcome: TurnOutcome,
178 #[serde(default, skip_serializing_if = "Option::is_none")]
179 cancellation: Option<TurnCancellationEvidence>,
180 #[serde(default, skip_serializing_if = "Option::is_none")]
181 session_revision: Option<u64>,
182 },
183 Failed {
184 error: ErrorEnvelope,
185 },
186}
187
188#[async_trait::async_trait]
190pub trait TurnAttach: Send + Sync {
191 async fn await_terminal(&self, address: &TurnAddress) -> Result<TurnTerminal, RuntimeError>;
192}
193
194#[derive(Clone)]
212pub struct TurnWorkDriver {
213 effect_host: Arc<dyn EffectHost>,
214 attach: Option<Arc<dyn TurnAttach>>,
215}
216
217impl TurnWorkDriver {
218 pub fn new(effect_host: Arc<dyn EffectHost>) -> Self {
219 Self {
220 effect_host,
221 attach: None,
222 }
223 }
224
225 pub fn with_attach(mut self, attach: Arc<dyn TurnAttach>) -> Self {
226 self.attach = Some(attach);
227 self
228 }
229
230 pub fn effect_host(&self) -> Arc<dyn EffectHost> {
231 Arc::clone(&self.effect_host)
232 }
233
234 pub async fn request_cancel(
235 &self,
236 request: TurnCancelRequest,
237 ) -> Result<TurnCancelReceipt, RuntimeError> {
238 request.validate()?;
239 let durability_tier = self.effect_host.durability_tier();
240 let key = match cancel_gate_key(self.effect_host.as_ref(), &request.address).await {
241 Ok(key) => key,
242 Err(err) if err.code.as_str() == "await_event_unknown_or_revoked" => {
243 return Ok(TurnCancelReceipt {
244 durability_tier,
245 outcome: TurnCancelOutcome::UnknownOrRevoked,
246 });
247 }
248 Err(err) => return Err(err),
249 };
250 let evidence = request.evidence();
251 let resolution = gate_resolution(TurnGateTerminal::CancelRequested(evidence.clone()))?;
252 let outcome = match self
253 .effect_host
254 .resolve_await_event(&key, resolution)
255 .await?
256 {
257 ResolveOutcome::Accepted => Ok(TurnCancelOutcome::Requested(evidence)),
258 ResolveOutcome::AlreadyResolved { terminal } => match decode_gate(terminal)? {
259 TurnGateTerminal::CancelRequested(existing) => {
260 Ok(TurnCancelOutcome::AlreadyRequested(existing))
261 }
262 TurnGateTerminal::CompletionSealed => Ok(TurnCancelOutcome::CompletionWonRace),
263 },
264 ResolveOutcome::UnknownOrRevoked => Ok(TurnCancelOutcome::UnknownOrRevoked),
265 }?;
266 Ok(TurnCancelReceipt {
267 durability_tier,
268 outcome,
269 })
270 }
271
272 pub async fn await_terminal(
273 &self,
274 address: &TurnAddress,
275 ) -> Result<TurnTerminal, RuntimeError> {
276 address.validate()?;
277 if let Some(attach) = self.attach.as_ref() {
278 return attach.await_terminal(address).await;
279 }
280 let key = terminal_key(self.effect_host.as_ref(), address).await?;
281 let resolution = self
282 .effect_host
283 .await_await_event(&key, CancellationToken::new(), None)
284 .await?;
285 decode_terminal(address, resolution)
286 }
287
288 pub async fn await_terminal_with_timeout(
293 &self,
294 address: &TurnAddress,
295 timeout: Duration,
296 ) -> Result<TurnTerminal, RuntimeError> {
297 tokio::time::timeout(timeout, self.await_terminal(address))
298 .await
299 .map_err(|_| {
300 RuntimeError::new(
301 "turn_terminal_await_timeout",
302 format!(
303 "timed out awaiting terminal for turn `{}` in session `{}` after {} ms",
304 address.turn_id,
305 address.session_id,
306 timeout.as_millis()
307 ),
308 )
309 })?
310 }
311}
312
313#[derive(Clone, Debug, Serialize, Deserialize)]
314#[serde(tag = "state", content = "cancellation", rename_all = "snake_case")]
315enum TurnGateTerminal {
316 CancelRequested(TurnCancellationEvidence),
317 CompletionSealed,
318}
319
320fn gate_resolution(value: TurnGateTerminal) -> Result<Resolution, RuntimeError> {
321 serde_json::to_value(value)
322 .map(Resolution::Ok)
323 .map_err(|err| RuntimeError::new("turn_cancel_gate_encode", err.to_string()))
324}
325
326fn decode_gate(resolution: Resolution) -> Result<TurnGateTerminal, RuntimeError> {
327 match resolution {
328 Resolution::Ok(value) => serde_json::from_value(value)
329 .map_err(|err| RuntimeError::new("turn_cancel_gate_decode", err.to_string())),
330 other => Err(RuntimeError::new(
331 "turn_cancel_gate_invalid_terminal",
332 format!("turn cancellation gate resolved with {other:?}"),
333 )),
334 }
335}
336
337fn terminal_resolution(value: &TurnTerminal) -> Result<Resolution, RuntimeError> {
338 serde_json::to_value(value)
339 .map(Resolution::Ok)
340 .map_err(|err| RuntimeError::new("turn_terminal_encode", err.to_string()))
341}
342
343fn decode_terminal(
344 address: &TurnAddress,
345 resolution: Resolution,
346) -> Result<TurnTerminal, RuntimeError> {
347 match resolution {
348 Resolution::Ok(value) => serde_json::from_value(value).map_err(|err| {
349 RuntimeError::new(
350 "turn_terminal_decode",
351 format!(
352 "invalid terminal result for turn `{}` in session `{}`: {err}",
353 address.turn_id, address.session_id
354 ),
355 )
356 }),
357 other => Err(RuntimeError::new(
358 "turn_terminal_invalid_resolution",
359 format!(
360 "terminal result for turn `{}` in session `{}` resolved with {other:?}",
361 address.turn_id, address.session_id
362 ),
363 )),
364 }
365}
366
367async fn cancel_gate_key(
368 resolver: &dyn AwaitEventResolver,
369 address: &TurnAddress,
370) -> Result<AwaitEventKey, RuntimeError> {
371 resolver
372 .await_event_key(&address.scope(), AwaitEventWaitIdentity::TurnCancelGate)
373 .await
374}
375
376async fn terminal_key(
377 resolver: &dyn AwaitEventResolver,
378 address: &TurnAddress,
379) -> Result<AwaitEventKey, RuntimeError> {
380 resolver
381 .await_event_key(&address.scope(), AwaitEventWaitIdentity::TurnTerminal)
382 .await
383}
384
385pub(crate) struct ActiveTurnControl {
388 address: TurnAddress,
389 cancel_key: AwaitEventKey,
390 terminal_key: AwaitEventKey,
391 evidence: Mutex<Option<TurnCancellationEvidence>>,
392 local_cancel_origin: TurnCancelOriginHint,
393}
394
395impl ActiveTurnControl {
396 pub(crate) async fn new(
397 resolver: &dyn AwaitEventResolver,
398 address: TurnAddress,
399 ) -> Result<Self, RuntimeError> {
400 address.validate()?;
401 Ok(Self {
402 cancel_key: cancel_gate_key(resolver, &address).await?,
403 terminal_key: terminal_key(resolver, &address).await?,
404 address,
405 evidence: Mutex::new(None),
406 local_cancel_origin: TurnCancelOriginHint::default(),
407 })
408 }
409
410 pub(crate) fn with_local_cancel_origin(mut self, origin: TurnCancelOriginHint) -> Self {
411 self.local_cancel_origin = origin;
412 self
413 }
414
415 pub(crate) async fn await_cancel(
416 &self,
417 resolver: &dyn AwaitEventResolver,
418 stop_wait: CancellationToken,
419 ) -> Result<Option<TurnCancellationEvidence>, RuntimeError> {
420 let resolution = resolver
421 .await_await_event(&self.cancel_key, stop_wait, None)
422 .await?;
423 match decode_gate(resolution)? {
424 TurnGateTerminal::CancelRequested(evidence) => {
425 self.remember(evidence.clone());
426 Ok(Some(evidence))
427 }
428 TurnGateTerminal::CompletionSealed => Ok(None),
429 }
430 }
431
432 pub(crate) async fn observe_pending_cancel(
433 &self,
434 controller: &dyn RuntimeEffectController,
435 identity: TurnCancelPeekIdentity,
436 ) -> Result<Option<TurnCancellationEvidence>, RuntimeError> {
437 let causal_identity = identity.as_str();
438 let invocation = RuntimeInvocation::effect(
439 RuntimeScope {
440 session_id: self.address.session_id.clone(),
441 turn_id: Some(self.address.turn_id.clone()),
442 turn_index: None,
443 protocol_iteration: None,
444 },
445 causal_identity,
446 RuntimeEffectKind::PeekAwaitEvent,
447 causal_identity,
448 );
449 let outcome = controller
450 .execute_effect(
451 RuntimeEffectEnvelope::new(
452 invocation,
453 RuntimeEffectCommand::PeekAwaitEvent {
454 key: self.cancel_key.clone(),
455 },
456 ),
457 RuntimeEffectLocalExecutor::unavailable(),
458 )
459 .await
460 .map_err(|err| RuntimeError::new(err.code, err.message))?;
461 let RuntimeEffectOutcome::PeekAwaitEvent { resolution } = outcome else {
462 return Err(RuntimeError::new(
463 "turn_control_peek_outcome",
464 format!("{causal_identity} returned a non-peek runtime effect outcome"),
465 ));
466 };
467 let Some(resolution) = resolution else {
468 return Ok(None);
469 };
470 match decode_gate(resolution)? {
471 TurnGateTerminal::CancelRequested(evidence) => {
472 self.remember(evidence.clone());
473 Ok(Some(evidence))
474 }
475 TurnGateTerminal::CompletionSealed => Ok(None),
476 }
477 }
478
479 pub(crate) async fn settle_before_commit(
480 &self,
481 resolver: &dyn AwaitEventResolver,
482 locally_cancelled: bool,
483 ) -> Result<Option<TurnCancellationEvidence>, RuntimeError> {
484 if let Some(evidence) = self.evidence() {
485 return Ok(Some(evidence));
486 }
487 let proposed = if locally_cancelled {
488 TurnGateTerminal::CancelRequested(self.internal_evidence())
489 } else {
490 TurnGateTerminal::CompletionSealed
491 };
492 let outcome = resolver
493 .resolve_await_event(&self.cancel_key, gate_resolution(proposed.clone())?)
494 .await?;
495 let terminal = match outcome {
496 ResolveOutcome::Accepted => proposed,
497 ResolveOutcome::AlreadyResolved { terminal } => decode_gate(terminal)?,
498 ResolveOutcome::UnknownOrRevoked => {
499 return Err(RuntimeError::new(
500 "turn_control_unknown_or_revoked",
501 format!(
502 "turn `{}` in session `{}` was revoked before final commit",
503 self.address.turn_id, self.address.session_id
504 ),
505 ));
506 }
507 };
508 match terminal {
509 TurnGateTerminal::CancelRequested(evidence) => {
510 self.remember(evidence.clone());
511 Ok(Some(evidence))
512 }
513 TurnGateTerminal::CompletionSealed => Ok(None),
514 }
515 }
516
517 pub(crate) async fn publish_terminal(
518 &self,
519 resolver: &dyn AwaitEventResolver,
520 terminal: &TurnTerminal,
521 ) -> Result<(), RuntimeError> {
522 match resolver
523 .resolve_await_event(&self.terminal_key, terminal_resolution(terminal)?)
524 .await?
525 {
526 ResolveOutcome::Accepted | ResolveOutcome::AlreadyResolved { .. } => Ok(()),
527 ResolveOutcome::UnknownOrRevoked => Err(RuntimeError::new(
528 "turn_terminal_unknown_or_revoked",
529 format!(
530 "terminal promise for turn `{}` in session `{}` was revoked",
531 self.address.turn_id, self.address.session_id
532 ),
533 )),
534 }
535 }
536
537 pub(crate) fn evidence(&self) -> Option<TurnCancellationEvidence> {
538 self.evidence
539 .lock()
540 .expect("turn cancellation evidence lock")
541 .clone()
542 }
543
544 fn remember(&self, evidence: TurnCancellationEvidence) {
545 *self
546 .evidence
547 .lock()
548 .expect("turn cancellation evidence lock") = Some(evidence);
549 }
550
551 fn internal_evidence(&self) -> TurnCancellationEvidence {
552 TurnCancellationEvidence {
553 request_id: format!("internal:{}", self.address.turn_id),
554 origin: self.local_cancel_origin.get(),
555 reason: None,
556 }
557 }
558}
559
560#[cfg(test)]
561mod tests {
562 use super::*;
563 use crate::{InlineEffectHost, TurnFinish, TurnStop};
564
565 fn address(label: &str) -> TurnAddress {
566 TurnAddress::new(
567 format!("turn-control-{label}-{}", uuid::Uuid::new_v4()),
568 "turn-a",
569 )
570 }
571
572 fn request(address: TurnAddress, request_id: &str) -> TurnCancelRequest {
573 TurnCancelRequest::new(address, request_id, Some("user".to_string()))
574 .with_reason("stop button")
575 }
576
577 #[tokio::test]
578 async fn cancel_before_start_duplicate_and_terminal_attach() {
579 let host = Arc::new(InlineEffectHost::default());
580 let driver = TurnWorkDriver::new(host.clone());
581 let address = address("before-start");
582
583 let first = driver
584 .request_cancel(request(address.clone(), "request-1"))
585 .await
586 .expect("request cancellation");
587 assert_eq!(first.durability_tier, crate::DurabilityTier::Inline);
588 let evidence = match first.outcome {
589 TurnCancelOutcome::Requested(evidence) => evidence,
590 other => panic!("expected requested, got {other:?}"),
591 };
592 assert_eq!(evidence.request_id, "request-1");
593
594 let duplicate = driver
595 .request_cancel(request(address.clone(), "request-2"))
596 .await
597 .expect("duplicate cancellation");
598 assert!(matches!(
599 duplicate.outcome,
600 TurnCancelOutcome::AlreadyRequested(TurnCancellationEvidence { ref request_id, .. })
601 if request_id == "request-1"
602 ));
603
604 let active = ActiveTurnControl::new(host.as_ref(), address.clone())
605 .await
606 .expect("active control");
607 let observed = active
608 .settle_before_commit(host.as_ref(), false)
609 .await
610 .expect("settle")
611 .expect("cancellation won");
612 assert_eq!(observed, evidence);
613 let terminal = TurnTerminal::Committed {
614 outcome: TurnOutcome::Stopped(TurnStop::Cancelled),
615 cancellation: Some(observed),
616 session_revision: Some(7),
617 };
618 active
619 .publish_terminal(host.as_ref(), &terminal)
620 .await
621 .expect("publish terminal");
622 let attached = driver
623 .await_terminal(&address)
624 .await
625 .expect("attach terminal");
626 assert!(matches!(
627 attached,
628 TurnTerminal::Committed {
629 outcome: TurnOutcome::Stopped(TurnStop::Cancelled),
630 cancellation: Some(_),
631 session_revision: Some(7),
632 }
633 ));
634 }
635
636 #[tokio::test]
637 async fn concurrent_completion_seal_vs_cancel_is_first_writer_wins() {
638 let host = Arc::new(InlineEffectHost::default());
639 let driver = TurnWorkDriver::new(host.clone());
640 let address = address("race");
641 let active = ActiveTurnControl::new(host.as_ref(), address.clone())
642 .await
643 .expect("active control");
644
645 let (seal, cancel) = tokio::join!(
646 active.settle_before_commit(host.as_ref(), false),
647 driver.request_cancel(request(address, "race-request")),
648 );
649 match (seal.expect("seal"), cancel.expect("cancel").outcome) {
650 (None, TurnCancelOutcome::CompletionWonRace) => {}
651 (Some(evidence), TurnCancelOutcome::Requested(requested)) => {
652 assert_eq!(evidence, requested);
653 }
654 other => panic!("inconsistent gate race result: {other:?}"),
655 }
656 }
657
658 #[tokio::test]
659 async fn recovered_owner_observes_pending_cancel_after_control_recreation() {
660 let host = Arc::new(InlineEffectHost::default());
661 let driver = TurnWorkDriver::new(host.clone());
662 let address = address("replay");
663 let requested = driver
664 .request_cancel(request(address.clone(), "request-before-replay"))
665 .await
666 .expect("request cancellation");
667 let expected = match requested.outcome {
668 TurnCancelOutcome::Requested(evidence) => evidence,
669 other => panic!("expected requested, got {other:?}"),
670 };
671
672 let scoped = host
673 .scoped(address.scope())
674 .expect("scope recovered turn controller");
675 let recovered = ActiveTurnControl::new(host.as_ref(), address)
676 .await
677 .expect("recreate active control under the recovered owner");
678 let observed = recovered
679 .observe_pending_cancel(scoped.controller(), TurnCancelPeekIdentity::StartGate)
680 .await
681 .expect("read recovered turn start gate")
682 .expect("pending cancellation is visible before recovered effects");
683 assert_eq!(observed, expected);
684 let settled = recovered
685 .settle_before_commit(host.as_ref(), false)
686 .await
687 .expect("settle recovered turn")
688 .expect("pending cancellation survives owner loss");
689 assert_eq!(settled, expected);
690 }
691
692 #[tokio::test]
693 async fn turn_control_is_exact_scope_and_excluded_from_wait_cancel_sweep() {
694 let host = Arc::new(InlineEffectHost::default());
695 let driver = TurnWorkDriver::new(host.clone());
696 let address_a = address("scope");
697 let address_b = TurnAddress::new(&address_a.session_id, "turn-b");
698 let address_future = TurnAddress::new(&address_a.session_id, "turn-future");
699
700 driver
701 .request_cancel(request(address_a.clone(), "request-a"))
702 .await
703 .expect("cancel a");
704
705 let tool_key = host
706 .await_event_key(
707 &ExecutionScope::turn(&address_a.session_id, "tool-turn"),
708 AwaitEventWaitIdentity::tool_completion("tool-call"),
709 )
710 .await
711 .expect("tool key");
712 let tool_host = host.clone();
713 let tool_wait = crate::task::spawn(async move {
714 tool_host
715 .await_await_event(&tool_key, CancellationToken::new(), None)
716 .await
717 });
718 tokio::task::yield_now().await;
719 host.cancel_await_events_for_session(&address_a.session_id)
720 .await
721 .expect("cancel durable waits");
722 assert!(matches!(
723 tool_wait
724 .await
725 .expect("tool wait task")
726 .expect("tool resolution"),
727 Resolution::Cancelled
728 ));
729
730 assert!(matches!(
731 driver
732 .request_cancel(request(address_a.clone(), "request-a-duplicate"))
733 .await
734 .expect("duplicate a")
735 .outcome,
736 TurnCancelOutcome::AlreadyRequested(_)
737 ));
738 assert!(matches!(
739 driver
740 .request_cancel(request(address_b, "request-b"))
741 .await
742 .expect("cancel b")
743 .outcome,
744 TurnCancelOutcome::Requested(_)
745 ));
746 assert!(matches!(
747 driver
748 .request_cancel(request(address_future, "request-future"))
749 .await
750 .expect("cancel future")
751 .outcome,
752 TurnCancelOutcome::Requested(_)
753 ));
754 }
755
756 #[tokio::test]
757 async fn session_deletion_revokes_control_promises() {
758 let host = Arc::new(InlineEffectHost::default());
759 let driver = TurnWorkDriver::new(host.clone());
760 let address = address("revoke");
761 host.revoke_await_events_for_session(&address.session_id)
762 .await
763 .expect("revoke session");
764 assert!(matches!(
765 driver
766 .request_cancel(request(address, "request-after-delete"))
767 .await
768 .expect("revoked outcome")
769 .outcome,
770 TurnCancelOutcome::UnknownOrRevoked
771 ));
772 }
773
774 #[tokio::test]
775 async fn terminal_attachment_timeout_does_not_poison_later_publication() {
776 let host = Arc::new(InlineEffectHost::default());
777 let driver = TurnWorkDriver::new(host.clone());
778 let address = address("terminal-timeout");
779 let error = driver
780 .await_terminal_with_timeout(&address, Duration::from_millis(1))
781 .await
782 .expect_err("unpublished terminal must time out");
783 assert_eq!(error.code.as_str(), "turn_terminal_await_timeout");
784
785 let active = ActiveTurnControl::new(host.as_ref(), address.clone())
786 .await
787 .expect("active control after timed-out attach");
788 active
789 .settle_before_commit(host.as_ref(), false)
790 .await
791 .expect("seal after timed-out attach");
792 active
793 .publish_terminal(
794 host.as_ref(),
795 &TurnTerminal::Committed {
796 outcome: TurnOutcome::Finished(TurnFinish::AssistantMessage {
797 text: "done".to_string(),
798 }),
799 cancellation: None,
800 session_revision: None,
801 },
802 )
803 .await
804 .expect("publish after timed-out attach");
805 assert!(matches!(
806 driver.await_terminal(&address).await.expect("late attach"),
807 TurnTerminal::Committed {
808 outcome: TurnOutcome::Finished(_),
809 cancellation: None,
810 ..
811 }
812 ));
813 }
814
815 #[test]
816 fn local_cancel_origin_hint_preserves_first_origin() {
817 let hint = TurnCancelOriginHint::default();
818 hint.set(Some("shutdown".to_string()));
819 hint.set(Some("user".to_string()));
820
821 assert_eq!(hint.get().as_deref(), Some("shutdown"));
822 }
823
824 #[test]
825 fn local_cancel_origin_hint_preserves_explicit_absence() {
826 let hint = TurnCancelOriginHint::default();
827 hint.set(None);
828 hint.set(Some("user".to_string()));
829
830 assert_eq!(hint.get(), None);
831 }
832
833 #[test]
834 fn terminal_success_has_no_cancellation_evidence() {
835 let terminal = TurnTerminal::Committed {
836 outcome: TurnOutcome::Finished(TurnFinish::AssistantMessage {
837 text: "done".to_string(),
838 }),
839 cancellation: None,
840 session_revision: None,
841 };
842 let encoded = terminal_resolution(&terminal).expect("encode terminal");
843 assert!(matches!(encoded, Resolution::Ok(_)));
844 }
845}