1use chrono::Utc;
2use std::sync::Arc;
3
4use crate::error::MacpError;
5use crate::extensions::ExtensionProviderRegistry;
6use crate::log_store::{EntryKind, LogEntry, LogStore};
7use crate::metrics::RuntimeMetrics;
8use crate::mode_registry::ModeRegistry;
9use crate::pb::{Envelope, ModeDescriptor};
10use crate::policy::registry::PolicyRegistry;
11use crate::policy::PolicyDefinition;
12use crate::registry::SessionRegistry;
13use crate::session::{
14 extract_ttl_ms, parse_session_start_payload, validate_canonical_session_start_payload_for_mode,
15 validate_session_id_for_acceptance, Session, SessionState,
16};
17use crate::storage::StorageBackend;
18use crate::stream_bus::SessionStreamBus;
19
20#[derive(Debug)]
21pub struct ProcessResult {
22 pub session_state: SessionState,
23 pub duplicate: bool,
24}
25
26#[derive(Clone, Debug)]
27pub enum SessionLifecycleEvent {
28 Created { session_id: String },
29 Resolved { session_id: String },
30 Expired { session_id: String },
31 Suspended { session_id: String },
32 Resumed { session_id: String },
33 Cancelled { session_id: String },
34}
35
36pub struct Runtime {
37 pub storage: Arc<dyn StorageBackend>,
38 pub registry: Arc<SessionRegistry>,
39 pub log_store: Arc<LogStore>,
40 stream_bus: Arc<SessionStreamBus>,
41 signal_bus: tokio::sync::broadcast::Sender<Envelope>,
42 session_lifecycle_bus: tokio::sync::broadcast::Sender<SessionLifecycleEvent>,
43 mode_registry: Arc<ModeRegistry>,
44 policy_registry: Arc<PolicyRegistry>,
45 #[allow(dead_code)] extensions: Arc<ExtensionProviderRegistry>,
47 metrics: Arc<RuntimeMetrics>,
48 checkpoint_interval: usize,
49}
50
51impl Runtime {
52 pub fn new(
53 storage: Arc<dyn StorageBackend>,
54 registry: Arc<SessionRegistry>,
55 log_store: Arc<LogStore>,
56 ) -> Self {
57 Self::with_mode_registry(
58 storage,
59 registry,
60 log_store,
61 Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
62 macp_policy::DefaultPolicyEvaluator,
63 ))),
64 )
65 }
66
67 pub fn with_mode_registry(
68 storage: Arc<dyn StorageBackend>,
69 registry: Arc<SessionRegistry>,
70 log_store: Arc<LogStore>,
71 mode_registry: Arc<ModeRegistry>,
72 ) -> Self {
73 Self::with_registries(
74 storage,
75 registry,
76 log_store,
77 mode_registry,
78 Arc::new(PolicyRegistry::new()),
79 )
80 }
81
82 pub fn with_registries(
83 storage: Arc<dyn StorageBackend>,
84 registry: Arc<SessionRegistry>,
85 log_store: Arc<LogStore>,
86 mode_registry: Arc<ModeRegistry>,
87 policy_registry: Arc<PolicyRegistry>,
88 ) -> Self {
89 let checkpoint_interval = std::env::var("MACP_CHECKPOINT_INTERVAL")
90 .ok()
91 .and_then(|v| v.parse().ok())
92 .unwrap_or(0); let (signal_tx, _) = tokio::sync::broadcast::channel(256);
94 let (session_lifecycle_tx, _) = tokio::sync::broadcast::channel(64);
95 Self {
96 storage,
97 registry,
98 log_store,
99 stream_bus: Arc::new(SessionStreamBus::default()),
100 signal_bus: signal_tx,
101 session_lifecycle_bus: session_lifecycle_tx,
102 mode_registry,
103 policy_registry,
104 extensions: Arc::new(ExtensionProviderRegistry::new()),
105 metrics: Arc::new(RuntimeMetrics::new()),
106 checkpoint_interval,
107 }
108 }
109
110 pub fn registered_mode_names(&self) -> Vec<String> {
113 self.mode_registry.all_mode_names()
114 }
115
116 pub fn standard_mode_descriptors(&self) -> Vec<ModeDescriptor> {
118 self.mode_registry.standard_mode_descriptors()
119 }
120
121 pub fn extension_mode_descriptors(&self) -> Vec<ModeDescriptor> {
123 self.mode_registry.extension_mode_descriptors()
124 }
125
126 pub fn register_extension(&self, descriptor: ModeDescriptor) -> Result<(), String> {
127 self.mode_registry.register_extension(descriptor)
128 }
129
130 pub fn unregister_extension(&self, mode: &str) -> Result<(), String> {
131 self.mode_registry.unregister_extension(mode)
132 }
133
134 pub fn promote_mode(&self, mode: &str, new_name: Option<&str>) -> Result<String, String> {
135 self.mode_registry.promote_mode(mode, new_name)
136 }
137
138 pub fn subscribe_mode_changes(&self) -> tokio::sync::broadcast::Receiver<()> {
139 self.mode_registry.subscribe_changes()
140 }
141
142 pub fn mode_registry(&self) -> &Arc<ModeRegistry> {
143 &self.mode_registry
144 }
145
146 pub fn register_policy(&self, definition: PolicyDefinition) -> Result<(), String> {
149 self.policy_registry.register(definition)
150 }
151
152 pub fn unregister_policy(&self, policy_id: &str) -> Result<(), String> {
153 self.policy_registry.unregister(policy_id)
154 }
155
156 pub fn get_policy(&self, policy_id: &str) -> Option<PolicyDefinition> {
157 self.policy_registry.get(policy_id)
158 }
159
160 pub fn list_policies(&self, mode_filter: Option<&str>) -> Vec<PolicyDefinition> {
161 self.policy_registry.list(mode_filter)
162 }
163
164 pub fn subscribe_policy_changes(&self) -> tokio::sync::broadcast::Receiver<()> {
165 self.policy_registry.subscribe_changes()
166 }
167
168 pub fn policy_registry(&self) -> &Arc<PolicyRegistry> {
169 &self.policy_registry
170 }
171
172 pub fn metrics(&self) -> &Arc<RuntimeMetrics> {
173 &self.metrics
174 }
175
176 pub fn subscribe_session_stream(
177 &self,
178 session_id: &str,
179 ) -> tokio::sync::broadcast::Receiver<Envelope> {
180 self.stream_bus.subscribe(session_id)
181 }
182
183 pub fn subscribe_signals(&self) -> tokio::sync::broadcast::Receiver<Envelope> {
184 self.signal_bus.subscribe()
185 }
186
187 pub fn subscribe_session_lifecycle(
188 &self,
189 ) -> tokio::sync::broadcast::Receiver<SessionLifecycleEvent> {
190 self.session_lifecycle_bus.subscribe()
191 }
192
193 pub async fn get_session_envelopes_after(
199 &self,
200 session_id: &str,
201 after_sequence: u64,
202 ) -> Result<Vec<Envelope>, u64> {
203 Ok(self
204 .log_store
205 .get_incoming_after(session_id, after_sequence)
206 .await?
207 .into_iter()
208 .map(|(_idx, entry)| Envelope {
209 macp_version: if entry.macp_version.is_empty() {
210 "1.0".into()
211 } else {
212 entry.macp_version
213 },
214 mode: entry.mode,
215 message_type: entry.message_type,
216 message_id: entry.message_id,
217 session_id: entry.session_id,
218 sender: entry.sender,
219 timestamp_unix_ms: if entry.timestamp_unix_ms != 0 {
220 entry.timestamp_unix_ms
221 } else {
222 entry.received_at_ms
223 },
224 payload: entry.raw_payload,
225 })
226 .collect())
227 }
228
229 fn publish_accepted_envelope(&self, env: &Envelope) {
230 if !env.session_id.is_empty() {
231 self.stream_bus.publish(&env.session_id, env.clone());
232 }
233 }
234
235 fn audit_verbose(session: &Session) -> bool {
238 session
239 .policy_definition
240 .as_ref()
241 .and_then(|p| p.rules.get("audit"))
242 .and_then(|a| a.get("level"))
243 .and_then(|l| l.as_str())
244 == Some("info")
245 }
246
247 fn make_incoming_entry(env: &Envelope, received_at_ms: i64) -> LogEntry {
248 LogEntry {
249 message_id: env.message_id.clone(),
250 received_at_ms,
251 sender: env.sender.clone(),
252 message_type: env.message_type.clone(),
253 raw_payload: env.payload.clone(),
254 entry_kind: EntryKind::Incoming,
255 session_id: env.session_id.clone(),
256 mode: env.mode.clone(),
257 macp_version: env.macp_version.clone(),
258 timestamp_unix_ms: env.timestamp_unix_ms,
259 bound_mode_version: None,
260 semantics_rev: 0,
261 bound_max_suspend_ms: None,
262 compacted_incoming_ordinals: 0,
263 }
264 }
265
266 fn make_internal_entry(
267 message_type: &str,
268 payload: &[u8],
269 session_id: &str,
270 mode: &str,
271 ) -> LogEntry {
272 let now = Utc::now().timestamp_millis();
273 LogEntry {
274 message_id: String::new(),
275 received_at_ms: now,
276 sender: "_runtime".into(),
277 message_type: message_type.into(),
278 raw_payload: payload.to_vec(),
279 entry_kind: EntryKind::Internal,
280 session_id: session_id.into(),
281 mode: mode.into(),
282 macp_version: "1.0".into(),
283 timestamp_unix_ms: now,
284 bound_mode_version: None,
285 semantics_rev: 0,
286 bound_max_suspend_ms: None,
287 compacted_incoming_ordinals: 0,
288 }
289 }
290
291 async fn save_session_to_storage(&self, session: &Session) {
292 if let Err(err) = self.storage.save_session(session).await {
293 tracing::warn!(
294 session_id = %session.session_id,
295 error = %err,
296 "failed to persist session snapshot"
297 );
298 }
299 }
300
301 async fn maybe_expire_session(
302 &self,
303 session_id: &str,
304 session: &mut Session,
305 ) -> Result<bool, MacpError> {
306 let now = Utc::now().timestamp_millis();
307 let expires = (session.state == SessionState::Open && now > session.ttl_expiry)
310 || (session.state == SessionState::Suspended && session.suspend_cap_exceeded(now));
311 if expires {
312 let entry = Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode);
313 self.storage
314 .append_log_entry(session_id, &entry)
315 .await
316 .map_err(|_| MacpError::StorageFailed)?;
317 self.log_store.append(session_id, entry).await;
318 session.state = SessionState::Expired;
319 session.suspended_at_ms = None;
320 self.metrics.record_session_expired(&session.mode);
321 tracing::info!(session_id, "session expired via TTL");
322 let _ = self
323 .session_lifecycle_bus
324 .send(SessionLifecycleEvent::Expired {
325 session_id: session_id.to_string(),
326 });
327 return Ok(true);
328 }
329 Ok(false)
330 }
331
332 pub async fn process(
333 &self,
334 env: &Envelope,
335 max_open_sessions: Option<usize>,
336 ) -> Result<ProcessResult, MacpError> {
337 match env.message_type.as_str() {
338 "SessionStart" => self.process_session_start(env, max_open_sessions).await,
339 "Signal" | "Progress" => self.process_signal(env).await,
340 _ => self.process_message(env).await,
341 }
342 }
343
344 async fn process_session_start(
345 &self,
346 env: &Envelope,
347 max_open_sessions: Option<usize>,
348 ) -> Result<ProcessResult, MacpError> {
349 if env.mode.trim().is_empty() {
350 return Err(MacpError::InvalidEnvelope);
351 }
352 validate_session_id_for_acceptance(&env.session_id)?;
353 let mode_name = env.mode.as_str();
354 let mode = self
355 .mode_registry
356 .get_mode(mode_name)
357 .ok_or(MacpError::UnknownMode)?;
358
359 let start_payload = parse_session_start_payload(&env.payload)?;
360 let require_complete_start = self.mode_registry.requires_strict_session_start(mode_name);
368 if require_complete_start {
369 validate_canonical_session_start_payload_for_mode(mode_name, &start_payload)?;
370 }
371
372 let descriptor_version = self.mode_registry.get_mode_version(mode_name);
380 if let Some(descriptor_version) = &descriptor_version {
381 if !start_payload.mode_version.is_empty()
382 && &start_payload.mode_version != descriptor_version
383 {
384 tracing::warn!(
385 mode = mode_name,
386 payload_version = %start_payload.mode_version,
387 descriptor_version = %descriptor_version,
388 "mode_version mismatch"
389 );
390 return Err(MacpError::InvalidEnvelope);
391 }
392 }
393 let bound_mode_version: Option<String> = if start_payload.mode_version.is_empty() {
394 descriptor_version
395 } else {
396 None
397 };
398 let effective_mode_version = bound_mode_version
399 .clone()
400 .unwrap_or_else(|| start_payload.mode_version.clone());
401
402 let ttl_ms = extract_ttl_ms(&start_payload)?;
403
404 if let Some(existing) = self.registry.get_shared(&env.session_id).await {
409 let existing = existing.lock().await;
410 if existing.seen_message_ids.contains(&env.message_id) {
411 return Ok(ProcessResult {
412 session_state: existing.state.clone(),
413 duplicate: true,
414 });
415 }
416 return Err(MacpError::SessionAlreadyExists);
417 }
418
419 let effective_policy_version = if start_payload.policy_version.is_empty() {
424 crate::policy::defaults::DEFAULT_POLICY_ID.to_string()
425 } else {
426 start_payload.policy_version.clone()
427 };
428 let policy_definition = match self.policy_registry.resolve(&effective_policy_version) {
429 Ok(policy) => {
430 if policy.mode != "*" && policy.mode != mode_name {
432 return Err(MacpError::InvalidPolicyDefinition);
433 }
434 Some(policy)
435 }
436 Err(_) => {
437 return Err(MacpError::UnknownPolicyVersion);
438 }
439 };
440
441 let accepted_at = Utc::now().timestamp_millis();
442 let ttl_base = if env.timestamp_unix_ms > 0 {
446 env.timestamp_unix_ms
447 } else {
448 accepted_at
449 };
450 let ttl_expiry = ttl_base.saturating_add(ttl_ms);
451 let bound_max_suspend_ms = if start_payload.max_suspend_ms > 0 {
456 start_payload.max_suspend_ms
457 } else {
458 macp_core::session::MAX_SUSPEND_MS
459 };
460 let session = Session::builder(env.session_id.clone(), mode_name, env.sender.clone())
461 .ttl_expiry(ttl_expiry)
462 .ttl_ms(ttl_ms)
463 .max_suspend_ms(bound_max_suspend_ms)
464 .started_at_unix_ms(accepted_at)
465 .participants(start_payload.participants.clone())
466 .intent(start_payload.intent.clone())
467 .mode_version(effective_mode_version)
468 .configuration_version(start_payload.configuration_version.clone())
469 .policy_version(effective_policy_version)
470 .context_id(start_payload.context_id.clone())
471 .extensions(start_payload.extensions.clone())
472 .roots(start_payload.roots.clone())
473 .policy_definition(policy_definition)
474 .build();
475
476 let response = mode.on_session_start(&session, env)?;
477 let semantics_rev = session.semantics_rev;
478
479 let shared = std::sync::Arc::new(tokio::sync::Mutex::new(session));
484 let mut session_guard = shared
487 .clone()
488 .try_lock_owned()
489 .expect("freshly created mutex is uncontended");
490 {
491 let mut map = self.registry.sessions.write().await;
492 if map.contains_key(&env.session_id) {
493 return Err(MacpError::SessionAlreadyExists);
495 }
496 if let Some(max_open) = max_open_sessions {
497 let now = Utc::now().timestamp_millis();
498 let mut count = 0usize;
499 for arc in map.values() {
500 let counts = match arc.try_lock() {
505 Ok(s) => {
506 s.initiator_sender == env.sender
507 && s.state == SessionState::Open
508 && now <= s.ttl_expiry
509 }
510 Err(_) => true,
511 };
512 if counts {
513 count += 1;
514 }
515 }
516 if count >= max_open {
517 return Err(MacpError::RateLimited);
518 }
519 }
520 map.insert(env.session_id.clone(), std::sync::Arc::clone(&shared));
521 }
522
523 let rollback = |runtime: &Self, session_guard: &mut Session| {
528 session_guard.state = SessionState::Expired;
529 let registry = std::sync::Arc::clone(&runtime.registry);
530 let sid = env.session_id.clone();
531 async move {
532 let mut map = registry.sessions.write().await;
533 map.remove(&sid);
534 }
535 };
536
537 if self
539 .storage
540 .create_session_storage(&env.session_id)
541 .await
542 .is_err()
543 {
544 rollback(self, &mut session_guard).await;
545 return Err(MacpError::StorageFailed);
546 }
547 let mut incoming_entry = Self::make_incoming_entry(env, accepted_at);
548 incoming_entry.bound_mode_version = bound_mode_version;
549 incoming_entry.semantics_rev = semantics_rev;
550 incoming_entry.bound_max_suspend_ms = Some(bound_max_suspend_ms);
551 if self
552 .storage
553 .append_log_entry(&env.session_id, &incoming_entry)
554 .await
555 .is_err()
556 {
557 rollback(self, &mut session_guard).await;
558 return Err(MacpError::StorageFailed);
559 }
560
561 self.log_store.create_session_log(&env.session_id).await;
563 self.log_store.append(&env.session_id, incoming_entry).await;
564
565 session_guard
566 .seen_message_ids
567 .insert(env.message_id.clone());
568 session_guard.apply_mode_response(response);
569
570 let result_state = session_guard.state.clone();
571 if let Err(err) = self.storage.save_session(&session_guard).await {
580 tracing::warn!(
581 session_id = %session_guard.session_id,
582 error = %err,
583 "failed to persist session snapshot at SessionStart (recoverable via replay)"
584 );
585 }
586 self.metrics.record_session_start(mode_name);
587 tracing::info!(
588 session_id = %env.session_id,
589 mode = mode_name,
590 sender = %env.sender,
591 "session started"
592 );
593 self.publish_accepted_envelope(env);
599 drop(session_guard);
600 let _ = self
601 .session_lifecycle_bus
602 .send(SessionLifecycleEvent::Created {
603 session_id: env.session_id.clone(),
604 });
605
606 Ok(ProcessResult {
607 session_state: result_state,
608 duplicate: false,
609 })
610 }
611
612 async fn process_message(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
620 let shared = self
627 .registry
628 .get_shared(&env.session_id)
629 .await
630 .ok_or(MacpError::UnknownSession)?;
631 let mut session_guard = shared.lock().await;
632 let session = &mut *session_guard;
633
634 let now_ms = chrono::Utc::now().timestamp_millis();
641 match macp_modes::step::check_preconditions(session, env, now_ms)? {
642 macp_modes::step::Precheck::Duplicate => {
643 return Ok(ProcessResult {
644 session_state: session.state.clone(),
645 duplicate: true,
646 });
647 }
648 macp_modes::step::Precheck::Expired => {
649 let expired = self.maybe_expire_session(&env.session_id, session).await?;
655 debug_assert!(expired, "check_preconditions reported Expired");
656 self.save_session_to_storage(session).await;
657 return Err(MacpError::TtlExpired);
658 }
659 macp_modes::step::Precheck::Proceed => {}
660 }
661
662 let mode = self
663 .mode_registry
664 .get_mode(&session.mode)
665 .ok_or(MacpError::UnknownMode)?;
666 mode.authorize_sender(session, env)?;
667 let accepted_at_ms = Utc::now().timestamp_millis();
670 let response = mode.on_message_at(
671 session,
672 env,
673 &macp_core::mode::MessageContext::new(accepted_at_ms),
674 )?;
675
676 let incoming_entry = Self::make_incoming_entry(env, accepted_at_ms);
678 self.storage
679 .append_log_entry(&env.session_id, &incoming_entry)
680 .await
681 .map_err(|_| MacpError::StorageFailed)?;
682
683 self.log_store.append(&env.session_id, incoming_entry).await;
687 let result_state = macp_modes::step::commit(session, env, response, now_ms);
688
689 self.metrics.record_message_accepted(&session.mode);
690 if env.message_type == "Commitment" {
691 self.metrics.record_commitment_accepted(&session.mode);
692 }
693
694 if Self::audit_verbose(session) {
699 tracing::info!(
700 session_id = %env.session_id,
701 message_type = %env.message_type,
702 sender = %env.sender,
703 state = ?result_state,
704 "message accepted (audit)"
705 );
706 } else {
707 tracing::debug!(
708 session_id = %env.session_id,
709 message_type = %env.message_type,
710 sender = %env.sender,
711 state = ?result_state,
712 "message accepted"
713 );
714 }
715
716 if result_state == SessionState::Resolved {
717 self.metrics.record_session_resolved(&session.mode);
718 tracing::info!(session_id = %env.session_id, mode = %session.mode, "session resolved");
719 let _ = self
720 .session_lifecycle_bus
721 .send(SessionLifecycleEvent::Resolved {
722 session_id: env.session_id.clone(),
723 });
724 }
725
726 self.save_session_to_storage(session).await;
728 if result_state == SessionState::Resolved {
729 if !self.maybe_compact_log(&env.session_id, session).await {
730 self.force_insert_checkpoint(&env.session_id, session).await;
731 }
732 } else {
733 self.maybe_insert_checkpoint(&env.session_id, session).await;
734 }
735 self.publish_accepted_envelope(env);
736
737 Ok(ProcessResult {
738 session_state: result_state,
739 duplicate: false,
740 })
741 }
742
743 async fn process_signal(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
747 if env.message_type == "Signal" && !env.payload.is_empty() {
750 let signal: crate::pb::SignalPayload =
751 prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
752 if signal.signal_type.trim().is_empty() {
753 return Err(MacpError::InvalidPayload);
754 }
755 }
756 if env.message_type == "Progress" && !env.payload.is_empty() {
758 let _: crate::pb::ProgressPayload =
759 prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
760 }
761 tracing::debug!(
762 sender = %env.sender,
763 message_id = %env.message_id,
764 message_type = %env.message_type,
765 "signal received"
766 );
767 let _ = self.signal_bus.send(env.clone());
768 Ok(ProcessResult {
769 session_state: SessionState::Open,
770 duplicate: false,
771 })
772 }
773
774 pub async fn get_session_checked(&self, session_id: &str) -> Option<Session> {
775 let shared = self.registry.get_shared(session_id).await?;
776 let mut session = shared.lock().await;
777 let changed = self
778 .maybe_expire_session(session_id, &mut session)
779 .await
780 .unwrap_or(false);
781 if changed {
782 self.save_session_to_storage(&session).await;
783 }
784 Some(session.clone())
785 }
786
787 pub async fn cancel_session(
791 &self,
792 session_id: &str,
793 reason: &str,
794 cancelled_by: &str,
795 ) -> Result<ProcessResult, MacpError> {
796 let shared = self
797 .registry
798 .get_shared(session_id)
799 .await
800 .ok_or(MacpError::UnknownSession)?;
801 let mut session_guard = shared.lock().await;
802 let session = &mut *session_guard;
803
804 self.maybe_expire_session(session_id, session).await?;
805
806 if session.state.is_terminal() {
809 let result_state = session.state.clone();
810 self.save_session_to_storage(session).await;
811 return Ok(ProcessResult {
812 session_state: result_state,
813 duplicate: false,
814 });
815 }
816
817 let cancel_payload = crate::pb::SessionCancelPayload {
820 reason: reason.to_string(),
821 cancelled_by: cancelled_by.to_string(),
822 };
823 let cancel_entry = Self::make_internal_entry(
824 "SessionCancel",
825 &prost::Message::encode_to_vec(&cancel_payload),
826 session_id,
827 &session.mode,
828 );
829 self.storage
830 .append_log_entry(session_id, &cancel_entry)
831 .await
832 .map_err(|_| MacpError::StorageFailed)?;
833 self.log_store.append(session_id, cancel_entry).await;
834 let _ = session.cancel();
837 self.save_session_to_storage(session).await;
838 if !self.maybe_compact_log(session_id, session).await {
839 self.force_insert_checkpoint(session_id, session).await;
840 }
841 self.metrics.record_session_cancelled(&session.mode);
842 tracing::info!(session_id, reason, "session cancelled");
843 let _ = self
844 .session_lifecycle_bus
845 .send(SessionLifecycleEvent::Cancelled {
846 session_id: session_id.to_string(),
847 });
848
849 Ok(ProcessResult {
850 session_state: SessionState::Cancelled,
851 duplicate: false,
852 })
853 }
854
855 pub async fn suspend_session(
859 &self,
860 session_id: &str,
861 reason: &str,
862 suspended_by: &str,
863 ) -> Result<ProcessResult, MacpError> {
864 let shared = self
865 .registry
866 .get_shared(session_id)
867 .await
868 .ok_or(MacpError::UnknownSession)?;
869 let mut session_guard = shared.lock().await;
870 let session = &mut *session_guard;
871
872 self.maybe_expire_session(session_id, session).await?;
873 if session.state != SessionState::Open {
874 return Err(MacpError::SessionNotOpen);
875 }
876
877 let now_ms = chrono::Utc::now().timestamp_millis();
878 let payload = crate::pb::SessionSuspendPayload {
879 reason: reason.to_string(),
880 suspended_by: suspended_by.to_string(),
881 };
882 let entry = Self::make_internal_entry(
883 "SessionSuspend",
884 &prost::Message::encode_to_vec(&payload),
885 session_id,
886 &session.mode,
887 );
888 self.storage
889 .append_log_entry(session_id, &entry)
890 .await
891 .map_err(|_| MacpError::StorageFailed)?;
892 self.log_store.append(session_id, entry).await;
893 session.suspend(now_ms)?;
894 self.save_session_to_storage(session).await;
895 self.metrics.record_session_suspended(&session.mode);
896 tracing::info!(session_id, reason, "session suspended");
897 let _ = self
898 .session_lifecycle_bus
899 .send(SessionLifecycleEvent::Suspended {
900 session_id: session_id.to_string(),
901 });
902
903 Ok(ProcessResult {
904 session_state: SessionState::Suspended,
905 duplicate: false,
906 })
907 }
908
909 pub async fn resume_session(
913 &self,
914 session_id: &str,
915 reason: &str,
916 resumed_by: &str,
917 ) -> Result<ProcessResult, MacpError> {
918 let shared = self
919 .registry
920 .get_shared(session_id)
921 .await
922 .ok_or(MacpError::UnknownSession)?;
923 let mut session_guard = shared.lock().await;
924 let session = &mut *session_guard;
925
926 if session.state != SessionState::Suspended {
927 return Err(MacpError::SessionNotOpen);
928 }
929
930 let now_ms = chrono::Utc::now().timestamp_millis();
931 let banked_before = session
932 .suspended_at_ms
933 .map(|at| (now_ms - at).max(0))
934 .unwrap_or(0);
935 let payload = crate::pb::SessionResumePayload {
936 reason: reason.to_string(),
937 resumed_by: resumed_by.to_string(),
938 banked_ms: banked_before,
939 };
940 let entry = Self::make_internal_entry(
941 "SessionResume",
942 &prost::Message::encode_to_vec(&payload),
943 session_id,
944 &session.mode,
945 );
946 self.storage
947 .append_log_entry(session_id, &entry)
948 .await
949 .map_err(|_| MacpError::StorageFailed)?;
950 self.log_store.append(session_id, entry).await;
951
952 match session.resume(now_ms) {
954 Ok(()) => {
955 self.save_session_to_storage(session).await;
956 self.metrics.record_session_resumed(&session.mode);
957 tracing::info!(session_id, reason, "session resumed");
958 let _ = self
959 .session_lifecycle_bus
960 .send(SessionLifecycleEvent::Resumed {
961 session_id: session_id.to_string(),
962 });
963 Ok(ProcessResult {
964 session_state: SessionState::Open,
965 duplicate: false,
966 })
967 }
968 Err(_) => {
969 self.save_session_to_storage(session).await;
971 self.metrics.record_session_expired(&session.mode);
972 let _ = self
973 .session_lifecycle_bus
974 .send(SessionLifecycleEvent::Expired {
975 session_id: session_id.to_string(),
976 });
977 Err(MacpError::TtlExpired)
978 }
979 }
980 }
981
982 async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
985 let discarded = match self.log_store.get_log(session_id).await {
989 Some(entries) => {
990 let prior_base: u64 = entries
991 .iter()
992 .filter(|e| e.entry_kind == EntryKind::Checkpoint)
993 .map(|e| e.compacted_incoming_ordinals)
994 .max()
995 .unwrap_or(0);
996 prior_base
997 + entries
998 .iter()
999 .filter(|e| e.entry_kind == EntryKind::Incoming)
1000 .count() as u64
1001 }
1002 None => 0,
1003 };
1004 match crate::storage::compaction::compact_session_log(
1005 &*self.storage,
1006 session_id,
1007 session,
1008 discarded,
1009 )
1010 .await
1011 {
1012 Ok(checkpoint) => {
1013 self.log_store
1017 .replace_session_log(session_id, vec![checkpoint])
1018 .await;
1019 true
1020 }
1021 Err(e) => {
1022 tracing::debug!(
1023 session_id,
1024 error = %e,
1025 "log compaction skipped (backend may not support it)"
1026 );
1027 false
1028 }
1029 }
1030 }
1031
1032 async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
1035 let persisted = crate::registry::PersistedSession::from(session);
1036 let raw_payload = match serde_json::to_vec(&persisted) {
1037 Ok(bytes) => bytes,
1038 Err(e) => {
1039 tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
1040 return;
1041 }
1042 };
1043 let now = Utc::now().timestamp_millis();
1044 let checkpoint = LogEntry {
1045 message_id: String::new(),
1046 received_at_ms: now,
1047 sender: "_runtime".into(),
1048 message_type: "Checkpoint".into(),
1049 raw_payload,
1050 entry_kind: EntryKind::Checkpoint,
1051 session_id: session_id.into(),
1052 mode: session.mode.clone(),
1053 macp_version: String::new(),
1054 timestamp_unix_ms: now,
1055 bound_mode_version: None,
1056 semantics_rev: 0,
1057 bound_max_suspend_ms: None,
1058 compacted_incoming_ordinals: 0,
1059 };
1060 if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
1061 tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
1062 return;
1063 }
1064 self.log_store.append(session_id, checkpoint).await;
1065 tracing::debug!(
1066 session_id,
1067 "forced checkpoint inserted for terminal session"
1068 );
1069 }
1070
1071 async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
1073 if self.checkpoint_interval == 0 {
1074 return;
1075 }
1076 let log_len = self
1077 .log_store
1078 .get_log(session_id)
1079 .await
1080 .map(|l| l.len())
1081 .unwrap_or(0);
1082 if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
1084 return;
1085 }
1086 self.force_insert_checkpoint(session_id, session).await;
1087 tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
1088 }
1089
1090 pub async fn cleanup_expired_sessions(&self) {
1094 let now = Utc::now().timestamp_millis();
1095 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1100 let guard = self.registry.sessions.read().await;
1101 guard
1102 .iter()
1103 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1104 .collect()
1105 };
1106
1107 let mut expired_count = 0usize;
1108 for (session_id, shared) in candidates {
1109 let mut session = shared.lock().await;
1110 if session.state != SessionState::Open || now <= session.ttl_expiry {
1111 continue;
1112 }
1113 let entry = Self::make_internal_entry("TtlExpired", b"", &session_id, &session.mode);
1114 if let Err(e) = self.storage.append_log_entry(&session_id, &entry).await {
1115 tracing::warn!(
1116 session_id,
1117 error = %e,
1118 "failed to write TTL expiry during cleanup"
1119 );
1120 continue;
1121 }
1122 self.log_store.append(&session_id, entry).await;
1123 session.state = SessionState::Expired;
1124 self.metrics.record_session_expired(&session.mode);
1125 self.save_session_to_storage(&session).await;
1126 if !self.maybe_compact_log(&session_id, &session).await {
1127 self.force_insert_checkpoint(&session_id, &session).await;
1128 }
1129 expired_count += 1;
1130 tracing::info!(session_id = %session_id, "session expired via background cleanup");
1131 let _ = self
1132 .session_lifecycle_bus
1133 .send(SessionLifecycleEvent::Expired {
1134 session_id: session_id.clone(),
1135 });
1136 }
1137
1138 if expired_count > 0 {
1139 tracing::info!(count = expired_count, "background cleanup expired sessions");
1140 }
1141 }
1142
1143 pub async fn gc_disk_sessions(&self, retention_secs: u64) -> usize {
1151 let now = Utc::now().timestamp_millis();
1152 let cutoff = now - (retention_secs as i64 * 1000);
1153 let ids = match self.storage.list_session_ids().await {
1154 Ok(ids) => ids,
1155 Err(e) => {
1156 tracing::warn!(error = %e, "disk GC: cannot list sessions");
1157 return 0;
1158 }
1159 };
1160 let mut removed = 0usize;
1161 for id in ids {
1162 let eligible = if let Some(shared) = self.registry.get_shared(&id).await {
1165 let s = shared.lock().await;
1166 s.state.is_terminal() && s.started_at_unix_ms < cutoff
1167 } else {
1168 match self.storage.load_session(&id).await {
1169 Ok(Some(s)) => s.state.is_terminal() && s.started_at_unix_ms < cutoff,
1170 _ => false,
1173 }
1174 };
1175 if !eligible {
1176 continue;
1177 }
1178 match self.storage.delete_session(&id).await {
1179 Ok(()) => {
1180 {
1181 let mut guard = self.registry.sessions.write().await;
1182 guard.remove(&id);
1183 }
1184 self.log_store.remove_session_log(&id).await;
1185 let _ = self.stream_bus.remove_if_unused(&id);
1186 removed += 1;
1187 }
1188 Err(e) => {
1189 tracing::warn!(session_id = %id, error = %e, "disk GC: delete failed");
1190 }
1191 }
1192 }
1193 if removed > 0 {
1194 tracing::info!(count = removed, "disk GC removed terminal sessions");
1195 }
1196 removed
1197 }
1198
1199 pub async fn evict_stale_sessions(&self, retention_secs: u64) {
1205 let now = Utc::now().timestamp_millis();
1206 let cutoff = now - (retention_secs as i64 * 1000);
1207
1208 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1209 let guard = self.registry.sessions.read().await;
1210 guard
1211 .iter()
1212 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1213 .collect()
1214 };
1215 let mut evict_ids = Vec::new();
1216 for (id, shared) in candidates {
1217 let session = shared.lock().await;
1218 if matches!(
1219 session.state,
1220 SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
1221 ) && session.started_at_unix_ms < cutoff
1222 {
1223 evict_ids.push(id);
1224 }
1225 }
1226
1227 if evict_ids.is_empty() {
1228 return;
1229 }
1230 {
1231 let mut guard = self.registry.sessions.write().await;
1232 for id in &evict_ids {
1233 guard.remove(id);
1234 }
1235 }
1236 for id in &evict_ids {
1237 self.log_store.remove_session_log(id).await;
1238 let _ = self.stream_bus.remove_if_unused(id);
1241 }
1242 tracing::info!(
1243 count = evict_ids.len(),
1244 "evicted stale sessions from memory (registry + log cache + stream bus)"
1245 );
1246 }
1247}
1248
1249#[cfg(test)]
1250mod tests {
1251 use super::*;
1252 use crate::decision_pb::ProposalPayload;
1253 use crate::pb::{CommitmentPayload, SessionStartPayload};
1254 use prost::Message;
1255
1256 fn new_sid() -> String {
1257 uuid::Uuid::new_v4().as_hyphenated().to_string()
1258 }
1259
1260 fn make_runtime() -> Runtime {
1261 let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
1262 let registry = Arc::new(SessionRegistry::new());
1263 let log_store = Arc::new(LogStore::new());
1264 Runtime::new(storage, registry, log_store)
1265 }
1266
1267 fn session_start(participants: Vec<String>) -> Vec<u8> {
1268 SessionStartPayload {
1269 intent: "intent".into(),
1270 participants,
1271 mode_version: "1.0.0".into(),
1272 configuration_version: "cfg-1".into(),
1273 policy_version: String::new(),
1274 ttl_ms: 1_000,
1275 context_id: String::new(),
1276 extensions: std::collections::HashMap::new(),
1277 roots: vec![],
1278 max_suspend_ms: 0,
1279 }
1280 .encode_to_vec()
1281 }
1282
1283 fn env(
1284 mode: &str,
1285 message_type: &str,
1286 message_id: &str,
1287 session_id: &str,
1288 sender: &str,
1289 payload: Vec<u8>,
1290 ) -> Envelope {
1291 Envelope {
1292 macp_version: "1.0".into(),
1293 mode: mode.into(),
1294 message_type: message_type.into(),
1295 message_id: message_id.into(),
1296 session_id: session_id.into(),
1297 sender: sender.into(),
1298 timestamp_unix_ms: Utc::now().timestamp_millis(),
1299 payload,
1300 }
1301 }
1302
1303 #[tokio::test]
1304 async fn standard_session_start_is_strict() {
1305 let rt = make_runtime();
1306 let sid = new_sid();
1307 let bad = SessionStartPayload {
1308 ttl_ms: 0,
1309 ..Default::default()
1310 }
1311 .encode_to_vec();
1312 let err = rt
1313 .process(
1314 &env(
1315 "macp.mode.decision.v1",
1316 "SessionStart",
1317 "m1",
1318 &sid,
1319 "agent://orchestrator",
1320 bad,
1321 ),
1322 None,
1323 )
1324 .await
1325 .unwrap_err();
1326 assert!(matches!(
1327 err,
1328 MacpError::InvalidPayload | MacpError::InvalidTtl
1329 ));
1330 }
1331
1332 #[tokio::test]
1345 async fn a_promoted_mode_still_gets_canonical_session_start_validation() {
1346 let mode_registry = Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
1347 macp_policy::DefaultPolicyEvaluator,
1348 )));
1349 mode_registry
1350 .register_extension(crate::pb::ModeDescriptor {
1351 mode: "ext.promoted.v1".into(),
1352 mode_version: "1.0.0".into(),
1353 title: "Promoted".into(),
1354 description: "promotion target".into(),
1355 determinism_class: "semantic-deterministic".into(),
1356 participant_model: "declared".into(),
1357 message_types: vec!["SessionStart".into(), "Commitment".into()],
1358 terminal_message_types: vec!["Commitment".into()],
1359 ..Default::default()
1360 })
1361 .expect("register extension");
1362 assert_eq!(
1363 mode_registry.promote_mode("ext.promoted.v1", None).unwrap(),
1364 "ext.promoted.v1"
1365 );
1366 assert!(
1367 mode_registry.requires_strict_session_start("ext.promoted.v1"),
1368 "promotion must mark the entry strict"
1369 );
1370 assert!(
1371 !crate::session::requires_strict_session_start("ext.promoted.v1"),
1372 "the core's static list must NOT know this name — that disagreement is the point"
1373 );
1374
1375 let rt = Runtime::with_mode_registry(
1376 Arc::new(crate::storage::MemoryBackend),
1377 Arc::new(SessionRegistry::new()),
1378 Arc::new(LogStore::new()),
1379 mode_registry,
1380 );
1381
1382 let err = rt
1384 .process(
1385 &env(
1386 "ext.promoted.v1",
1387 "SessionStart",
1388 "m1",
1389 &new_sid(),
1390 "agent://orchestrator",
1391 session_start(vec![]),
1392 ),
1393 None,
1394 )
1395 .await
1396 .unwrap_err();
1397 assert_eq!(err.to_string(), "InvalidPayload");
1398
1399 let no_versions = SessionStartPayload {
1402 participants: vec!["agent://fraud".into()],
1403 ttl_ms: 1_000,
1404 ..Default::default()
1405 }
1406 .encode_to_vec();
1407 let err = rt
1408 .process(
1409 &env(
1410 "ext.promoted.v1",
1411 "SessionStart",
1412 "m2",
1413 &new_sid(),
1414 "agent://orchestrator",
1415 no_versions,
1416 ),
1417 None,
1418 )
1419 .await
1420 .unwrap_err();
1421 assert_eq!(err.to_string(), "InvalidPayload");
1422
1423 rt.process(
1426 &env(
1427 "ext.promoted.v1",
1428 "SessionStart",
1429 "m3",
1430 &new_sid(),
1431 "agent://orchestrator",
1432 session_start(vec!["agent://fraud".into()]),
1433 ),
1434 None,
1435 )
1436 .await
1437 .expect("a complete SessionStart must still be accepted for a promoted mode");
1438 }
1439
1440 #[tokio::test]
1441 async fn empty_mode_is_rejected() {
1442 let rt = make_runtime();
1443 let sid = new_sid();
1444 let err = rt
1445 .process(
1446 &env(
1447 "",
1448 "SessionStart",
1449 "m1",
1450 &sid,
1451 "agent://orchestrator",
1452 session_start(vec!["agent://fraud".into()]),
1453 ),
1454 None,
1455 )
1456 .await
1457 .unwrap_err();
1458 assert_eq!(err.to_string(), "InvalidEnvelope");
1459 }
1460
1461 #[tokio::test]
1462 async fn rejected_messages_do_not_enter_dedup_state() {
1463 let rt = make_runtime();
1464 let sid = new_sid();
1465 rt.process(
1466 &env(
1467 "macp.mode.decision.v1",
1468 "SessionStart",
1469 "m1",
1470 &sid,
1471 "agent://orchestrator",
1472 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1473 ),
1474 None,
1475 )
1476 .await
1477 .unwrap();
1478
1479 let bad = rt
1480 .process(
1481 &env(
1482 "macp.mode.decision.v1",
1483 "Proposal",
1484 "m2",
1485 &sid,
1486 "agent://fraud",
1487 b"not-protobuf".to_vec(),
1488 ),
1489 None,
1490 )
1491 .await
1492 .unwrap_err();
1493 assert_eq!(bad.to_string(), "InvalidPayload");
1494
1495 let good = ProposalPayload {
1496 proposal_id: "p1".into(),
1497 option: "step-up".into(),
1498 rationale: "risk".into(),
1499 supporting_data: vec![],
1500 }
1501 .encode_to_vec();
1502 let result = rt
1503 .process(
1504 &env(
1505 "macp.mode.decision.v1",
1506 "Proposal",
1507 "m2",
1508 &sid,
1509 "agent://orchestrator",
1510 good,
1511 ),
1512 None,
1513 )
1514 .await
1515 .unwrap();
1516 assert!(!result.duplicate);
1517 }
1518
1519 #[tokio::test]
1520 async fn get_session_transitions_expired_sessions() {
1521 let rt = make_runtime();
1522 let sid = new_sid();
1523 let payload = SessionStartPayload {
1524 intent: "intent".into(),
1525 participants: vec!["agent://fraud".into()],
1526 mode_version: "1.0.0".into(),
1527 configuration_version: "cfg-1".into(),
1528 policy_version: String::new(),
1529 ttl_ms: 1,
1530 context_id: String::new(),
1531 extensions: std::collections::HashMap::new(),
1532 roots: vec![],
1533 max_suspend_ms: 0,
1534 }
1535 .encode_to_vec();
1536 rt.process(
1537 &env(
1538 "macp.mode.decision.v1",
1539 "SessionStart",
1540 "m1",
1541 &sid,
1542 "agent://orchestrator",
1543 payload,
1544 ),
1545 None,
1546 )
1547 .await
1548 .unwrap();
1549 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1550 let session = rt.get_session_checked(&sid).await.unwrap();
1551 assert_eq!(session.state, SessionState::Expired);
1552 }
1553
1554 #[tokio::test]
1555 async fn multi_round_requires_standard_session_start() {
1556 let rt = make_runtime();
1557 let sid = new_sid();
1558 let payload = SessionStartPayload {
1560 participants: vec!["creator".into(), "other".into()],
1561 ..Default::default()
1562 }
1563 .encode_to_vec();
1564 let err = rt
1565 .process(
1566 &env(
1567 "ext.multi_round.v1",
1568 "SessionStart",
1569 "m1",
1570 &sid,
1571 "creator",
1572 payload,
1573 ),
1574 None,
1575 )
1576 .await
1577 .unwrap_err();
1578 assert!(matches!(
1579 err,
1580 MacpError::InvalidPayload | MacpError::InvalidTtl
1581 ));
1582 }
1583
1584 #[tokio::test]
1585 async fn multi_round_valid_session_start() {
1586 let rt = make_runtime();
1587 let sid = new_sid();
1588 let payload = session_start(vec!["alice".into(), "bob".into()]);
1589 rt.process(
1590 &env(
1591 "ext.multi_round.v1",
1592 "SessionStart",
1593 "m1",
1594 &sid,
1595 "coordinator",
1596 payload,
1597 ),
1598 None,
1599 )
1600 .await
1601 .unwrap();
1602 let session = rt.get_session_checked(&sid).await.unwrap();
1603 assert_eq!(session.mode, "ext.multi_round.v1");
1604 assert_eq!(session.participants, vec!["alice", "bob"]);
1605 }
1606
1607 #[tokio::test]
1608 async fn duplicate_session_start_message_id_returns_duplicate() {
1609 let rt = make_runtime();
1610 let sid = new_sid();
1611 let payload = session_start(vec!["agent://fraud".into()]);
1612 rt.process(
1613 &env(
1614 "macp.mode.decision.v1",
1615 "SessionStart",
1616 "m1",
1617 &sid,
1618 "agent://orchestrator",
1619 payload.clone(),
1620 ),
1621 None,
1622 )
1623 .await
1624 .unwrap();
1625
1626 let result = rt
1627 .process(
1628 &env(
1629 "macp.mode.decision.v1",
1630 "SessionStart",
1631 "m1",
1632 &sid,
1633 "agent://orchestrator",
1634 payload,
1635 ),
1636 None,
1637 )
1638 .await
1639 .unwrap();
1640 assert!(result.duplicate);
1641 }
1642
1643 #[tokio::test]
1644 async fn non_start_mode_mismatch_rejected() {
1645 let rt = make_runtime();
1646 let sid = new_sid();
1647 rt.process(
1648 &env(
1649 "macp.mode.decision.v1",
1650 "SessionStart",
1651 "m1",
1652 &sid,
1653 "agent://orchestrator",
1654 session_start(vec!["agent://fraud".into()]),
1655 ),
1656 None,
1657 )
1658 .await
1659 .unwrap();
1660
1661 let proposal = ProposalPayload {
1662 proposal_id: "p1".into(),
1663 option: "step-up".into(),
1664 rationale: "risk".into(),
1665 supporting_data: vec![],
1666 }
1667 .encode_to_vec();
1668 let err = rt
1669 .process(
1670 &env(
1671 "macp.mode.task.v1",
1672 "Proposal",
1673 "m2",
1674 &sid,
1675 "agent://orchestrator",
1676 proposal,
1677 ),
1678 None,
1679 )
1680 .await
1681 .unwrap_err();
1682 assert_eq!(err.to_string(), "InvalidEnvelope");
1683 }
1684
1685 #[tokio::test]
1686 async fn cancel_idempotent_on_already_expired() {
1687 let rt = make_runtime();
1688 let sid = new_sid();
1689 let payload = SessionStartPayload {
1690 intent: "intent".into(),
1691 participants: vec!["agent://fraud".into()],
1692 mode_version: "1.0.0".into(),
1693 configuration_version: "cfg-1".into(),
1694 policy_version: String::new(),
1695 ttl_ms: 1,
1696 context_id: String::new(),
1697 extensions: std::collections::HashMap::new(),
1698 roots: vec![],
1699 max_suspend_ms: 0,
1700 }
1701 .encode_to_vec();
1702 rt.process(
1703 &env(
1704 "macp.mode.decision.v1",
1705 "SessionStart",
1706 "m1",
1707 &sid,
1708 "agent://orchestrator",
1709 payload,
1710 ),
1711 None,
1712 )
1713 .await
1714 .unwrap();
1715 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1716 let result = rt
1717 .cancel_session(&sid, "cleanup", "agent://orchestrator")
1718 .await
1719 .unwrap();
1720 assert_eq!(result.session_state, SessionState::Expired);
1721 }
1722
1723 #[tokio::test]
1724 async fn accepted_envelopes_are_published_in_order() {
1725 let rt = make_runtime();
1726 let sid = new_sid();
1727 let mut events = rt.subscribe_session_stream(&sid);
1728
1729 let start = env(
1730 "macp.mode.decision.v1",
1731 "SessionStart",
1732 "m1",
1733 &sid,
1734 "agent://orchestrator",
1735 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1736 );
1737 rt.process(&start, None).await.unwrap();
1738 let first = events.recv().await.unwrap();
1739 assert_eq!(first.message_id, "m1");
1740 assert_eq!(first.message_type, "SessionStart");
1741
1742 let proposal = ProposalPayload {
1743 proposal_id: "p1".into(),
1744 option: "step-up".into(),
1745 rationale: "risk".into(),
1746 supporting_data: vec![],
1747 }
1748 .encode_to_vec();
1749 let proposal_env = env(
1750 "macp.mode.decision.v1",
1751 "Proposal",
1752 "m2",
1753 &sid,
1754 "agent://orchestrator",
1755 proposal,
1756 );
1757 rt.process(&proposal_env, None).await.unwrap();
1758 let second = events.recv().await.unwrap();
1759 assert_eq!(second.message_id, "m2");
1760 assert_eq!(second.message_type, "Proposal");
1761 }
1762
1763 #[tokio::test]
1764 async fn commitment_versions_are_carried_into_resolution() {
1765 let rt = make_runtime();
1766 let sid = new_sid();
1767 rt.process(
1768 &env(
1769 "macp.mode.proposal.v1",
1770 "SessionStart",
1771 "m1",
1772 &sid,
1773 "agent://buyer",
1774 session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
1775 ),
1776 None,
1777 )
1778 .await
1779 .unwrap();
1780
1781 let proposal = crate::proposal_pb::ProposalPayload {
1782 proposal_id: "p1".into(),
1783 title: "offer".into(),
1784 summary: "summary".into(),
1785 details: vec![],
1786 tags: vec![],
1787 }
1788 .encode_to_vec();
1789 rt.process(
1790 &env(
1791 "macp.mode.proposal.v1",
1792 "Proposal",
1793 "m2",
1794 &sid,
1795 "agent://seller",
1796 proposal,
1797 ),
1798 None,
1799 )
1800 .await
1801 .unwrap();
1802 let accept = crate::proposal_pb::AcceptPayload {
1803 proposal_id: "p1".into(),
1804 reason: String::new(),
1805 }
1806 .encode_to_vec();
1807 rt.process(
1808 &env(
1809 "macp.mode.proposal.v1",
1810 "Accept",
1811 "m3",
1812 &sid,
1813 "agent://seller",
1814 accept.clone(),
1815 ),
1816 None,
1817 )
1818 .await
1819 .unwrap();
1820 rt.process(
1821 &env(
1822 "macp.mode.proposal.v1",
1823 "Accept",
1824 "m4",
1825 &sid,
1826 "agent://buyer",
1827 accept,
1828 ),
1829 None,
1830 )
1831 .await
1832 .unwrap();
1833 let commitment = CommitmentPayload {
1834 commitment_id: "c1".into(),
1835 action: "proposal.accepted".into(),
1836 authority_scope: "commercial".into(),
1837 reason: "bound".into(),
1838 mode_version: "1.0.0".into(),
1839 policy_version: "policy.default".into(),
1840 configuration_version: "cfg-1".into(),
1841 outcome_positive: true,
1842 supersedes: None,
1843 }
1844 .encode_to_vec();
1845 let result = rt
1846 .process(
1847 &env(
1848 "macp.mode.proposal.v1",
1849 "Commitment",
1850 "m5",
1851 &sid,
1852 "agent://buyer",
1853 commitment,
1854 ),
1855 None,
1856 )
1857 .await
1858 .unwrap();
1859 assert_eq!(result.session_state, SessionState::Resolved);
1860 }
1861
1862 #[tokio::test]
1863 async fn max_open_sessions_enforced_under_write_lock() {
1864 let rt = make_runtime();
1865 let sid1 = new_sid();
1866 let sid2 = new_sid();
1867 let sid3 = new_sid();
1868 rt.process(
1869 &env(
1870 "macp.mode.decision.v1",
1871 "SessionStart",
1872 "m1",
1873 &sid1,
1874 "agent://orchestrator",
1875 session_start(vec!["agent://fraud".into()]),
1876 ),
1877 Some(1),
1878 )
1879 .await
1880 .unwrap();
1881
1882 let err = rt
1883 .process(
1884 &env(
1885 "macp.mode.decision.v1",
1886 "SessionStart",
1887 "m2",
1888 &sid2,
1889 "agent://orchestrator",
1890 session_start(vec!["agent://fraud".into()]),
1891 ),
1892 Some(1),
1893 )
1894 .await
1895 .unwrap_err();
1896 assert!(matches!(err, MacpError::RateLimited));
1897
1898 rt.process(
1899 &env(
1900 "macp.mode.decision.v1",
1901 "SessionStart",
1902 "m3",
1903 &sid3,
1904 "agent://other",
1905 session_start(vec!["agent://fraud".into()]),
1906 ),
1907 Some(1),
1908 )
1909 .await
1910 .unwrap();
1911 }
1912
1913 #[tokio::test]
1914 async fn weak_session_id_rejected() {
1915 let rt = make_runtime();
1916 let err = rt
1917 .process(
1918 &env(
1919 "macp.mode.decision.v1",
1920 "SessionStart",
1921 "m1",
1922 "s1",
1923 "agent://orchestrator",
1924 session_start(vec!["agent://fraud".into()]),
1925 ),
1926 None,
1927 )
1928 .await
1929 .unwrap_err();
1930 assert_eq!(err.to_string(), "InvalidSessionId");
1931 }
1932
1933 #[tokio::test]
1934 async fn log_append_failure_rejects_session_start() {
1935 use std::io;
1936 struct FailingBackend;
1937 #[async_trait::async_trait]
1938 impl StorageBackend for FailingBackend {
1939 async fn save_session(&self, _: &Session) -> io::Result<()> {
1940 Ok(())
1941 }
1942 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
1943 Ok(None)
1944 }
1945 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
1946 Ok(vec![])
1947 }
1948 async fn delete_session(&self, _: &str) -> io::Result<()> {
1949 Ok(())
1950 }
1951 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
1952 Ok(vec![])
1953 }
1954 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
1955 Err(io::Error::other("disk full"))
1956 }
1957 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
1958 Ok(vec![])
1959 }
1960 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
1961 Ok(())
1962 }
1963 }
1964
1965 let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
1966 let registry = Arc::new(SessionRegistry::new());
1967 let log_store = Arc::new(LogStore::new());
1968 let rt = Runtime::new(storage, registry, log_store);
1969 let sid = new_sid();
1970
1971 let err = rt
1972 .process(
1973 &env(
1974 "macp.mode.decision.v1",
1975 "SessionStart",
1976 "m1",
1977 &sid,
1978 "agent://orchestrator",
1979 session_start(vec!["agent://fraud".into()]),
1980 ),
1981 None,
1982 )
1983 .await
1984 .unwrap_err();
1985 assert_eq!(err.to_string(), "StorageFailed");
1986 }
1987
1988 #[tokio::test]
1989 async fn log_append_failure_rejects_in_session_message() {
1990 use std::io;
1991 use std::sync::atomic::{AtomicUsize, Ordering};
1992
1993 struct FailOnSecondAppend {
1994 count: AtomicUsize,
1995 }
1996 #[async_trait::async_trait]
1997 impl StorageBackend for FailOnSecondAppend {
1998 async fn save_session(&self, _: &Session) -> io::Result<()> {
1999 Ok(())
2000 }
2001 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2002 Ok(None)
2003 }
2004 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2005 Ok(vec![])
2006 }
2007 async fn delete_session(&self, _: &str) -> io::Result<()> {
2008 Ok(())
2009 }
2010 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2011 Ok(vec![])
2012 }
2013 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2014 let n = self.count.fetch_add(1, Ordering::SeqCst);
2015 if n >= 1 {
2016 Err(io::Error::other("disk full"))
2017 } else {
2018 Ok(())
2019 }
2020 }
2021 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2022 Ok(vec![])
2023 }
2024 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2025 Ok(())
2026 }
2027 }
2028
2029 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2030 count: AtomicUsize::new(0),
2031 });
2032 let registry = Arc::new(SessionRegistry::new());
2033 let log_store = Arc::new(LogStore::new());
2034 let rt = Runtime::new(storage, registry, log_store);
2035 let sid = new_sid();
2036
2037 rt.process(
2039 &env(
2040 "macp.mode.decision.v1",
2041 "SessionStart",
2042 "m1",
2043 &sid,
2044 "agent://orchestrator",
2045 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2046 ),
2047 None,
2048 )
2049 .await
2050 .unwrap();
2051
2052 let proposal = ProposalPayload {
2054 proposal_id: "p1".into(),
2055 option: "step-up".into(),
2056 rationale: "risk".into(),
2057 supporting_data: vec![],
2058 }
2059 .encode_to_vec();
2060 let err = rt
2061 .process(
2062 &env(
2063 "macp.mode.decision.v1",
2064 "Proposal",
2065 "m2",
2066 &sid,
2067 "agent://orchestrator",
2068 proposal,
2069 ),
2070 None,
2071 )
2072 .await
2073 .unwrap_err();
2074 assert_eq!(err.to_string(), "StorageFailed");
2075
2076 let session = rt.get_session_checked(&sid).await.unwrap();
2078 assert!(!session.seen_message_ids.contains("m2"));
2079 }
2080
2081 #[tokio::test]
2082 async fn cancel_session_fails_if_log_append_fails() {
2083 use std::io;
2084 use std::sync::atomic::{AtomicUsize, Ordering};
2085
2086 struct FailOnSecondAppend {
2087 count: AtomicUsize,
2088 }
2089 #[async_trait::async_trait]
2090 impl StorageBackend for FailOnSecondAppend {
2091 async fn save_session(&self, _: &Session) -> io::Result<()> {
2092 Ok(())
2093 }
2094 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2095 Ok(None)
2096 }
2097 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2098 Ok(vec![])
2099 }
2100 async fn delete_session(&self, _: &str) -> io::Result<()> {
2101 Ok(())
2102 }
2103 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2104 Ok(vec![])
2105 }
2106 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2107 let n = self.count.fetch_add(1, Ordering::SeqCst);
2108 if n >= 1 {
2109 Err(io::Error::other("disk full"))
2110 } else {
2111 Ok(())
2112 }
2113 }
2114 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2115 Ok(vec![])
2116 }
2117 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2118 Ok(())
2119 }
2120 }
2121
2122 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2123 count: AtomicUsize::new(0),
2124 });
2125 let registry = Arc::new(SessionRegistry::new());
2126 let log_store = Arc::new(LogStore::new());
2127 let rt = Runtime::new(storage, registry, log_store);
2128 let sid = new_sid();
2129
2130 rt.process(
2131 &env(
2132 "macp.mode.decision.v1",
2133 "SessionStart",
2134 "m1",
2135 &sid,
2136 "agent://orchestrator",
2137 session_start(vec!["agent://fraud".into()]),
2138 ),
2139 None,
2140 )
2141 .await
2142 .unwrap();
2143
2144 let err = rt
2145 .cancel_session(&sid, "test cancel", "agent://orchestrator")
2146 .await
2147 .unwrap_err();
2148 assert_eq!(err.to_string(), "StorageFailed");
2149 }
2150
2151 #[tokio::test]
2152 async fn ttl_expiration_rejects_message() {
2153 let rt = make_runtime();
2154 let sid = new_sid();
2155 let payload = SessionStartPayload {
2156 intent: "intent".into(),
2157 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
2158 mode_version: "1.0.0".into(),
2159 configuration_version: "cfg-1".into(),
2160 policy_version: String::new(),
2161 ttl_ms: 1,
2162 context_id: String::new(),
2163 extensions: std::collections::HashMap::new(),
2164 roots: vec![],
2165 max_suspend_ms: 0,
2166 }
2167 .encode_to_vec();
2168 rt.process(
2169 &env(
2170 "macp.mode.decision.v1",
2171 "SessionStart",
2172 "m1",
2173 &sid,
2174 "agent://orchestrator",
2175 payload,
2176 ),
2177 None,
2178 )
2179 .await
2180 .unwrap();
2181 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2182 let proposal = ProposalPayload {
2183 proposal_id: "p1".into(),
2184 option: "step-up".into(),
2185 rationale: "risk".into(),
2186 supporting_data: vec![],
2187 }
2188 .encode_to_vec();
2189 let err = rt
2190 .process(
2191 &env(
2192 "macp.mode.decision.v1",
2193 "Proposal",
2194 "m2",
2195 &sid,
2196 "agent://orchestrator",
2197 proposal,
2198 ),
2199 None,
2200 )
2201 .await
2202 .unwrap_err();
2203 assert_eq!(err.to_string(), "TtlExpired");
2204 }
2205
2206 #[tokio::test]
2207 async fn cleanup_expired_sessions_marks_expired() {
2208 let rt = make_runtime();
2209 let sid = new_sid();
2210 let payload = SessionStartPayload {
2211 intent: "intent".into(),
2212 participants: vec!["agent://fraud".into()],
2213 mode_version: "1.0.0".into(),
2214 configuration_version: "cfg-1".into(),
2215 policy_version: String::new(),
2216 ttl_ms: 1,
2217 context_id: String::new(),
2218 extensions: std::collections::HashMap::new(),
2219 roots: vec![],
2220 max_suspend_ms: 0,
2221 }
2222 .encode_to_vec();
2223 rt.process(
2224 &env(
2225 "macp.mode.decision.v1",
2226 "SessionStart",
2227 "m1",
2228 &sid,
2229 "agent://orchestrator",
2230 payload,
2231 ),
2232 None,
2233 )
2234 .await
2235 .unwrap();
2236 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2237 rt.cleanup_expired_sessions().await;
2238 let session = rt.get_session_checked(&sid).await.unwrap();
2239 assert_eq!(session.state, SessionState::Expired);
2240 }
2241
2242 #[tokio::test]
2243 async fn evict_stale_sessions_removes_resolved() {
2244 let rt = make_runtime();
2245 let sid = new_sid();
2246 rt.process(
2248 &env(
2249 "macp.mode.decision.v1",
2250 "SessionStart",
2251 "m1",
2252 &sid,
2253 "agent://orchestrator",
2254 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2255 ),
2256 None,
2257 )
2258 .await
2259 .unwrap();
2260 let proposal = ProposalPayload {
2262 proposal_id: "p1".into(),
2263 option: "step-up".into(),
2264 rationale: "risk".into(),
2265 supporting_data: vec![],
2266 }
2267 .encode_to_vec();
2268 rt.process(
2269 &env(
2270 "macp.mode.decision.v1",
2271 "Proposal",
2272 "m2",
2273 &sid,
2274 "agent://orchestrator",
2275 proposal,
2276 ),
2277 None,
2278 )
2279 .await
2280 .unwrap();
2281 let commitment = CommitmentPayload {
2283 commitment_id: "c1".into(),
2284 action: "decision.selected".into(),
2285 authority_scope: "payments".into(),
2286 reason: "bound".into(),
2287 mode_version: "1.0.0".into(),
2288 policy_version: "policy.default".into(),
2289 configuration_version: "cfg-1".into(),
2290 outcome_positive: true,
2291 supersedes: None,
2292 }
2293 .encode_to_vec();
2294 let result = rt
2295 .process(
2296 &env(
2297 "macp.mode.decision.v1",
2298 "Commitment",
2299 "m3",
2300 &sid,
2301 "agent://orchestrator",
2302 commitment,
2303 ),
2304 None,
2305 )
2306 .await
2307 .unwrap();
2308 assert_eq!(result.session_state, SessionState::Resolved);
2309 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2311 rt.evict_stale_sessions(0).await;
2313 assert!(rt.registry.get_session(&sid).await.is_none());
2315 }
2316
2317 #[tokio::test]
2318 async fn session_start_with_wrong_mode_version_rejected() {
2319 let rt = make_runtime();
2320 let sid = new_sid();
2321 let payload = SessionStartPayload {
2322 intent: "test".into(),
2323 participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
2324 mode_version: "99.0.0".into(), configuration_version: "cfg-1".into(),
2326 policy_version: String::new(),
2327 ttl_ms: 60_000,
2328 context_id: String::new(),
2329 extensions: std::collections::HashMap::new(),
2330 roots: vec![],
2331 max_suspend_ms: 0,
2332 }
2333 .encode_to_vec();
2334
2335 let err = rt
2336 .process(
2337 &env(
2338 "macp.mode.decision.v1",
2339 "SessionStart",
2340 "m1",
2341 &sid,
2342 "agent://orchestrator",
2343 payload,
2344 ),
2345 None,
2346 )
2347 .await
2348 .unwrap_err();
2349 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2350 }
2351
2352 #[tokio::test]
2353 async fn signal_empty_signal_type_rejected() {
2354 let rt = make_runtime();
2355 let signal_payload = crate::pb::SignalPayload {
2357 signal_type: String::new(),
2358 data: b"some data".to_vec(),
2359 confidence: 0.0,
2360 correlation_session_id: String::new(),
2361 }
2362 .encode_to_vec();
2363 let signal = Envelope {
2364 macp_version: "1.0".into(),
2365 mode: String::new(),
2366 message_type: "Signal".into(),
2367 message_id: "sig-1".into(),
2368 session_id: String::new(),
2369 sender: "agent://a".into(),
2370 timestamp_unix_ms: 0,
2371 payload: signal_payload,
2372 };
2373 let err = rt.process_signal(&signal).await.unwrap_err();
2374 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2375 }
2376
2377 #[tokio::test]
2378 async fn signal_valid_payload_accepted() {
2379 let rt = make_runtime();
2380 let signal_payload = crate::pb::SignalPayload {
2381 signal_type: "heartbeat".into(),
2382 data: vec![],
2383 confidence: 0.8,
2384 correlation_session_id: String::new(),
2385 }
2386 .encode_to_vec();
2387 let signal = Envelope {
2388 macp_version: "1.0".into(),
2389 mode: String::new(),
2390 message_type: "Signal".into(),
2391 message_id: "sig-2".into(),
2392 session_id: String::new(),
2393 sender: "agent://a".into(),
2394 timestamp_unix_ms: 0,
2395 payload: signal_payload,
2396 };
2397 rt.process_signal(&signal).await.unwrap();
2398 }
2399
2400 #[tokio::test]
2401 async fn signal_empty_payload_accepted() {
2402 let rt = make_runtime();
2403 let signal = Envelope {
2404 macp_version: "1.0".into(),
2405 mode: String::new(),
2406 message_type: "Signal".into(),
2407 message_id: "sig-3".into(),
2408 session_id: String::new(),
2409 sender: "agent://a".into(),
2410 timestamp_unix_ms: 0,
2411 payload: vec![],
2412 };
2413 rt.process_signal(&signal).await.unwrap();
2414 }
2415
2416 #[tokio::test]
2422 async fn ext_mode_empty_version_binds_descriptor_version() {
2423 let rt = make_runtime();
2424 rt.register_extension(ModeDescriptor {
2425 mode: "ext.dyn.v1".into(),
2426 mode_version: "2.5.0".into(),
2427 message_types: vec!["SessionStart".into(), "Note".into(), "Commitment".into()],
2428 terminal_message_types: vec!["Commitment".into()],
2429 ..Default::default()
2430 })
2431 .unwrap();
2432
2433 let sid = new_sid();
2434 let payload = SessionStartPayload {
2435 participants: vec!["alice".into()],
2436 configuration_version: "cfg-1".into(),
2437 ttl_ms: 60_000,
2438 ..Default::default()
2439 }
2440 .encode_to_vec();
2441 rt.process(
2442 &env("ext.dyn.v1", "SessionStart", "m1", &sid, "alice", payload),
2443 None,
2444 )
2445 .await
2446 .unwrap();
2447
2448 let session = rt.get_session_checked(&sid).await.unwrap();
2450 assert_eq!(session.mode_version, "2.5.0");
2451
2452 let bad = CommitmentPayload {
2454 commitment_id: "c1".into(),
2455 action: "work.completed".into(),
2456 authority_scope: "test".into(),
2457 reason: "done".into(),
2458 mode_version: String::new(),
2459 policy_version: "policy.default".into(),
2460 configuration_version: "cfg-1".into(),
2461 outcome_positive: true,
2462 supersedes: None,
2463 }
2464 .encode_to_vec();
2465 let err = rt
2466 .process(
2467 &env("ext.dyn.v1", "Commitment", "m2", &sid, "alice", bad),
2468 None,
2469 )
2470 .await
2471 .unwrap_err();
2472 assert_eq!(err.to_string(), "InvalidPayload");
2473
2474 let good = CommitmentPayload {
2476 commitment_id: "c1".into(),
2477 action: "work.completed".into(),
2478 authority_scope: "test".into(),
2479 reason: "done".into(),
2480 mode_version: "2.5.0".into(),
2481 policy_version: "policy.default".into(),
2482 configuration_version: "cfg-1".into(),
2483 outcome_positive: true,
2484 supersedes: None,
2485 }
2486 .encode_to_vec();
2487 let result = rt
2488 .process(
2489 &env("ext.dyn.v1", "Commitment", "m3", &sid, "alice", good),
2490 None,
2491 )
2492 .await
2493 .unwrap();
2494 assert_eq!(result.session_state, SessionState::Resolved);
2495 }
2496
2497 #[tokio::test]
2500 async fn ext_mode_binding_recorded_on_session_start_log_entry() {
2501 let rt = make_runtime();
2502 rt.register_extension(ModeDescriptor {
2503 mode: "ext.dyn2.v1".into(),
2504 mode_version: "3.0.0".into(),
2505 message_types: vec!["SessionStart".into(), "Commitment".into()],
2506 terminal_message_types: vec!["Commitment".into()],
2507 ..Default::default()
2508 })
2509 .unwrap();
2510
2511 let sid = new_sid();
2512 let payload = SessionStartPayload {
2513 participants: vec!["alice".into()],
2514 configuration_version: "cfg-1".into(),
2515 ttl_ms: 60_000,
2516 ..Default::default()
2517 }
2518 .encode_to_vec();
2519 rt.process(
2520 &env("ext.dyn2.v1", "SessionStart", "m1", &sid, "alice", payload),
2521 None,
2522 )
2523 .await
2524 .unwrap();
2525
2526 let log = rt.log_store.get_log(&sid).await.unwrap();
2527 assert_eq!(log[0].message_type, "SessionStart");
2528 assert_eq!(log[0].bound_mode_version.as_deref(), Some("3.0.0"));
2529
2530 let sid2 = new_sid();
2532 let payload2 = SessionStartPayload {
2533 participants: vec!["alice".into()],
2534 mode_version: "3.0.0".into(),
2535 configuration_version: "cfg-1".into(),
2536 ttl_ms: 60_000,
2537 ..Default::default()
2538 }
2539 .encode_to_vec();
2540 rt.process(
2541 &env(
2542 "ext.dyn2.v1",
2543 "SessionStart",
2544 "m1",
2545 &sid2,
2546 "alice",
2547 payload2,
2548 ),
2549 None,
2550 )
2551 .await
2552 .unwrap();
2553 let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2554 assert_eq!(log2[0].bound_mode_version, None);
2555 }
2556
2557 #[tokio::test]
2562 async fn session_start_binds_and_records_max_suspend_cap() {
2563 let rt = make_runtime();
2564
2565 let sid = new_sid();
2567 let payload = SessionStartPayload {
2568 participants: vec!["alice".into(), "bob".into()],
2569 mode_version: "1.0.0".into(),
2570 configuration_version: "cfg-1".into(),
2571 ttl_ms: 60_000,
2572 max_suspend_ms: 12_345,
2573 ..Default::default()
2574 }
2575 .encode_to_vec();
2576 rt.process(
2577 &env(
2578 "macp.mode.decision.v1",
2579 "SessionStart",
2580 "m1",
2581 &sid,
2582 "alice",
2583 payload,
2584 ),
2585 None,
2586 )
2587 .await
2588 .unwrap();
2589 let log = rt.log_store.get_log(&sid).await.unwrap();
2590 assert_eq!(log[0].bound_max_suspend_ms, Some(12_345));
2591
2592 let sid2 = new_sid();
2594 let payload2 = SessionStartPayload {
2595 participants: vec!["alice".into(), "bob".into()],
2596 mode_version: "1.0.0".into(),
2597 configuration_version: "cfg-1".into(),
2598 ttl_ms: 60_000,
2599 max_suspend_ms: 0,
2600 ..Default::default()
2601 }
2602 .encode_to_vec();
2603 rt.process(
2604 &env(
2605 "macp.mode.decision.v1",
2606 "SessionStart",
2607 "m2",
2608 &sid2,
2609 "alice",
2610 payload2,
2611 ),
2612 None,
2613 )
2614 .await
2615 .unwrap();
2616 let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2617 assert_eq!(
2618 log2[0].bound_max_suspend_ms,
2619 Some(macp_core::session::MAX_SUSPEND_MS)
2620 );
2621 }
2622
2623 #[test]
2624 fn audit_verbosity_reads_policy_rules() {
2625 let mut session = Session::builder("s1", "macp.mode.decision.v1", "a").build();
2626 assert!(!Runtime::audit_verbose(&session));
2627
2628 session.policy_definition = Some(macp_core::policy::PolicyDefinition {
2629 policy_id: "policy.test.audit".into(),
2630 mode: "*".into(),
2631 description: "audited".into(),
2632 rules: serde_json::json!({ "audit": { "level": "info" } }),
2633 schema_version: 1,
2634 });
2635 assert!(Runtime::audit_verbose(&session));
2636
2637 session.policy_definition.as_mut().unwrap().rules =
2638 serde_json::json!({ "audit": { "level": "debug" } });
2639 assert!(!Runtime::audit_verbose(&session));
2640 }
2641
2642 #[tokio::test]
2648 async fn session_start_snapshot_failure_is_nonfatal_after_commit_point() {
2649 use std::io;
2650
2651 struct FailSnapshotBackend;
2652 #[async_trait::async_trait]
2653 impl StorageBackend for FailSnapshotBackend {
2654 async fn create_session_storage(&self, _s: &str) -> io::Result<()> {
2655 Ok(())
2656 }
2657 async fn save_session(&self, _s: &Session) -> io::Result<()> {
2658 Err(io::Error::other("snapshot disk full"))
2659 }
2660 async fn load_session(&self, _s: &str) -> io::Result<Option<Session>> {
2661 Ok(None)
2662 }
2663 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2664 Ok(vec![])
2665 }
2666 async fn delete_session(&self, _s: &str) -> io::Result<()> {
2667 Ok(())
2668 }
2669 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2670 Ok(vec![])
2671 }
2672 async fn append_log_entry(
2673 &self,
2674 _s: &str,
2675 _e: &crate::log_store::LogEntry,
2676 ) -> io::Result<()> {
2677 Ok(())
2678 }
2679 async fn load_log(&self, _s: &str) -> io::Result<Vec<crate::log_store::LogEntry>> {
2680 Ok(vec![])
2681 }
2682 }
2683
2684 let rt = Runtime::new(
2685 Arc::new(FailSnapshotBackend),
2686 Arc::new(SessionRegistry::new()),
2687 Arc::new(LogStore::new()),
2688 );
2689 let sid = new_sid();
2690 let result = rt
2691 .process(
2692 &env(
2693 "macp.mode.decision.v1",
2694 "SessionStart",
2695 "m1",
2696 &sid,
2697 "agent://orchestrator",
2698 session_start(vec!["agent://orchestrator".into()]),
2699 ),
2700 None,
2701 )
2702 .await
2703 .expect("start must succeed: the log append (commit point) succeeded");
2704 assert!(!result.duplicate);
2705 assert!(rt.get_session_checked(&sid).await.is_some());
2707 }
2708}