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, 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/// Stable routing identity for one foreground turn.
32///
33/// These identifiers select work; they are not authorization credentials.
34/// Hosts exposing turn control to untrusted callers must authenticate and
35/// authorize the request before calling Lash.
36#[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/// Shared origin hint for a process-local cancellation token.
60///
61/// The outer option records whether a local entry point supplied a hint; the
62/// inner option is the opaque host origin, which may intentionally be absent.
63/// It is not a durable cancellation request and must not be used as
64/// authorization.
65#[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    /// Opaque host-domain data. Lash records and returns it unchanged.
92    #[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    /// Opaque host-domain data. Lash never interprets this value.
103    #[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/// Result of addressing one turn-cancellation gate.
158///
159/// [`durability_tier`](Self::durability_tier) describes the keyed-promise
160/// deployment that accepted the request. [`crate::DurabilityTier::Inline`]
161/// receipts are process-local: they do not prove that an owner in another OS
162/// process observed the request. Durable cross-process cancellation requires
163/// a [`crate::DurabilityTier::Durable`] effect-host deployment that passes the
164/// full cold-instance AwaitEvent conformance across independent processes. A
165/// durable journal paired with process-local waits must issue `Inline`
166/// receipts.
167#[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/// Backend-specific terminal attachment for a foreground turn.
189#[async_trait::async_trait]
190pub trait TurnAttach: Send + Sync {
191    async fn await_terminal(&self, address: &TurnAddress) -> Result<TurnTerminal, RuntimeError>;
192}
193
194/// Cooperative, exact-turn control compiled onto Lash's keyed-promise seam.
195///
196/// `Requested` means the cancellation request won this driver's keyed-promise
197/// gate. On an inline effect host that promise is process-local, so the driver
198/// only reaches owners in the same OS process. A durable effect-host deployment
199/// is required for another process or replayed owner to observe the request.
200/// The returned [`TurnCancelReceipt`] exposes that tier so hosts can gate their
201/// UX rather than treating an inline receipt as cross-process proof.
202///
203/// Lash asks the running or replayed owner to unwind and commit a cancelled
204/// result; it cannot guarantee that detached tasks, subprocesses, or
205/// non-cooperative providers have stopped. Engine invocation cancellation
206/// remains a host-owned break-glass action and is never proof of a Lash
207/// `Cancelled` result.
208///
209/// Session and turn ids are routing identity, not authorization. Hosts must
210/// enforce authorization before exposing this driver across a trust boundary.
211#[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    /// Await a terminal publication for at most `timeout`.
289    ///
290    /// Timing out only stops this caller's attachment. It never resolves or
291    /// poisons the turn's first-writer-wins keyed promises.
292    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
385/// Per-execution bridge between the durable gate and the turn's internal
386/// cancellation token.
387pub(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}