Skip to main content

lash_core/runtime/
turn_control.rs

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, RuntimeError,
12};
13
14/// Stable routing identity for one foreground turn.
15///
16/// These identifiers select work; they are not authorization credentials.
17/// Hosts exposing turn control to untrusted callers must authenticate and
18/// authorize the request before calling Lash.
19#[derive(Clone, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
20pub struct TurnAddress {
21    pub session_id: String,
22    pub turn_id: String,
23}
24
25impl TurnAddress {
26    pub fn new(session_id: impl Into<String>, turn_id: impl Into<String>) -> Self {
27        Self {
28            session_id: session_id.into(),
29            turn_id: turn_id.into(),
30        }
31    }
32
33    fn scope(&self) -> ExecutionScope {
34        ExecutionScope::turn(&self.session_id, &self.turn_id)
35    }
36
37    fn validate(&self) -> Result<(), RuntimeError> {
38        self.scope().validate()
39    }
40}
41
42/// Shared origin hint for a process-local cancellation token.
43///
44/// The outer option records whether a local entry point supplied a hint; the
45/// inner option is the opaque host origin, which may intentionally be absent.
46/// It is not a durable cancellation request and must not be used as
47/// authorization.
48#[doc(hidden)]
49#[derive(Clone, Default)]
50pub struct TurnCancelOriginHint {
51    origin: Arc<Mutex<Option<Option<String>>>>,
52}
53
54impl TurnCancelOriginHint {
55    pub fn set(&self, origin: Option<String>) {
56        let mut hint = self.origin.lock().expect("turn cancel origin hint lock");
57        if hint.is_none() {
58            *hint = Some(origin);
59        }
60    }
61
62    pub(crate) fn get(&self) -> Option<String> {
63        self.origin
64            .lock()
65            .expect("turn cancel origin hint lock")
66            .clone()
67            .flatten()
68    }
69}
70
71#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
72pub struct TurnCancellationEvidence {
73    pub request_id: String,
74    /// Opaque host-domain data. Lash records and returns it unchanged.
75    #[serde(default, skip_serializing_if = "Option::is_none")]
76    pub origin: Option<String>,
77    #[serde(default, skip_serializing_if = "Option::is_none")]
78    pub reason: Option<String>,
79}
80
81#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
82pub struct TurnCancelRequest {
83    pub address: TurnAddress,
84    pub request_id: String,
85    /// Opaque host-domain data. Lash never interprets this value.
86    #[serde(default, skip_serializing_if = "Option::is_none")]
87    pub origin: Option<String>,
88    #[serde(default, skip_serializing_if = "Option::is_none")]
89    pub reason: Option<String>,
90}
91
92impl TurnCancelRequest {
93    pub fn new(
94        address: TurnAddress,
95        request_id: impl Into<String>,
96        origin: Option<String>,
97    ) -> Self {
98        Self {
99            address,
100            request_id: request_id.into(),
101            origin,
102            reason: None,
103        }
104    }
105
106    pub fn with_reason(mut self, reason: impl Into<String>) -> Self {
107        self.reason = Some(reason.into());
108        self
109    }
110
111    fn validate(&self) -> Result<(), RuntimeError> {
112        self.address.validate()?;
113        if self.request_id.trim().is_empty() {
114            return Err(RuntimeError::new(
115                "invalid_turn_cancel_request",
116                "turn cancellation requires a non-empty request id",
117            ));
118        }
119        Ok(())
120    }
121
122    fn evidence(&self) -> TurnCancellationEvidence {
123        TurnCancellationEvidence {
124            request_id: self.request_id.clone(),
125            origin: self.origin.clone(),
126            reason: self.reason.clone(),
127        }
128    }
129}
130
131#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
132#[serde(tag = "outcome", content = "cancellation", rename_all = "snake_case")]
133pub enum TurnCancelOutcome {
134    Requested(TurnCancellationEvidence),
135    AlreadyRequested(TurnCancellationEvidence),
136    CompletionWonRace,
137    UnknownOrRevoked,
138}
139
140/// Result of addressing one turn-cancellation gate.
141///
142/// [`durability_tier`](Self::durability_tier) describes the keyed-promise
143/// deployment that accepted the request. [`crate::DurabilityTier::Inline`]
144/// receipts are process-local: they do not prove that an owner in another OS
145/// process observed the request. Durable cross-process cancellation requires
146/// a [`crate::DurabilityTier::Durable`] effect-host deployment.
147#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
148pub struct TurnCancelReceipt {
149    pub durability_tier: crate::DurabilityTier,
150    pub outcome: TurnCancelOutcome,
151}
152
153#[derive(Clone, Debug, Serialize, Deserialize)]
154#[serde(tag = "status", rename_all = "snake_case")]
155pub enum TurnTerminal {
156    Committed {
157        outcome: TurnOutcome,
158        #[serde(default, skip_serializing_if = "Option::is_none")]
159        cancellation: Option<TurnCancellationEvidence>,
160        #[serde(default, skip_serializing_if = "Option::is_none")]
161        session_revision: Option<u64>,
162    },
163    Failed {
164        error: ErrorEnvelope,
165    },
166}
167
168/// Backend-specific terminal attachment for a foreground turn.
169#[async_trait::async_trait]
170pub trait TurnAttach: Send + Sync {
171    async fn await_terminal(&self, address: &TurnAddress) -> Result<TurnTerminal, RuntimeError>;
172}
173
174/// Cooperative, exact-turn control compiled onto Lash's keyed-promise seam.
175///
176/// `Requested` means the cancellation request won this driver's keyed-promise
177/// gate. On an inline effect host that promise is process-local, so the driver
178/// only reaches owners in the same OS process. A durable effect-host deployment
179/// is required for another process or replayed owner to observe the request.
180/// The returned [`TurnCancelReceipt`] exposes that tier so hosts can gate their
181/// UX rather than treating an inline receipt as cross-process proof.
182///
183/// Lash asks the running or replayed owner to unwind and commit a cancelled
184/// result; it cannot guarantee that detached tasks, subprocesses, or
185/// non-cooperative providers have stopped. Engine invocation cancellation
186/// remains a host-owned break-glass action and is never proof of a Lash
187/// `Cancelled` result.
188///
189/// Session and turn ids are routing identity, not authorization. Hosts must
190/// enforce authorization before exposing this driver across a trust boundary.
191#[derive(Clone)]
192pub struct TurnWorkDriver {
193    effect_host: Arc<dyn EffectHost>,
194    attach: Option<Arc<dyn TurnAttach>>,
195}
196
197impl TurnWorkDriver {
198    pub fn new(effect_host: Arc<dyn EffectHost>) -> Self {
199        Self {
200            effect_host,
201            attach: None,
202        }
203    }
204
205    pub fn with_attach(mut self, attach: Arc<dyn TurnAttach>) -> Self {
206        self.attach = Some(attach);
207        self
208    }
209
210    pub fn effect_host(&self) -> Arc<dyn EffectHost> {
211        Arc::clone(&self.effect_host)
212    }
213
214    pub async fn request_cancel(
215        &self,
216        request: TurnCancelRequest,
217    ) -> Result<TurnCancelReceipt, RuntimeError> {
218        request.validate()?;
219        let durability_tier = self.effect_host.durability_tier();
220        let key = cancel_gate_key(self.effect_host.as_ref(), &request.address).await?;
221        let evidence = request.evidence();
222        let resolution = gate_resolution(TurnGateTerminal::CancelRequested(evidence.clone()))?;
223        let outcome = match self
224            .effect_host
225            .resolve_await_event(&key, resolution)
226            .await?
227        {
228            ResolveOutcome::Accepted => Ok(TurnCancelOutcome::Requested(evidence)),
229            ResolveOutcome::AlreadyResolved { terminal } => match decode_gate(terminal)? {
230                TurnGateTerminal::CancelRequested(existing) => {
231                    Ok(TurnCancelOutcome::AlreadyRequested(existing))
232                }
233                TurnGateTerminal::CompletionSealed => Ok(TurnCancelOutcome::CompletionWonRace),
234            },
235            ResolveOutcome::UnknownOrRevoked => Ok(TurnCancelOutcome::UnknownOrRevoked),
236        }?;
237        Ok(TurnCancelReceipt {
238            durability_tier,
239            outcome,
240        })
241    }
242
243    pub async fn await_terminal(
244        &self,
245        address: &TurnAddress,
246    ) -> Result<TurnTerminal, RuntimeError> {
247        address.validate()?;
248        if let Some(attach) = self.attach.as_ref() {
249            return attach.await_terminal(address).await;
250        }
251        let key = terminal_key(self.effect_host.as_ref(), address).await?;
252        let resolution = self
253            .effect_host
254            .await_await_event(&key, CancellationToken::new(), None)
255            .await?;
256        decode_terminal(address, resolution)
257    }
258
259    /// Await a terminal publication for at most `timeout`.
260    ///
261    /// Timing out only stops this caller's attachment. It never resolves or
262    /// poisons the turn's first-writer-wins keyed promises.
263    pub async fn await_terminal_with_timeout(
264        &self,
265        address: &TurnAddress,
266        timeout: Duration,
267    ) -> Result<TurnTerminal, RuntimeError> {
268        tokio::time::timeout(timeout, self.await_terminal(address))
269            .await
270            .map_err(|_| {
271                RuntimeError::new(
272                    "turn_terminal_await_timeout",
273                    format!(
274                        "timed out awaiting terminal for turn `{}` in session `{}` after {} ms",
275                        address.turn_id,
276                        address.session_id,
277                        timeout.as_millis()
278                    ),
279                )
280            })?
281    }
282}
283
284#[derive(Clone, Debug, Serialize, Deserialize)]
285#[serde(tag = "state", content = "cancellation", rename_all = "snake_case")]
286enum TurnGateTerminal {
287    CancelRequested(TurnCancellationEvidence),
288    CompletionSealed,
289}
290
291fn gate_resolution(value: TurnGateTerminal) -> Result<Resolution, RuntimeError> {
292    serde_json::to_value(value)
293        .map(Resolution::Ok)
294        .map_err(|err| RuntimeError::new("turn_cancel_gate_encode", err.to_string()))
295}
296
297fn decode_gate(resolution: Resolution) -> Result<TurnGateTerminal, RuntimeError> {
298    match resolution {
299        Resolution::Ok(value) => serde_json::from_value(value)
300            .map_err(|err| RuntimeError::new("turn_cancel_gate_decode", err.to_string())),
301        other => Err(RuntimeError::new(
302            "turn_cancel_gate_invalid_terminal",
303            format!("turn cancellation gate resolved with {other:?}"),
304        )),
305    }
306}
307
308fn terminal_resolution(value: &TurnTerminal) -> Result<Resolution, RuntimeError> {
309    serde_json::to_value(value)
310        .map(Resolution::Ok)
311        .map_err(|err| RuntimeError::new("turn_terminal_encode", err.to_string()))
312}
313
314fn decode_terminal(
315    address: &TurnAddress,
316    resolution: Resolution,
317) -> Result<TurnTerminal, RuntimeError> {
318    match resolution {
319        Resolution::Ok(value) => serde_json::from_value(value).map_err(|err| {
320            RuntimeError::new(
321                "turn_terminal_decode",
322                format!(
323                    "invalid terminal result for turn `{}` in session `{}`: {err}",
324                    address.turn_id, address.session_id
325                ),
326            )
327        }),
328        other => Err(RuntimeError::new(
329            "turn_terminal_invalid_resolution",
330            format!(
331                "terminal result for turn `{}` in session `{}` resolved with {other:?}",
332                address.turn_id, address.session_id
333            ),
334        )),
335    }
336}
337
338async fn cancel_gate_key(
339    resolver: &dyn AwaitEventResolver,
340    address: &TurnAddress,
341) -> Result<AwaitEventKey, RuntimeError> {
342    resolver
343        .await_event_key(&address.scope(), AwaitEventWaitIdentity::TurnCancelGate)
344        .await
345}
346
347async fn terminal_key(
348    resolver: &dyn AwaitEventResolver,
349    address: &TurnAddress,
350) -> Result<AwaitEventKey, RuntimeError> {
351    resolver
352        .await_event_key(&address.scope(), AwaitEventWaitIdentity::TurnTerminal)
353        .await
354}
355
356/// Per-execution bridge between the durable gate and the turn's internal
357/// cancellation token.
358pub(crate) struct ActiveTurnControl {
359    address: TurnAddress,
360    cancel_key: AwaitEventKey,
361    terminal_key: AwaitEventKey,
362    evidence: Mutex<Option<TurnCancellationEvidence>>,
363    local_cancel_origin: TurnCancelOriginHint,
364}
365
366impl ActiveTurnControl {
367    pub(crate) async fn new(
368        resolver: &dyn AwaitEventResolver,
369        address: TurnAddress,
370    ) -> Result<Self, RuntimeError> {
371        address.validate()?;
372        Ok(Self {
373            cancel_key: cancel_gate_key(resolver, &address).await?,
374            terminal_key: terminal_key(resolver, &address).await?,
375            address,
376            evidence: Mutex::new(None),
377            local_cancel_origin: TurnCancelOriginHint::default(),
378        })
379    }
380
381    pub(crate) fn with_local_cancel_origin(mut self, origin: TurnCancelOriginHint) -> Self {
382        self.local_cancel_origin = origin;
383        self
384    }
385
386    pub(crate) async fn await_cancel(
387        &self,
388        resolver: &dyn AwaitEventResolver,
389        stop_wait: CancellationToken,
390    ) -> Result<Option<TurnCancellationEvidence>, RuntimeError> {
391        let resolution = resolver
392            .await_await_event(&self.cancel_key, stop_wait, None)
393            .await?;
394        match decode_gate(resolution)? {
395            TurnGateTerminal::CancelRequested(evidence) => {
396                self.remember(evidence.clone());
397                Ok(Some(evidence))
398            }
399            TurnGateTerminal::CompletionSealed => Ok(None),
400        }
401    }
402
403    pub(crate) async fn observe_pending_cancel(
404        &self,
405        resolver: &dyn AwaitEventResolver,
406    ) -> Result<Option<TurnCancellationEvidence>, RuntimeError> {
407        let Some(resolution) = resolver.peek_await_event(&self.cancel_key).await? else {
408            return Ok(None);
409        };
410        match decode_gate(resolution)? {
411            TurnGateTerminal::CancelRequested(evidence) => {
412                self.remember(evidence.clone());
413                Ok(Some(evidence))
414            }
415            TurnGateTerminal::CompletionSealed => Ok(None),
416        }
417    }
418
419    pub(crate) async fn settle_before_commit(
420        &self,
421        resolver: &dyn AwaitEventResolver,
422        locally_cancelled: bool,
423    ) -> Result<Option<TurnCancellationEvidence>, RuntimeError> {
424        if let Some(evidence) = self.evidence() {
425            return Ok(Some(evidence));
426        }
427        let proposed = if locally_cancelled {
428            TurnGateTerminal::CancelRequested(self.internal_evidence())
429        } else {
430            TurnGateTerminal::CompletionSealed
431        };
432        let outcome = resolver
433            .resolve_await_event(&self.cancel_key, gate_resolution(proposed.clone())?)
434            .await?;
435        let terminal = match outcome {
436            ResolveOutcome::Accepted => proposed,
437            ResolveOutcome::AlreadyResolved { terminal } => decode_gate(terminal)?,
438            ResolveOutcome::UnknownOrRevoked => {
439                return Err(RuntimeError::new(
440                    "turn_control_unknown_or_revoked",
441                    format!(
442                        "turn `{}` in session `{}` was revoked before final commit",
443                        self.address.turn_id, self.address.session_id
444                    ),
445                ));
446            }
447        };
448        match terminal {
449            TurnGateTerminal::CancelRequested(evidence) => {
450                self.remember(evidence.clone());
451                Ok(Some(evidence))
452            }
453            TurnGateTerminal::CompletionSealed => Ok(None),
454        }
455    }
456
457    pub(crate) async fn publish_terminal(
458        &self,
459        resolver: &dyn AwaitEventResolver,
460        terminal: &TurnTerminal,
461    ) -> Result<(), RuntimeError> {
462        match resolver
463            .resolve_await_event(&self.terminal_key, terminal_resolution(terminal)?)
464            .await?
465        {
466            ResolveOutcome::Accepted | ResolveOutcome::AlreadyResolved { .. } => Ok(()),
467            ResolveOutcome::UnknownOrRevoked => Err(RuntimeError::new(
468                "turn_terminal_unknown_or_revoked",
469                format!(
470                    "terminal promise for turn `{}` in session `{}` was revoked",
471                    self.address.turn_id, self.address.session_id
472                ),
473            )),
474        }
475    }
476
477    pub(crate) fn evidence(&self) -> Option<TurnCancellationEvidence> {
478        self.evidence
479            .lock()
480            .expect("turn cancellation evidence lock")
481            .clone()
482    }
483
484    fn remember(&self, evidence: TurnCancellationEvidence) {
485        *self
486            .evidence
487            .lock()
488            .expect("turn cancellation evidence lock") = Some(evidence);
489    }
490
491    fn internal_evidence(&self) -> TurnCancellationEvidence {
492        TurnCancellationEvidence {
493            request_id: format!("internal:{}", self.address.turn_id),
494            origin: self.local_cancel_origin.get(),
495            reason: None,
496        }
497    }
498}
499
500#[cfg(test)]
501mod tests {
502    use super::*;
503    use crate::{InlineEffectHost, TurnFinish, TurnStop};
504
505    fn address(label: &str) -> TurnAddress {
506        TurnAddress::new(
507            format!("turn-control-{label}-{}", uuid::Uuid::new_v4()),
508            "turn-a",
509        )
510    }
511
512    fn request(address: TurnAddress, request_id: &str) -> TurnCancelRequest {
513        TurnCancelRequest::new(address, request_id, Some("user".to_string()))
514            .with_reason("stop button")
515    }
516
517    #[tokio::test]
518    async fn cancel_before_start_duplicate_and_terminal_attach() {
519        let host = Arc::new(InlineEffectHost::default());
520        let driver = TurnWorkDriver::new(host.clone());
521        let address = address("before-start");
522
523        let first = driver
524            .request_cancel(request(address.clone(), "request-1"))
525            .await
526            .expect("request cancellation");
527        assert_eq!(first.durability_tier, crate::DurabilityTier::Inline);
528        let evidence = match first.outcome {
529            TurnCancelOutcome::Requested(evidence) => evidence,
530            other => panic!("expected requested, got {other:?}"),
531        };
532        assert_eq!(evidence.request_id, "request-1");
533
534        let duplicate = driver
535            .request_cancel(request(address.clone(), "request-2"))
536            .await
537            .expect("duplicate cancellation");
538        assert!(matches!(
539            duplicate.outcome,
540            TurnCancelOutcome::AlreadyRequested(TurnCancellationEvidence { ref request_id, .. })
541                if request_id == "request-1"
542        ));
543
544        let active = ActiveTurnControl::new(host.as_ref(), address.clone())
545            .await
546            .expect("active control");
547        let observed = active
548            .settle_before_commit(host.as_ref(), false)
549            .await
550            .expect("settle")
551            .expect("cancellation won");
552        assert_eq!(observed, evidence);
553        let terminal = TurnTerminal::Committed {
554            outcome: TurnOutcome::Stopped(TurnStop::Cancelled),
555            cancellation: Some(observed),
556            session_revision: Some(7),
557        };
558        active
559            .publish_terminal(host.as_ref(), &terminal)
560            .await
561            .expect("publish terminal");
562        let attached = driver
563            .await_terminal(&address)
564            .await
565            .expect("attach terminal");
566        assert!(matches!(
567            attached,
568            TurnTerminal::Committed {
569                outcome: TurnOutcome::Stopped(TurnStop::Cancelled),
570                cancellation: Some(_),
571                session_revision: Some(7),
572            }
573        ));
574    }
575
576    #[tokio::test]
577    async fn concurrent_completion_seal_vs_cancel_is_first_writer_wins() {
578        let host = Arc::new(InlineEffectHost::default());
579        let driver = TurnWorkDriver::new(host.clone());
580        let address = address("race");
581        let active = ActiveTurnControl::new(host.as_ref(), address.clone())
582            .await
583            .expect("active control");
584
585        let (seal, cancel) = tokio::join!(
586            active.settle_before_commit(host.as_ref(), false),
587            driver.request_cancel(request(address, "race-request")),
588        );
589        match (seal.expect("seal"), cancel.expect("cancel").outcome) {
590            (None, TurnCancelOutcome::CompletionWonRace) => {}
591            (Some(evidence), TurnCancelOutcome::Requested(requested)) => {
592                assert_eq!(evidence, requested);
593            }
594            other => panic!("inconsistent gate race result: {other:?}"),
595        }
596    }
597
598    #[tokio::test]
599    async fn recovered_owner_observes_pending_cancel_after_control_recreation() {
600        let host = Arc::new(InlineEffectHost::default());
601        let driver = TurnWorkDriver::new(host.clone());
602        let address = address("replay");
603        let requested = driver
604            .request_cancel(request(address.clone(), "request-before-replay"))
605            .await
606            .expect("request cancellation");
607        let expected = match requested.outcome {
608            TurnCancelOutcome::Requested(evidence) => evidence,
609            other => panic!("expected requested, got {other:?}"),
610        };
611
612        let recovered = ActiveTurnControl::new(host.as_ref(), address)
613            .await
614            .expect("recreate active control under the recovered owner");
615        let observed = recovered
616            .observe_pending_cancel(host.as_ref())
617            .await
618            .expect("read recovered turn start gate")
619            .expect("pending cancellation is visible before recovered effects");
620        assert_eq!(observed, expected);
621        let settled = recovered
622            .settle_before_commit(host.as_ref(), false)
623            .await
624            .expect("settle recovered turn")
625            .expect("pending cancellation survives owner loss");
626        assert_eq!(settled, expected);
627    }
628
629    #[tokio::test]
630    async fn turn_control_is_exact_scope_and_excluded_from_wait_cancel_sweep() {
631        let host = Arc::new(InlineEffectHost::default());
632        let driver = TurnWorkDriver::new(host.clone());
633        let address_a = address("scope");
634        let address_b = TurnAddress::new(&address_a.session_id, "turn-b");
635        let address_future = TurnAddress::new(&address_a.session_id, "turn-future");
636
637        driver
638            .request_cancel(request(address_a.clone(), "request-a"))
639            .await
640            .expect("cancel a");
641
642        let tool_key = host
643            .await_event_key(
644                &ExecutionScope::turn(&address_a.session_id, "tool-turn"),
645                AwaitEventWaitIdentity::tool_completion("tool-call"),
646            )
647            .await
648            .expect("tool key");
649        let tool_host = host.clone();
650        let tool_wait = tokio::spawn(async move {
651            tool_host
652                .await_await_event(&tool_key, CancellationToken::new(), None)
653                .await
654        });
655        tokio::task::yield_now().await;
656        host.cancel_await_events_for_session(&address_a.session_id)
657            .await
658            .expect("cancel durable waits");
659        assert!(matches!(
660            tool_wait
661                .await
662                .expect("tool wait task")
663                .expect("tool resolution"),
664            Resolution::Cancelled
665        ));
666
667        assert!(matches!(
668            driver
669                .request_cancel(request(address_a.clone(), "request-a-duplicate"))
670                .await
671                .expect("duplicate a")
672                .outcome,
673            TurnCancelOutcome::AlreadyRequested(_)
674        ));
675        assert!(matches!(
676            driver
677                .request_cancel(request(address_b, "request-b"))
678                .await
679                .expect("cancel b")
680                .outcome,
681            TurnCancelOutcome::Requested(_)
682        ));
683        assert!(matches!(
684            driver
685                .request_cancel(request(address_future, "request-future"))
686                .await
687                .expect("cancel future")
688                .outcome,
689            TurnCancelOutcome::Requested(_)
690        ));
691    }
692
693    #[tokio::test]
694    async fn session_deletion_revokes_control_promises() {
695        let host = Arc::new(InlineEffectHost::default());
696        let driver = TurnWorkDriver::new(host.clone());
697        let address = address("revoke");
698        host.revoke_await_events_for_session(&address.session_id)
699            .await
700            .expect("revoke session");
701        assert!(matches!(
702            driver
703                .request_cancel(request(address, "request-after-delete"))
704                .await
705                .expect("revoked outcome")
706                .outcome,
707            TurnCancelOutcome::UnknownOrRevoked
708        ));
709    }
710
711    #[tokio::test]
712    async fn terminal_attachment_timeout_does_not_poison_later_publication() {
713        let host = Arc::new(InlineEffectHost::default());
714        let driver = TurnWorkDriver::new(host.clone());
715        let address = address("terminal-timeout");
716        let error = driver
717            .await_terminal_with_timeout(&address, Duration::from_millis(1))
718            .await
719            .expect_err("unpublished terminal must time out");
720        assert_eq!(error.code.as_str(), "turn_terminal_await_timeout");
721
722        let active = ActiveTurnControl::new(host.as_ref(), address.clone())
723            .await
724            .expect("active control after timed-out attach");
725        active
726            .settle_before_commit(host.as_ref(), false)
727            .await
728            .expect("seal after timed-out attach");
729        active
730            .publish_terminal(
731                host.as_ref(),
732                &TurnTerminal::Committed {
733                    outcome: TurnOutcome::Finished(TurnFinish::AssistantMessage {
734                        text: "done".to_string(),
735                    }),
736                    cancellation: None,
737                    session_revision: None,
738                },
739            )
740            .await
741            .expect("publish after timed-out attach");
742        assert!(matches!(
743            driver.await_terminal(&address).await.expect("late attach"),
744            TurnTerminal::Committed {
745                outcome: TurnOutcome::Finished(_),
746                cancellation: None,
747                ..
748            }
749        ));
750    }
751
752    #[test]
753    fn local_cancel_origin_hint_preserves_first_origin() {
754        let hint = TurnCancelOriginHint::default();
755        hint.set(Some("shutdown".to_string()));
756        hint.set(Some("user".to_string()));
757
758        assert_eq!(hint.get().as_deref(), Some("shutdown"));
759    }
760
761    #[test]
762    fn local_cancel_origin_hint_preserves_explicit_absence() {
763        let hint = TurnCancelOriginHint::default();
764        hint.set(None);
765        hint.set(Some("user".to_string()));
766
767        assert_eq!(hint.get(), None);
768    }
769
770    #[test]
771    fn terminal_success_has_no_cancellation_evidence() {
772        let terminal = TurnTerminal::Committed {
773            outcome: TurnOutcome::Finished(TurnFinish::AssistantMessage {
774                text: "done".to_string(),
775            }),
776            cancellation: None,
777            session_revision: None,
778        };
779        let encoded = terminal_resolution(&terminal).expect("encode terminal");
780        assert!(matches!(encoded, Resolution::Ok(_)));
781    }
782}