Skip to main content

meerkat_runtime/
mob_adapter.rs

1//! MobRuntimeAdapter — bridges mob provisioning to v9 RuntimeDriver lifecycle.
2//!
3//! When a mob member is spawned, the adapter registers a RuntimeDriver for that
4//! session. When retired, the adapter retires/unregisters the driver. Flow steps
5//! are delivered as FlowStepInput through accept_input().
6//!
7//! This adapter is optional — mob works without it (existing SessionService path).
8//! When present, it enables v9 input lifecycle tracking for mob members.
9
10use meerkat_core::lifecycle::InputId;
11use meerkat_core::types::ContentInput;
12use meerkat_core::types::SessionId;
13
14use crate::MeerkatMachine;
15use crate::input::{
16    FlowStepInput, Input, InputDurability, InputHeader, InputOrigin, InputVisibility,
17};
18#[allow(unused_imports)]
19use crate::service_ext::SessionServiceRuntimeExt as _;
20use crate::traits::{RuntimeControlPlaneError, RuntimeDriverError};
21
22/// Create a FlowStepInput for a mob flow step.
23pub fn create_flow_step_input(
24    step_id: &str,
25    instructions: ContentInput,
26    flow_id: &str,
27    step_index: usize,
28    turn_metadata: Option<meerkat_core::lifecycle::run_primitive::RuntimeTurnMetadata>,
29) -> Input {
30    Input::FlowStep(FlowStepInput {
31        header: InputHeader {
32            id: InputId::new(),
33            timestamp: chrono::Utc::now(),
34            source: InputOrigin::Flow {
35                flow_id: flow_id.into(),
36                step_index,
37            },
38            durability: InputDurability::Durable,
39            visibility: InputVisibility::default(),
40            idempotency_key: None,
41            supersession_key: None,
42            correlation_id: None,
43        },
44        step_id: step_id.into(),
45        content: instructions,
46        turn_metadata,
47    })
48}
49
50/// Register a mob member's session with the runtime adapter.
51///
52/// Registration is a control-plane prerequisite. A failed register is propagated
53/// as a typed error rather than swallowed so the mob provisioning caller can abort
54/// instead of proceeding as if the member's runtime exists.
55pub async fn register_mob_member(
56    adapter: &MeerkatMachine,
57    session_id: SessionId,
58) -> Result<(), RuntimeControlPlaneError> {
59    adapter.register_session(session_id).await
60}
61
62/// Unregister a mob member's session from the runtime adapter.
63pub async fn unregister_mob_member(adapter: &MeerkatMachine, session_id: &SessionId) {
64    adapter.unregister_session(session_id).await;
65}
66
67/// Deliver a flow step to a mob member through the runtime path.
68pub async fn deliver_flow_step(
69    adapter: &MeerkatMachine,
70    session_id: &SessionId,
71    step_id: &str,
72    instructions: impl Into<ContentInput>,
73    flow_id: &str,
74    step_index: usize,
75) -> Result<crate::AcceptOutcome, RuntimeDriverError> {
76    let input = create_flow_step_input(step_id, instructions.into(), flow_id, step_index, None);
77    adapter.accept_input(session_id, input).await
78}
79
80/// Retire a mob member's runtime.
81///
82/// If the session is attached to a live `RuntimeLoop`, queued inputs remain
83/// pending for drain. For plain registered sessions without a loop, retirement
84/// abandons queued work because nothing can execute the drain path.
85pub async fn retire_mob_member(
86    adapter: &MeerkatMachine,
87    session_id: &SessionId,
88) -> Result<crate::traits::RetireReport, RuntimeDriverError> {
89    adapter.retire_runtime(session_id).await
90}
91
92#[cfg(test)]
93#[allow(clippy::unwrap_used)]
94mod tests {
95    use super::*;
96    use crate::policy_table::DefaultPolicyTable;
97    use std::sync::Arc;
98
99    #[tokio::test]
100    async fn spawn_creates_runtime_driver_session() {
101        let adapter = Arc::new(MeerkatMachine::ephemeral());
102        let sid = SessionId::new();
103
104        register_mob_member(&adapter, sid.clone()).await.unwrap();
105
106        // Session should have a runtime driver
107        let state = adapter.runtime_state(&sid).await.unwrap();
108        assert_eq!(state, crate::RuntimeState::Idle);
109    }
110
111    #[tokio::test]
112    async fn flow_step_delivered_as_input() {
113        let adapter = Arc::new(MeerkatMachine::ephemeral());
114        let sid = SessionId::new();
115        register_mob_member(&adapter, sid.clone()).await.unwrap();
116
117        let outcome = deliver_flow_step(&adapter, &sid, "step-1", "analyze the data", "flow-1", 0)
118            .await
119            .unwrap();
120
121        assert!(outcome.is_accepted());
122
123        // Verify policy: flow_step → StageRunStart + WakeIfIdle
124        let input = create_flow_step_input("s", "i".into(), "f", 0, None);
125        let policy = DefaultPolicyTable::resolve(&input, true);
126        assert_eq!(policy.apply_mode, crate::ApplyMode::StageRunStart);
127        assert_eq!(policy.wake_mode, crate::WakeMode::WakeIfIdle);
128    }
129
130    #[tokio::test]
131    async fn retire_without_runtime_loop_abandons_pending_inputs() {
132        let adapter = Arc::new(MeerkatMachine::ephemeral());
133        let sid = SessionId::new();
134        register_mob_member(&adapter, sid.clone()).await.unwrap();
135
136        // Accept an input first
137        deliver_flow_step(&adapter, &sid, "s1", "do it", "f1", 0)
138            .await
139            .unwrap();
140
141        // No RuntimeLoop is attached for plain registration, so retirement
142        // abandons queued work instead of leaving it pending forever.
143        let report = retire_mob_member(&adapter, &sid).await.unwrap();
144        assert_eq!(report.inputs_abandoned, 1);
145        assert_eq!(report.inputs_pending_drain, 0);
146    }
147
148    #[tokio::test]
149    async fn create_flow_step_input_preserves_multimodal_blocks() -> Result<(), String> {
150        let input = create_flow_step_input(
151            "s",
152            ContentInput::Blocks(vec![
153                meerkat_core::types::ContentBlock::Text {
154                    text: "inspect image".into(),
155                },
156                meerkat_core::types::ContentBlock::Image {
157                    media_type: "image/png".into(),
158                    data: "abc123".into(),
159                },
160            ]),
161            "f",
162            0,
163            None,
164        );
165
166        let flow_step = match input {
167            Input::FlowStep(flow_step) => flow_step,
168            other => return Err(format!("expected flow step input, got {other:?}")),
169        };
170        assert_eq!(
171            flow_step.content.text_content(),
172            "inspect image\n[image: image/png]"
173        );
174        assert!(matches!(
175            &flow_step.content,
176            meerkat_core::types::ContentInput::Blocks(blocks) if blocks.len() == 2
177        ));
178        Ok(())
179    }
180
181    #[tokio::test]
182    async fn unregister_removes_driver() {
183        let adapter = Arc::new(MeerkatMachine::ephemeral());
184        let sid = SessionId::new();
185        register_mob_member(&adapter, sid.clone()).await.unwrap();
186
187        unregister_mob_member(&adapter, &sid).await;
188
189        // Should fail now
190        let result = adapter.runtime_state(&sid).await;
191        assert!(result.is_err());
192    }
193
194    #[tokio::test]
195    async fn register_exposes_driver_state() {
196        // Mob-member registration is owned by the runtime control plane:
197        // the registered session must have a live driver entry that surfaces
198        // through typed runtime-state queries.
199        let adapter = Arc::new(MeerkatMachine::ephemeral());
200        let sid = SessionId::new();
201        register_mob_member(&adapter, sid.clone()).await.unwrap();
202
203        let active = adapter.list_active_inputs(&sid).await.unwrap();
204        assert!(active.is_empty()); // No inputs yet
205    }
206
207    /// Gate (#99/#277): a control-plane registration that fails (the session was
208    /// destroyed, so the RegisterSession command returns `Destroyed`) must surface
209    /// a typed `Err` from `register_session` rather than being laundered to success.
210    /// Pre-fix the helper dropped the result with `let _ = ...` and returned `()`,
211    /// so the failure was invisible to callers.
212    #[tokio::test]
213    async fn register_session_surfaces_failure_on_destroyed_session() {
214        let adapter = Arc::new(MeerkatMachine::ephemeral());
215        let sid = SessionId::new();
216
217        // Establish the session, then destroy it so re-registration must fail.
218        adapter.register_session(sid.clone()).await.unwrap();
219        let runtime_id = MeerkatMachine::logical_runtime_id(&sid);
220        crate::traits::RuntimeControlPlane::destroy(&*adapter, &runtime_id)
221            .await
222            .unwrap();
223
224        let result = adapter.register_session(sid.clone()).await;
225        assert!(
226            matches!(result, Err(RuntimeControlPlaneError::Internal(_))),
227            "register_session on a destroyed session must surface a typed control-plane error, got {result:?}"
228        );
229    }
230
231    /// Gate (#99): `register_mob_member` must propagate the typed registration
232    /// failure rather than swallow it — a mob member whose runtime cannot be
233    /// registered must not be reported as provisioned.
234    #[tokio::test]
235    async fn register_mob_member_propagates_failure_on_destroyed_session() {
236        let adapter = Arc::new(MeerkatMachine::ephemeral());
237        let sid = SessionId::new();
238
239        register_mob_member(&adapter, sid.clone()).await.unwrap();
240        let runtime_id = MeerkatMachine::logical_runtime_id(&sid);
241        crate::traits::RuntimeControlPlane::destroy(&*adapter, &runtime_id)
242            .await
243            .unwrap();
244
245        let result = register_mob_member(&adapter, sid.clone()).await;
246        assert!(
247            result.is_err(),
248            "register_mob_member must propagate the typed registration failure, got {result:?}"
249        );
250    }
251
252    /// Gate (#41): `set_session_silent_intents` must surface the inner command's
253    /// typed `Err` to the caller instead of dropping it with `let _ = ...`. A
254    /// destroyed session makes the SetSilentIntents command return `Destroyed`.
255    #[tokio::test]
256    async fn set_session_silent_intents_surfaces_failure_on_destroyed_session() {
257        let adapter = Arc::new(MeerkatMachine::ephemeral());
258        let sid = SessionId::new();
259
260        adapter.register_session(sid.clone()).await.unwrap();
261        let runtime_id = MeerkatMachine::logical_runtime_id(&sid);
262        crate::traits::RuntimeControlPlane::destroy(&*adapter, &runtime_id)
263            .await
264            .unwrap();
265
266        let result = adapter
267            .set_session_silent_intents(&sid, vec!["status".to_string()])
268            .await;
269        assert!(
270            matches!(result, Err(RuntimeDriverError::Destroyed)),
271            "set_session_silent_intents must surface the inner command failure, got {result:?}"
272        );
273    }
274}