1use 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
22pub 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
50pub async fn register_mob_member(
56 adapter: &MeerkatMachine,
57 session_id: SessionId,
58) -> Result<(), RuntimeControlPlaneError> {
59 adapter.register_session(session_id).await
60}
61
62pub async fn unregister_mob_member(adapter: &MeerkatMachine, session_id: &SessionId) {
64 adapter.unregister_session(session_id).await;
65}
66
67pub 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
80pub 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 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 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 deliver_flow_step(&adapter, &sid, "s1", "do it", "f1", 0)
138 .await
139 .unwrap();
140
141 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 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 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()); }
206
207 #[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 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 #[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 #[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}