Skip to main content

meerkat_runtime/
lib.rs

1//! meerkat-runtime — v9 runtime control-plane for Meerkat agent lifecycle.
2//!
3//! This crate implements the runtime/control-plane layer of the v9 Canonical
4//! Lifecycle specification. It sits between surfaces (CLI, RPC, REST, MCP)
5//! and core (`meerkat-core`), managing:
6//!
7//! - Input acceptance, validation, and queueing
8//! - InputState lifecycle tracking
9//! - Policy resolution (what to do with each input)
10//! - Runtime state machine (Initializing ↔ Idle ↔ Attached ↔ Running ↔ Retired/Stopped/Destroyed)
11//! - Retire/recycle/reset lifecycle operations
12//! - RuntimeEvent observability
13//!
14//! Core-facing types (RunPrimitive, RunEvent, CoreExecutor, etc.) live in
15//! `meerkat-core::lifecycle`. This crate contains everything else.
16
17#![cfg_attr(
18    test,
19    allow(
20        dead_code,
21        unused_imports,
22        clippy::expect_used,
23        clippy::large_futures,
24        clippy::needless_borrow,
25        clippy::panic,
26        clippy::redundant_closure_for_method_calls,
27        clippy::redundant_clone,
28        clippy::type_complexity,
29        clippy::unnecessary_to_owned,
30        clippy::unwrap_used
31    )
32)]
33
34#[cfg(target_arch = "wasm32")]
35pub mod tokio {
36    pub use tokio_with_wasm::alias::*;
37}
38
39#[cfg(not(target_arch = "wasm32"))]
40pub use ::tokio;
41
42pub mod accept;
43pub mod auth_machine;
44pub mod coalescing;
45pub mod comms_bridge;
46pub mod comms_drain;
47pub mod comms_trust_reconcile;
48pub mod completion;
49pub mod composition;
50pub(crate) mod control_plane;
51pub mod delivery_inbox;
52pub mod driver;
53pub(crate) mod effect;
54#[doc(hidden)]
55pub mod generated;
56pub mod handles;
57pub mod identifiers;
58pub mod ingress_types;
59pub mod input;
60pub mod input_ledger;
61pub mod input_scope;
62pub mod input_state;
63pub mod interrupt_public_result;
64pub mod meerkat_machine;
65pub(crate) mod meerkat_machine_types;
66pub mod member_live;
67pub mod member_observation;
68pub mod mob_adapter;
69pub mod mob_operator_authority;
70pub mod ops_lifecycle;
71pub(crate) mod panic_boundary;
72pub mod peer_handling_mode;
73pub mod policy;
74pub mod policy_table;
75#[allow(unused_imports)]
76#[path = "generated/protocol_auth_lease_lifecycle_publication.rs"]
77pub mod protocol_auth_lease_lifecycle_publication;
78#[allow(unused_imports)]
79#[path = "generated/protocol_auth_release_oauth_flow_drain.rs"]
80pub mod protocol_auth_release_oauth_flow_drain;
81#[allow(unused_imports)]
82#[path = "generated/protocol_comms_trust_reconcile.rs"]
83pub mod protocol_comms_trust_reconcile;
84#[allow(unused_imports)]
85#[path = "generated/protocol_supervisor_trust_publish.rs"]
86pub mod protocol_supervisor_trust_publish;
87#[allow(unused_imports)]
88#[path = "generated/protocol_supervisor_trust_revoke.rs"]
89pub mod protocol_supervisor_trust_revoke;
90pub(crate) mod queue;
91pub mod recovery;
92pub mod runtime_event;
93pub(crate) mod runtime_loop;
94pub mod runtime_state;
95pub mod service_ext;
96pub(crate) mod silent_intent;
97pub mod stack_relief;
98pub mod store;
99pub mod terminal_status;
100pub mod traits;
101
102use meerkat_core::lifecycle::run_primitive::RuntimeTurnMetadata as RuntimeStampedTurnMetadata;
103use std::any::Any;
104use std::sync::Arc;
105
106pub(crate) struct SessionRuntimeBindingsAuthority {
107    pub(crate) session_id: meerkat_core::SessionId,
108    pub(crate) epoch_id: meerkat_core::RuntimeEpochId,
109    pub(crate) dsl_authority: Arc<std::sync::Mutex<meerkat_machine::dsl::MeerkatMachineAuthority>>,
110    pub(crate) teardown_gate: Arc<handles::HandleTeardownGate>,
111    pub(crate) materialization_claim_id: Option<uuid::Uuid>,
112    pub(crate) materialization_claim_state:
113        Arc<std::sync::Mutex<RuntimeActorMaterializationClaimState>>,
114    /// Compatibility capability for cloneable `prepare_bindings()` results.
115    /// It is minted only while the registration is unattached and does not
116    /// itself reserve the exact materialization claim. `begin_*` atomically
117    /// converts it into a one-shot claim if the window is still vacant.
118    pub(crate) legacy_actor_materialization_generation: Option<u64>,
119    pub(crate) release_materialization_claim_on_drop: bool,
120}
121
122#[derive(Debug, Clone, Copy, PartialEq, Eq)]
123pub(crate) enum RuntimeActorMaterializationClaimPhase {
124    Vacant,
125    Prepared,
126    Staged,
127    ActorCreating,
128    ActorMaterializedPendingCommit,
129    RetainedActor,
130    Aborting,
131}
132
133pub(crate) struct RuntimeActorMaterializationClaimState {
134    pub(crate) current: Option<uuid::Uuid>,
135    pub(crate) phase: RuntimeActorMaterializationClaimPhase,
136    /// True only when this registry entry was inserted by a unique,
137    /// rollback-owning actor materialization that has not yet attached an
138    /// executor. A successor unique prepare inherits this exact rollback
139    /// authority; cloneable compatibility bindings and pre-existing committed
140    /// registrations never gain it.
141    pub(crate) rollback_registration_available: bool,
142    /// Monotonic incarnation of the cloneable compatibility binding window.
143    /// Executor attachment increments this under the same mutex as the claim
144    /// phase, permanently fencing bindings that escaped an older window.
145    pub(crate) legacy_capability_generation: u64,
146    pub(crate) changed: Arc<crate::tokio::sync::Notify>,
147}
148
149impl RuntimeActorMaterializationClaimState {
150    pub(crate) fn new(rollback_registration_available: bool) -> Self {
151        Self {
152            current: None,
153            phase: RuntimeActorMaterializationClaimPhase::Vacant,
154            rollback_registration_available,
155            legacy_capability_generation: 0,
156            changed: Arc::new(crate::tokio::sync::Notify::new()),
157        }
158    }
159
160    pub(crate) fn exact_claim_is(
161        &self,
162        claim_id: uuid::Uuid,
163        phases: &[RuntimeActorMaterializationClaimPhase],
164    ) -> bool {
165        self.current == Some(claim_id) && phases.contains(&self.phase)
166    }
167}
168
169impl Drop for SessionRuntimeBindingsAuthority {
170    fn drop(&mut self) {
171        if !self.release_materialization_claim_on_drop {
172            return;
173        }
174        let Some(claim_id) = self.materialization_claim_id else {
175            return;
176        };
177        let changed = {
178            let mut state = self
179                .materialization_claim_state
180                .lock()
181                .unwrap_or_else(std::sync::PoisonError::into_inner);
182            if !state.exact_claim_is(
183                claim_id,
184                &[
185                    RuntimeActorMaterializationClaimPhase::Prepared,
186                    RuntimeActorMaterializationClaimPhase::Staged,
187                ],
188            ) {
189                return;
190            }
191            state.current = None;
192            state.phase = RuntimeActorMaterializationClaimPhase::Vacant;
193            Arc::clone(&state.changed)
194        };
195        changed.notify_waiters();
196    }
197}
198
199// Constructor mirrors the opaque session-binding authority payload exactly;
200// keeping each carrier explicit prevents partial or reordered minting.
201#[allow(clippy::too_many_arguments)]
202pub(crate) fn session_runtime_bindings_authority(
203    session_id: meerkat_core::SessionId,
204    epoch_id: meerkat_core::RuntimeEpochId,
205    dsl_authority: Arc<std::sync::Mutex<meerkat_machine::dsl::MeerkatMachineAuthority>>,
206    teardown_gate: Arc<handles::HandleTeardownGate>,
207    materialization_claim_id: Option<uuid::Uuid>,
208    materialization_claim_state: Arc<std::sync::Mutex<RuntimeActorMaterializationClaimState>>,
209    legacy_actor_materialization_generation: Option<u64>,
210    release_materialization_claim_on_drop: bool,
211) -> Arc<dyn Any + Send + Sync> {
212    Arc::new(SessionRuntimeBindingsAuthority {
213        session_id,
214        epoch_id,
215        dsl_authority,
216        teardown_gate,
217        materialization_claim_id,
218        materialization_claim_state,
219        legacy_actor_materialization_generation,
220        release_materialization_claim_on_drop,
221    })
222}
223
224#[allow(clippy::too_many_arguments)]
225pub(crate) fn local_session_runtime_bindings_authority(
226    session_id: meerkat_core::SessionId,
227    epoch_id: meerkat_core::RuntimeEpochId,
228    dsl_authority: Arc<std::sync::Mutex<meerkat_machine::dsl::MeerkatMachineAuthority>>,
229    teardown_gate: Arc<handles::HandleTeardownGate>,
230    materialization_claim_id: Option<uuid::Uuid>,
231    materialization_claim_state: Arc<std::sync::Mutex<RuntimeActorMaterializationClaimState>>,
232    legacy_actor_materialization_generation: Option<u64>,
233    release_materialization_claim_on_drop: bool,
234) -> Arc<dyn Any + Send + Sync> {
235    session_runtime_bindings_authority(
236        session_id,
237        epoch_id,
238        dsl_authority,
239        teardown_gate,
240        materialization_claim_id,
241        materialization_claim_state,
242        legacy_actor_materialization_generation,
243        release_materialization_claim_on_drop,
244    )
245}
246
247pub fn session_runtime_bindings_have_machine_authority(
248    bindings: &meerkat_core::SessionRuntimeBindings,
249) -> bool {
250    bindings
251        .__runtime_authority()
252        .is::<SessionRuntimeBindingsAuthority>()
253}
254
255#[derive(Debug, thiserror::Error)]
256pub enum RuntimeActorMaterializationError {
257    #[error("invalid runtime binding materialization authority: {0}")]
258    InvalidAuthority(String),
259    #[error("runtime binding registration no longer admits actor materialization")]
260    RegistrationClosed,
261}
262
263/// Exclusive actor-create permit for one exact prepared runtime binding.
264///
265/// The persistent session service acquires this immediately before it starts
266/// building the live actor. Dropping an uncommitted permit restores the prior
267/// prepared/staged phase; committing it records that the actor exists but is
268/// still owned by the surrounding materialization transaction until executor
269/// attachment or an explicit retained-actor commit.
270pub struct RuntimeActorMaterializationPermit {
271    claim_id: uuid::Uuid,
272    claim_state: Arc<std::sync::Mutex<RuntimeActorMaterializationClaimState>>,
273    session_id: meerkat_core::SessionId,
274    epoch_id: meerkat_core::RuntimeEpochId,
275    dsl_authority: Arc<std::sync::Mutex<meerkat_machine::dsl::MeerkatMachineAuthority>>,
276    teardown_gate: Arc<handles::HandleTeardownGate>,
277    previous_phase: RuntimeActorMaterializationClaimPhase,
278    transactional: bool,
279    phase_policy: RuntimeActorMaterializationPhasePolicy,
280    _mutation_guard: Option<crate::tokio::sync::OwnedMutexGuard<()>>,
281    committed: bool,
282}
283
284#[derive(Debug, Clone, Copy, PartialEq, Eq)]
285enum RuntimeActorMaterializationPhasePolicy {
286    RejectRetired,
287    RequireArchivedRevivalBoundary,
288}
289
290impl RuntimeActorMaterializationPermit {
291    pub fn commit(mut self) -> Result<(), RuntimeActorMaterializationError> {
292        let generated_authority = self
293            .dsl_authority
294            .lock()
295            .unwrap_or_else(std::sync::PoisonError::into_inner);
296        validate_materialization_registration_authority(
297            &self.session_id,
298            &self.epoch_id,
299            &self.teardown_gate,
300            &generated_authority,
301            self.phase_policy,
302        )?;
303        let mut state = self
304            .claim_state
305            .lock()
306            .unwrap_or_else(std::sync::PoisonError::into_inner);
307        if !state.exact_claim_is(
308            self.claim_id,
309            &[RuntimeActorMaterializationClaimPhase::ActorCreating],
310        ) {
311            return Err(RuntimeActorMaterializationError::RegistrationClosed);
312        }
313        let changed = if self.transactional {
314            state.phase = RuntimeActorMaterializationClaimPhase::ActorMaterializedPendingCommit;
315            None
316        } else {
317            state.current = None;
318            state.phase = RuntimeActorMaterializationClaimPhase::RetainedActor;
319            state.rollback_registration_available = false;
320            Some(Arc::clone(&state.changed))
321        };
322        self.committed = true;
323        drop(state);
324        if let Some(changed) = changed {
325            changed.notify_waiters();
326        }
327        Ok(())
328    }
329}
330
331impl Drop for RuntimeActorMaterializationPermit {
332    fn drop(&mut self) {
333        if self.committed {
334            return;
335        }
336        let mut state = self
337            .claim_state
338            .lock()
339            .unwrap_or_else(std::sync::PoisonError::into_inner);
340        if state.exact_claim_is(
341            self.claim_id,
342            &[RuntimeActorMaterializationClaimPhase::ActorCreating],
343        ) {
344            if self.transactional {
345                state.phase = self.previous_phase;
346            } else {
347                state.current = None;
348                state.phase = RuntimeActorMaterializationClaimPhase::Vacant;
349                Arc::clone(&state.changed).notify_waiters();
350            }
351        }
352    }
353}
354
355fn validated_session_runtime_bindings_authority(
356    bindings: &meerkat_core::SessionRuntimeBindings,
357) -> Result<&SessionRuntimeBindingsAuthority, RuntimeActorMaterializationError> {
358    let authority = bindings
359        .__runtime_authority()
360        .downcast_ref::<SessionRuntimeBindingsAuthority>()
361        .ok_or_else(|| {
362            RuntimeActorMaterializationError::InvalidAuthority(
363                "session runtime bindings lack MeerkatMachine authority".to_string(),
364            )
365        })?;
366    if bindings.session_id() != &authority.session_id || bindings.epoch_id() != &authority.epoch_id
367    {
368        return Err(RuntimeActorMaterializationError::InvalidAuthority(
369            "session runtime binding identity does not match its machine authority".into(),
370        ));
371    }
372    Ok(authority)
373}
374
375fn validate_materialization_registration_authority(
376    session_id: &meerkat_core::SessionId,
377    epoch_id: &meerkat_core::RuntimeEpochId,
378    teardown_gate: &Arc<handles::HandleTeardownGate>,
379    generated_authority: &meerkat_machine::dsl::MeerkatMachineAuthority,
380    phase_policy: RuntimeActorMaterializationPhasePolicy,
381) -> Result<(), RuntimeActorMaterializationError> {
382    let state = generated_authority.state();
383    let expected_session_id = meerkat_machine::dsl::SessionId::from_domain(session_id);
384    let expected_epoch_id = meerkat_machine::dsl::RuntimeEpochId::from_domain(epoch_id);
385    let runtime_phase =
386        meerkat_machine::dsl_authority::runtime_phase_from_authority(generated_authority);
387    if !teardown_gate.is_open()
388        || state.session_id.as_ref() != Some(&expected_session_id)
389        || state.registration_phase == meerkat_machine::dsl::RegistrationPhase::Draining
390        || matches!(
391            runtime_phase,
392            crate::runtime_state::RuntimeState::Stopped
393                | crate::runtime_state::RuntimeState::Destroyed
394        )
395        || match phase_policy {
396            RuntimeActorMaterializationPhasePolicy::RejectRetired => {
397                runtime_phase == crate::runtime_state::RuntimeState::Retired
398            }
399            RuntimeActorMaterializationPhasePolicy::RequireArchivedRevivalBoundary => !matches!(
400                runtime_phase,
401                crate::runtime_state::RuntimeState::Retired
402                    | crate::runtime_state::RuntimeState::Idle
403            ),
404        }
405        || state
406            .active_runtime_epoch_id
407            .as_ref()
408            .is_some_and(|epoch_id| epoch_id != &expected_epoch_id)
409    {
410        return Err(RuntimeActorMaterializationError::RegistrationClosed);
411    }
412    Ok(())
413}
414
415/// Begin exclusive construction of the live actor for an exact prepared
416/// binding. This is the cancellation-safe successor to the read-only validator.
417pub fn begin_session_runtime_actor_materialization(
418    bindings: &meerkat_core::SessionRuntimeBindings,
419) -> Result<RuntimeActorMaterializationPermit, RuntimeActorMaterializationError> {
420    begin_session_runtime_actor_materialization_with_phase_policy(
421        bindings,
422        RuntimeActorMaterializationPhasePolicy::RejectRetired,
423        None,
424    )
425}
426
427/// Begin actor construction for the exact machine-authorized
428/// Archived+Retired revival midpoint.
429///
430/// The public archived-resume service path remains closed; only a caller that
431/// already holds MeerkatMachine's session-control capability may construct the
432/// temporary live actor that the same revival transaction will promote to
433/// Active+Idle before executor attachment.
434pub async fn begin_session_runtime_actor_materialization_for_archived_resume(
435    bindings: &meerkat_core::SessionRuntimeBindings,
436    authorization: crate::meerkat_machine::ArchivedSessionActorMaterializationAuthorization,
437) -> Result<RuntimeActorMaterializationPermit, RuntimeActorMaterializationError> {
438    authorization.begin(bindings).await
439}
440
441fn begin_session_runtime_actor_materialization_with_phase_policy(
442    bindings: &meerkat_core::SessionRuntimeBindings,
443    phase_policy: RuntimeActorMaterializationPhasePolicy,
444    mutation_guard: Option<crate::tokio::sync::OwnedMutexGuard<()>>,
445) -> Result<RuntimeActorMaterializationPermit, RuntimeActorMaterializationError> {
446    let authority = validated_session_runtime_bindings_authority(bindings)?;
447    let generated_authority = authority
448        .dsl_authority
449        .lock()
450        .unwrap_or_else(std::sync::PoisonError::into_inner);
451    validate_materialization_registration_authority(
452        &authority.session_id,
453        &authority.epoch_id,
454        &authority.teardown_gate,
455        &generated_authority,
456        phase_policy,
457    )?;
458    let (claim_id, previous_phase, transactional) = {
459        let mut state = authority
460            .materialization_claim_state
461            .lock()
462            .unwrap_or_else(std::sync::PoisonError::into_inner);
463        if let Some(claim_id) = authority.materialization_claim_id {
464            if !state.exact_claim_is(
465                claim_id,
466                &[
467                    RuntimeActorMaterializationClaimPhase::Prepared,
468                    RuntimeActorMaterializationClaimPhase::Staged,
469                ],
470            ) {
471                return Err(RuntimeActorMaterializationError::RegistrationClosed);
472            }
473            let previous = state.phase;
474            state.phase = RuntimeActorMaterializationClaimPhase::ActorCreating;
475            (claim_id, previous, true)
476        } else if authority.legacy_actor_materialization_generation
477            == Some(state.legacy_capability_generation)
478            && state.current.is_none()
479            && state.phase == RuntimeActorMaterializationClaimPhase::Vacant
480        {
481            let claim_id = uuid::Uuid::new_v4();
482            state.current = Some(claim_id);
483            state.phase = RuntimeActorMaterializationClaimPhase::ActorCreating;
484            (
485                claim_id,
486                RuntimeActorMaterializationClaimPhase::Vacant,
487                false,
488            )
489        } else {
490            return Err(RuntimeActorMaterializationError::RegistrationClosed);
491        }
492    };
493    drop(generated_authority);
494    Ok(RuntimeActorMaterializationPermit {
495        claim_id,
496        claim_state: Arc::clone(&authority.materialization_claim_state),
497        session_id: authority.session_id.clone(),
498        epoch_id: authority.epoch_id.clone(),
499        dsl_authority: Arc::clone(&authority.dsl_authority),
500        teardown_gate: Arc::clone(&authority.teardown_gate),
501        previous_phase,
502        transactional,
503        phase_policy,
504        _mutation_guard: mutation_guard,
505        committed: false,
506    })
507}
508
509// Re-exports for convenience
510pub use accept::{AcceptOutcome, RejectReason};
511pub use coalescing::{
512    AggregateDescriptor, CoalescingResult, SupersessionScope, check_supersession,
513    create_aggregate_input, is_coalescing_eligible,
514};
515pub use completion::{
516    CompletionCleanupObservation, CompletionHandle, CompletionOutcome, CompletionWaitError,
517};
518pub use delivery_inbox::{
519    RuntimeDeliveryError, RuntimeDeliveryId, RuntimeDeliveryInbox, RuntimeDeliveryKind,
520    RuntimeDeliveryReceipt, RuntimeDeliveryRecord, RuntimeDeliverySubmission,
521};
522pub use driver::{EphemeralRuntimeDriver, PersistentRuntimeDriver, PostAdmissionSignal};
523pub use handles::{
524    HandleDslAuthority, RuntimeAuthLeaseHandle, RuntimeCommsDrainHandle,
525    RuntimeExternalToolSurfaceHandle, RuntimeInteractionStreamHandle,
526    RuntimeMcpServerLifecycleHandle, RuntimeModelRoutingHandle, RuntimePeerCommsHandle,
527    RuntimePeerInteractionHandle, RuntimeSessionAdmissionHandle, RuntimeSessionContextHandle,
528    RuntimeTurnStateHandle,
529};
530pub use identifiers::{
531    CausationId, ConversationId, CorrelationId, EventCodeId, IdempotencyKey, InputKind, KindId,
532    LogicalRuntimeId, PolicyVersion, ProjectionRuleId, RuntimeEventId, SchemaId, SupersessionKey,
533};
534pub use ingress_types::{ContentShape, RequestId, ReservationKey};
535pub use input::{
536    ContinuationInput, ContinuationKind, ExternalEventInput, FlowStepInput, Input, InputDurability,
537    InputHeader, InputOrigin, InputVisibility, OperationInput, PeerConvention, PeerInput,
538    PromptInput, ResponseProgressPhase, ResponseTerminalStatus, peer_response_terminal_input,
539    response_terminal_status_from_wire,
540};
541pub use input_ledger::InputLedger;
542pub use input_scope::InputScope;
543pub use input_state::{
544    InputAbandonReason, InputLifecycleState, InputState, InputStateEvent, InputStateHistoryEntry,
545    InputTerminalOutcome, PolicySnapshot, ReconstructionSource,
546};
547pub use meerkat_core::types::HandlingMode;
548#[cfg(not(target_arch = "wasm32"))]
549pub use meerkat_machine::ProviderAuthRuntimeAuthority;
550pub use meerkat_machine::{
551    ArchivedSessionActorMaterializationAuthorization, AuthorizedArchivedResumeCommitLease,
552    CommittedRuntimeExecutorAttachmentPublicationLease, CommsDrainMode, CommsDrainPhase,
553    DrainExitReason, EnsureRuntimeExecutorAttachment, LocalSessionMaterializationMode,
554    MachineServiceTurnCommitLease, MachineServiceTurnIdentity, MachineSessionArchiveLease,
555    MachineSessionControlAuthority, MeerkatConsumerSurface, MeerkatMachine, PeerIngressOwner,
556    PendingRuntimeExecutorAttachment, PreparedArchivedResumeCommitLease,
557    PreparedAttachedSessionActorRecovery, PreparedRuntimeExecutorAttachmentRetirement,
558    PreparedSessionMaterialization, RuntimeBindingsError, RuntimeCleanupTaskSpawner,
559    RuntimeExecutorAttachmentRetirementCompletion, RuntimeExecutorAttachmentWitness,
560    RuntimeLifecycleFacts, RuntimeLoopQueueAdmissionPlan, RuntimeSessionLifecycleObservation,
561    RuntimeSessionRegistrationOutcome, RuntimeSessionRegistrationWitness,
562    StandaloneSessionRuntimeAuthorities, classify_runtime_lifecycle_state,
563    classify_runtime_loop_queue_admission, standalone_session_runtime_authorities,
564    standalone_tool_visibility_owner,
565};
566pub use meerkat_machine_types::{
567    HydratedSessionLlmState, ImageOperationRoutingRequest, ImageOperationRoutingResult,
568    ModelRoutingApprovalDisposition, ModelRoutingRealtimePolicy, ResolvedSessionLlmReconfigure,
569    SessionLlmCapabilitySurface, SessionLlmCapabilitySurfaceStatus, SessionLlmReconfigureHost,
570    SessionLlmReconfigureReport, SessionLlmReconfigureRequest, SessionToolVisibilityDelta,
571};
572#[doc(hidden)]
573pub use meerkat_machine_types::{
574    MeerkatAdmittedInputSnapshot, MeerkatArchiveSnapshot, MeerkatBindingSnapshot,
575    MeerkatCompletionWaiterSnapshot, MeerkatCompletionWaitersSnapshot, MeerkatControlSnapshot,
576    MeerkatCursorSnapshot, MeerkatDrainSnapshot, MeerkatDriverKind, MeerkatInputsSnapshot,
577    MeerkatMachineCatalogInput, MeerkatMachineCommandClassification,
578    MeerkatMachineCommandClassificationRecord, MeerkatMachineCommandVariant,
579    MeerkatMachineFieldlessRuntimeInternalInput, MeerkatMachineRuntimeInternalClassificationRecord,
580    MeerkatMachineRuntimeInternalInput, MeerkatMachineRuntimeInternalReason,
581    MeerkatMachineShellMechanicReason, MeerkatMachineSpineSnapshot, MeerkatOpsSnapshot,
582    OffDrainResponder, SupervisorBridgeCommandAdmissionRoute,
583    SupervisorBridgeCommandClassificationRecord, SupervisorBridgeCommandKind,
584    SupervisorBridgeCommandRealization, canonical_meerkat_machine_command_classifications,
585    canonical_meerkat_machine_command_input_variant_manifest,
586    canonical_meerkat_machine_command_manifest,
587    canonical_meerkat_machine_runtime_internal_classifications,
588    canonical_meerkat_machine_runtime_internal_fieldless_input_variant_manifest,
589    canonical_meerkat_machine_runtime_internal_input_variant_manifest,
590    canonical_meerkat_machine_runtime_internal_manifest,
591    canonical_supervisor_bridge_command_classifications,
592};
593pub use ops_lifecycle::{
594    OpsLifecycleConfig, OpsLifecyclePersistenceRequest, PersistedOpsSnapshot,
595    RuntimeOpsLifecycleRegistry,
596};
597
598#[cfg(all(not(target_arch = "wasm32"), any(test, feature = "test-support")))]
599#[doc(hidden)]
600pub fn test_peer_comms_handle() -> Arc<dyn meerkat_core::handles::PeerCommsHandle> {
601    test_peer_comms_handle_with_silent(std::iter::empty::<String>())
602}
603
604#[cfg(all(not(target_arch = "wasm32"), any(test, feature = "test-support")))]
605#[doc(hidden)]
606#[allow(clippy::expect_used)]
607pub fn test_peer_comms_handle_with_silent<I, S>(
608    silent_intents: I,
609) -> Arc<dyn meerkat_core::handles::PeerCommsHandle>
610where
611    I: IntoIterator<Item = S>,
612    S: Into<String>,
613{
614    let silent_intents = silent_intents
615        .into_iter()
616        .map(Into::into)
617        .collect::<Vec<_>>();
618    std::thread::spawn(move || {
619        let runtime = tokio::runtime::Builder::new_current_thread()
620            .enable_all()
621            .build()
622            .expect("test peer-comms runtime should build");
623        runtime.block_on(async move {
624            let machine = MeerkatMachine::ephemeral();
625            let session_id = meerkat_core::SessionId::new();
626            let bindings = machine
627                .prepare_bindings(session_id.clone())
628                .await
629                .expect("generated MeerkatMachine should prepare test peer-comms bindings");
630            if !silent_intents.is_empty() {
631                machine
632                    .set_session_silent_intents(&session_id, silent_intents)
633                    .await
634                    .expect("set silent intents");
635            }
636            Arc::clone(bindings.peer_comms())
637        })
638    })
639    .join()
640    .expect("test peer-comms authority thread should finish")
641}
642
643#[cfg(all(not(target_arch = "wasm32"), any(test, feature = "test-support")))]
644#[doc(hidden)]
645#[allow(clippy::expect_used)]
646pub fn test_peer_input_candidate_from_interaction(
647    interaction: meerkat_core::interaction::InboxInteraction,
648    peer_id: meerkat_core::comms::PeerId,
649) -> meerkat_core::interaction::PeerInputCandidate {
650    use meerkat_core::interaction::{
651        InteractionContent, InteractionId, PeerIngressEnvelopeFacts, PeerIngressEnvelopeKind,
652        PeerIngressFact, PeerIngressIdentity,
653    };
654
655    let handle = test_peer_comms_handle();
656    let facts = PeerIngressEnvelopeFacts {
657        item_id: interaction.id.to_string(),
658        from_peer: interaction.from.clone(),
659        from_peer_id: peer_id,
660        kind: match &interaction.content {
661            InteractionContent::Message { body, .. }
662            | InteractionContent::IncarnationFencedMessage { body, .. } => {
663                PeerIngressEnvelopeKind::Message { body: body.clone() }
664            }
665            InteractionContent::Request { intent, params, .. } => {
666                PeerIngressEnvelopeKind::Request {
667                    intent: intent.clone(),
668                    params: params.clone(),
669                }
670            }
671            InteractionContent::Response {
672                in_reply_to,
673                status,
674                result,
675                ..
676            } => PeerIngressEnvelopeKind::Response {
677                in_reply_to: in_reply_to.to_string(),
678                status: *status,
679                result: result.clone(),
680            },
681        },
682    };
683    let admission = handle
684        .classify_external_envelope(facts)
685        .expect("generated peer-comms authority should classify test interaction");
686    // R084: the admitted sender identity comes from the machine-echoed
687    // canonical peer id on the classification effect, not the local input.
688    let canonical_from_peer_id = admission
689        .from_peer_id
690        .expect("generated envelope classification should echo the canonical sender peer id");
691    let classification = admission.classification;
692    let convention = match &interaction.content {
693        InteractionContent::Message { .. }
694        | InteractionContent::IncarnationFencedMessage { .. } => {
695            meerkat_core::PeerIngressConvention::Message
696        }
697        InteractionContent::Request { intent, .. } => {
698            if let Some(kind) = classification.lifecycle_kind {
699                let peer = admission
700                    .lifecycle_peer
701                    .clone()
702                    .expect("generated lifecycle classification should include a peer subject");
703                meerkat_core::PeerIngressConvention::Lifecycle { kind, peer }
704            } else {
705                let request_id = admission
706                    .request_id
707                    .clone()
708                    .expect("generated request classification should include request id");
709                meerkat_core::PeerIngressConvention::Request {
710                    request_id,
711                    intent: intent.clone(),
712                }
713            }
714        }
715        InteractionContent::Response { status, .. } => {
716            let in_reply_to = admission
717                .request_id
718                .as_deref()
719                .and_then(|id| uuid::Uuid::parse_str(id).ok())
720                .map(InteractionId)
721                .expect("generated response classification should include in-reply-to id");
722            meerkat_core::PeerIngressConvention::Response {
723                in_reply_to,
724                status: *status,
725            }
726        }
727    };
728    let ingress = PeerIngressFact::peer(
729        interaction.id,
730        classification.class,
731        classification.kind,
732        Some(classification.auth),
733        PeerIngressIdentity::new(canonical_from_peer_id, interaction.from.clone(), convention),
734    );
735    let mut candidate = meerkat_core::interaction::PeerInputCandidate::new(
736        interaction,
737        ingress,
738        admission.lifecycle_peer,
739    );
740    candidate.response_terminality = classification.response_terminality;
741    candidate
742}
743
744/// Stamp prompt turn metadata with the runtime-owned input semantics.
745///
746/// This helper exists for runtime-backed service-turn paths that already hold
747/// machine admission and must pass a runtime-classified prompt turn into the
748/// session layer. New prompt materialization should prefer `MeerkatMachine`
749/// input admission so the machine creates this metadata directly.
750pub fn runtime_stamped_prompt_turn_metadata(
751    metadata: Option<RuntimeStampedTurnMetadata>,
752) -> RuntimeStampedTurnMetadata {
753    let input = Input::Prompt(PromptInput::from_content_input(
754        meerkat_core::ContentInput::Text(String::new()),
755        metadata,
756    ));
757    let semantics = runtime_prompt_semantics_from_machine(&input);
758    runtime_loop::for_input(&input, semantics)
759}
760
761#[allow(clippy::expect_used)]
762fn runtime_prompt_semantics_from_machine(input: &Input) -> ingress_types::RuntimeInputSemantics {
763    let mut authority = meerkat_machine::dsl_authority::new_initialized_authority(
764        "generated runtime prompt machine authority must initialize",
765    );
766    let transition = meerkat_machine::dsl::MeerkatMachineMutator::apply(
767        &mut authority,
768        meerkat_machine::dsl::MeerkatMachineInput::ResolveAdmissionPlan {
769            input_id: input.id().to_string(),
770            input_kind: meerkat_machine::dsl::AdmissionInputKind::from(input.kind()),
771            requested_lane: input
772                .handling_mode()
773                .map(meerkat_machine::dsl::InputLane::from),
774            continuation_kind: meerkat_machine::dsl::AdmissionContinuationKind::from(
775                input.continuation_kind(),
776            ),
777            silent_intent_match: false,
778            existing_superseded_input_id: None,
779            runtime_running: false,
780            active_turn_boundary_available: false,
781            without_wake: false,
782        },
783    )
784    .expect("generated admission authority must accept runtime prompt metadata");
785
786    transition
787        .into_effects()
788        .into_iter()
789        .find_map(|effect| match effect {
790            meerkat_machine::dsl::MeerkatMachineEffect::AdmissionResolved {
791                runtime_boundary,
792                runtime_execution_kind,
793                runtime_peer_response_terminal_apply_intent,
794                live_interrupt_required,
795                ..
796            } => Some(ingress_types::RuntimeInputSemantics {
797                boundary: runtime_boundary.into(),
798                execution_kind: runtime_execution_kind.into(),
799                execution_handling_mode: None,
800                peer_response_terminal_apply_intent: runtime_peer_response_terminal_apply_intent
801                    .map(Into::into),
802                live_interrupt_required,
803            }),
804            _ => None,
805        })
806        .expect("generated admission authority must emit prompt runtime semantics")
807}
808
809#[cfg(test)]
810mod runtime_prompt_metadata_tests {
811    #[test]
812    fn runtime_stamped_prompt_turn_metadata_uses_generated_prompt_semantics() {
813        let metadata = super::runtime_stamped_prompt_turn_metadata(None);
814        assert_eq!(
815            metadata.execution_kind,
816            Some(meerkat_core::lifecycle::RuntimeExecutionKind::ContentTurn)
817        );
818        assert!(metadata.peer_response_terminal_apply_intent.is_none());
819    }
820}
821
822#[doc(hidden)]
823pub mod machine_schema_exports {
824    pub fn meerkat_machine_schema() -> meerkat_machine_schema::MachineSchema {
825        meerkat_machine_schema::catalog::dsl::meerkat_machine_schema_metadata()
826            .attach_to(crate::meerkat_machine::dsl::MeerkatMachineState::schema())
827    }
828
829    pub fn auth_machine_schema() -> meerkat_machine_schema::MachineSchema {
830        meerkat_machine_schema::catalog::dsl::auth_machine_schema_metadata()
831            .attach_to(crate::auth_machine::dsl::AuthMachineState::schema())
832    }
833}
834pub use interrupt_public_result::{
835    UserInterruptObservation, UserInterruptPublicResult, resolve_user_interrupt_public_result,
836};
837pub use peer_handling_mode::{PeerHandlingModeError, validate_peer_handling_mode};
838pub use policy::{
839    ApplyMode, ConsumePoint, DrainPolicy, PolicyDecision, QueueMode, RoutingDisposition, WakeMode,
840};
841pub use policy_table::{DefaultPolicyTable, generated_default_policy_version};
842pub use runtime_event::{
843    InputLifecycleEvent, RunLifecycleEvent, RuntimeEvent, RuntimeEventEnvelope,
844    RuntimeProjectionEvent, RuntimeStateChangeEvent, RuntimeTopologyEvent,
845};
846pub use runtime_state::{RuntimeState, RuntimeStateTransitionError};
847pub use service_ext::SessionServiceRuntimeExt;
848#[cfg(feature = "sqlite-store")]
849pub use store::SqliteRuntimeStore;
850pub use store::{
851    CommittedRecoveryBoundary, CommittedWholeBlobProvisionalTail, CommittedWholeBlobSnapshot,
852    HeadCanonicalProvisionalTailAuthority, HeadCanonicalStoreAuthority, InMemoryRuntimeStore,
853    InputStateRow, PreparedDurableTailRecoverySource, PreparedHeadCanonicalProvisionalPromotion,
854    PreparedHeadCanonicalProvisionalTail, PreparedRecoveryEvidence,
855    PreparedRecoveryReceiptDigestEnrichment, PreparedRecoveryReceiptSource,
856    PreparedRuntimeSessionCommit, PreparedRuntimeSessionCommitKind,
857    PreparedRuntimeSessionCommitOutcome, PreparedRuntimeSessionCommitResult,
858    PreparedWholeBlobProvisionalTail, PreparedWholeBlobRewriteBoundary,
859    PreparedWholeBlobRewriteStoreParts, RecoveryCommitStatus, RuntimeDeliveryAuthorityCasOutcome,
860    RuntimeDeliveryAuthorityRecord, RuntimeDeliveryStoreRecord, RuntimeSessionAuthority,
861    RuntimeSessionPersistenceProfile, RuntimeStore, RuntimeStoreError, RuntimeStoreWriteFence,
862    RuntimeStoreWriteFenceOutcome, SerializedSessionSnapshot, VerifiedCommittedWholeBlobPayload,
863    WholeBlobProvisionalTailAuthority, WholeBlobStoreAuthority,
864};
865pub use traits::{
866    DestroyReport, RecoveryReport, RecycleReport, ResetReport, RetireReport, RuntimeControlPlane,
867    RuntimeControlPlaneError, RuntimeDriver, RuntimeDriverError,
868};