1#![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 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 pub(crate) rollback_registration_available: bool,
142 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#[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
263pub 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
415pub 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
427pub 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
509pub 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 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
744pub 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};