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,
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(
197 &self,
198 session_id: &str,
199 after_sequence: u64,
200 ) -> Vec<Envelope> {
201 self.log_store
202 .get_incoming_after(session_id, after_sequence)
203 .await
204 .into_iter()
205 .map(|(_idx, entry)| Envelope {
206 macp_version: if entry.macp_version.is_empty() {
207 "1.0".into()
208 } else {
209 entry.macp_version
210 },
211 mode: entry.mode,
212 message_type: entry.message_type,
213 message_id: entry.message_id,
214 session_id: entry.session_id,
215 sender: entry.sender,
216 timestamp_unix_ms: if entry.timestamp_unix_ms != 0 {
217 entry.timestamp_unix_ms
218 } else {
219 entry.received_at_ms
220 },
221 payload: entry.raw_payload,
222 })
223 .collect()
224 }
225
226 fn publish_accepted_envelope(&self, env: &Envelope) {
227 if !env.session_id.is_empty() {
228 self.stream_bus.publish(&env.session_id, env.clone());
229 }
230 }
231
232 fn make_incoming_entry(env: &Envelope) -> LogEntry {
233 LogEntry {
234 message_id: env.message_id.clone(),
235 received_at_ms: Utc::now().timestamp_millis(),
236 sender: env.sender.clone(),
237 message_type: env.message_type.clone(),
238 raw_payload: env.payload.clone(),
239 entry_kind: EntryKind::Incoming,
240 session_id: env.session_id.clone(),
241 mode: env.mode.clone(),
242 macp_version: env.macp_version.clone(),
243 timestamp_unix_ms: env.timestamp_unix_ms,
244 }
245 }
246
247 fn make_internal_entry(
248 message_type: &str,
249 payload: &[u8],
250 session_id: &str,
251 mode: &str,
252 ) -> LogEntry {
253 let now = Utc::now().timestamp_millis();
254 LogEntry {
255 message_id: String::new(),
256 received_at_ms: now,
257 sender: "_runtime".into(),
258 message_type: message_type.into(),
259 raw_payload: payload.to_vec(),
260 entry_kind: EntryKind::Internal,
261 session_id: session_id.into(),
262 mode: mode.into(),
263 macp_version: "1.0".into(),
264 timestamp_unix_ms: now,
265 }
266 }
267
268 async fn save_session_to_storage(&self, session: &Session) {
269 if let Err(err) = self.storage.save_session(session).await {
270 tracing::warn!(
271 session_id = %session.session_id,
272 error = %err,
273 "failed to persist session snapshot"
274 );
275 }
276 }
277
278 async fn maybe_expire_session(
279 &self,
280 session_id: &str,
281 session: &mut Session,
282 ) -> Result<bool, MacpError> {
283 let now = Utc::now().timestamp_millis();
284 let expires = (session.state == SessionState::Open && now > session.ttl_expiry)
287 || (session.state == SessionState::Suspended && session.suspend_cap_exceeded(now));
288 if expires {
289 let entry = Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode);
290 self.storage
291 .append_log_entry(session_id, &entry)
292 .await
293 .map_err(|_| MacpError::StorageFailed)?;
294 self.log_store.append(session_id, entry).await;
295 session.state = SessionState::Expired;
296 session.suspended_at_ms = None;
297 self.metrics.record_session_expired(&session.mode);
298 tracing::info!(session_id, "session expired via TTL");
299 let _ = self
300 .session_lifecycle_bus
301 .send(SessionLifecycleEvent::Expired {
302 session_id: session_id.to_string(),
303 });
304 return Ok(true);
305 }
306 Ok(false)
307 }
308
309 pub async fn process(
310 &self,
311 env: &Envelope,
312 max_open_sessions: Option<usize>,
313 ) -> Result<ProcessResult, MacpError> {
314 match env.message_type.as_str() {
315 "SessionStart" => self.process_session_start(env, max_open_sessions).await,
316 "Signal" | "Progress" => self.process_signal(env).await,
317 _ => self.process_message(env).await,
318 }
319 }
320
321 async fn process_session_start(
322 &self,
323 env: &Envelope,
324 max_open_sessions: Option<usize>,
325 ) -> Result<ProcessResult, MacpError> {
326 if env.mode.trim().is_empty() {
327 return Err(MacpError::InvalidEnvelope);
328 }
329 validate_session_id_for_acceptance(&env.session_id)?;
330 let mode_name = env.mode.as_str();
331 let mode = self
332 .mode_registry
333 .get_mode(mode_name)
334 .ok_or(MacpError::UnknownMode)?;
335
336 let start_payload = parse_session_start_payload(&env.payload)?;
337 let require_complete_start = self.mode_registry.requires_strict_session_start(mode_name);
338 if require_complete_start {
339 validate_canonical_session_start_payload(&start_payload)?;
340 }
341
342 if let Some(descriptor_version) = self.mode_registry.get_mode_version(mode_name) {
344 if !start_payload.mode_version.is_empty()
345 && start_payload.mode_version != descriptor_version
346 {
347 tracing::warn!(
348 mode = mode_name,
349 payload_version = %start_payload.mode_version,
350 descriptor_version = %descriptor_version,
351 "mode_version mismatch"
352 );
353 return Err(MacpError::InvalidEnvelope);
354 }
355 }
356
357 let ttl_ms = extract_ttl_ms(&start_payload)?;
358
359 let mut guard = self.registry.sessions.write().await;
360 if let Some(existing) = guard.get(&env.session_id) {
361 if existing.seen_message_ids.contains(&env.message_id) {
362 return Ok(ProcessResult {
363 session_state: existing.state.clone(),
364 duplicate: true,
365 });
366 }
367 return Err(MacpError::SessionAlreadyExists);
368 }
369
370 if let Some(max_open) = max_open_sessions {
374 let now = Utc::now().timestamp_millis();
375 let count = guard
376 .values()
377 .filter(|s| {
378 s.initiator_sender == env.sender
379 && s.state == SessionState::Open
380 && now <= s.ttl_expiry
381 })
382 .count();
383 if count >= max_open {
384 return Err(MacpError::RateLimited);
385 }
386 }
387
388 let effective_policy_version = if start_payload.policy_version.is_empty() {
393 crate::policy::defaults::DEFAULT_POLICY_ID.to_string()
394 } else {
395 start_payload.policy_version.clone()
396 };
397 let policy_definition = match self.policy_registry.resolve(&effective_policy_version) {
398 Ok(policy) => {
399 if policy.mode != "*" && policy.mode != mode_name {
401 return Err(MacpError::InvalidPolicyDefinition);
402 }
403 Some(policy)
404 }
405 Err(_) => {
406 return Err(MacpError::UnknownPolicyVersion);
407 }
408 };
409
410 let accepted_at = Utc::now().timestamp_millis();
411 let ttl_base = if env.timestamp_unix_ms > 0 {
415 env.timestamp_unix_ms
416 } else {
417 accepted_at
418 };
419 let ttl_expiry = ttl_base.saturating_add(ttl_ms);
420 let session = Session {
421 session_id: env.session_id.clone(),
422 state: SessionState::Open,
423 ttl_expiry,
424 ttl_ms,
425 started_at_unix_ms: accepted_at,
426 resolution: None,
427 mode: mode_name.to_string(),
428 mode_state: vec![],
429 participants: start_payload.participants.clone(),
430 seen_message_ids: std::collections::HashSet::new(),
431 intent: start_payload.intent.clone(),
432 mode_version: start_payload.mode_version.clone(),
433 configuration_version: start_payload.configuration_version.clone(),
434 policy_version: effective_policy_version,
435 context_id: start_payload.context_id.clone(),
436 extensions: start_payload.extensions.clone(),
437 roots: start_payload.roots.clone(),
438 initiator_sender: env.sender.clone(),
439 participant_message_counts: std::collections::HashMap::new(),
440 participant_last_seen: std::collections::HashMap::new(),
441 policy_definition,
442 suspended_at_ms: None,
443 accumulated_suspended_ms: 0,
444 };
445
446 let response = mode.on_session_start(&session, env)?;
447
448 self.storage
450 .create_session_storage(&env.session_id)
451 .await
452 .map_err(|_| MacpError::StorageFailed)?;
453 let incoming_entry = Self::make_incoming_entry(env);
454 self.storage
455 .append_log_entry(&env.session_id, &incoming_entry)
456 .await
457 .map_err(|_| MacpError::StorageFailed)?;
458
459 self.log_store.create_session_log(&env.session_id).await;
461 self.log_store.append(&env.session_id, incoming_entry).await;
462
463 let mut session = session;
464 session.seen_message_ids.insert(env.message_id.clone());
465 session.apply_mode_response(response);
466
467 let result_state = session.state.clone();
468 if let Err(err) = self.storage.save_session(&session).await {
472 tracing::error!(
473 session_id = %session.session_id,
474 error = %err,
475 "failed to persist session snapshot at SessionStart"
476 );
477 return Err(MacpError::StorageFailed);
478 }
479 self.metrics.record_session_start(mode_name);
480 tracing::info!(
481 session_id = %env.session_id,
482 mode = mode_name,
483 sender = %env.sender,
484 "session started"
485 );
486 guard.insert(env.session_id.clone(), session);
487 self.publish_accepted_envelope(env);
488 let _ = self
489 .session_lifecycle_bus
490 .send(SessionLifecycleEvent::Created {
491 session_id: env.session_id.clone(),
492 });
493
494 Ok(ProcessResult {
495 session_state: result_state,
496 duplicate: false,
497 })
498 }
499
500 async fn process_message(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
508 let mut guard = self.registry.sessions.write().await;
509 let session = guard
510 .get_mut(&env.session_id)
511 .ok_or(MacpError::UnknownSession)?;
512
513 let now_ms = chrono::Utc::now().timestamp_millis();
520 match macp_modes::step::check_preconditions(session, env, now_ms)? {
521 macp_modes::step::Precheck::Duplicate => {
522 return Ok(ProcessResult {
523 session_state: session.state.clone(),
524 duplicate: true,
525 });
526 }
527 macp_modes::step::Precheck::Expired => {
528 let expired = self.maybe_expire_session(&env.session_id, session).await?;
534 debug_assert!(expired, "check_preconditions reported Expired");
535 self.save_session_to_storage(session).await;
536 return Err(MacpError::TtlExpired);
537 }
538 macp_modes::step::Precheck::Proceed => {}
539 }
540
541 let mode = self
542 .mode_registry
543 .get_mode(&session.mode)
544 .ok_or(MacpError::UnknownMode)?;
545 mode.authorize_sender(session, env)?;
546 let response = mode.on_message(session, env)?;
547
548 let incoming_entry = Self::make_incoming_entry(env);
550 self.storage
551 .append_log_entry(&env.session_id, &incoming_entry)
552 .await
553 .map_err(|_| MacpError::StorageFailed)?;
554
555 self.log_store.append(&env.session_id, incoming_entry).await;
559 let result_state = macp_modes::step::commit(session, env, response, now_ms);
560
561 self.metrics.record_message_accepted(&session.mode);
562 if env.message_type == "Commitment" {
563 self.metrics.record_commitment_accepted(&session.mode);
564 }
565
566 tracing::debug!(
567 session_id = %env.session_id,
568 message_type = %env.message_type,
569 sender = %env.sender,
570 state = ?result_state,
571 "message accepted"
572 );
573
574 if result_state == SessionState::Resolved {
575 self.metrics.record_session_resolved(&session.mode);
576 tracing::info!(session_id = %env.session_id, mode = %session.mode, "session resolved");
577 let _ = self
578 .session_lifecycle_bus
579 .send(SessionLifecycleEvent::Resolved {
580 session_id: env.session_id.clone(),
581 });
582 }
583
584 self.save_session_to_storage(session).await;
586 if result_state == SessionState::Resolved {
587 if !self.maybe_compact_log(&env.session_id, session).await {
588 self.force_insert_checkpoint(&env.session_id, session).await;
589 }
590 } else {
591 self.maybe_insert_checkpoint(&env.session_id, session).await;
592 }
593 self.publish_accepted_envelope(env);
594
595 Ok(ProcessResult {
596 session_state: result_state,
597 duplicate: false,
598 })
599 }
600
601 async fn process_signal(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
605 if env.message_type == "Signal" && !env.payload.is_empty() {
608 let signal: crate::pb::SignalPayload =
609 prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
610 if signal.signal_type.trim().is_empty() {
611 return Err(MacpError::InvalidPayload);
612 }
613 }
614 if env.message_type == "Progress" && !env.payload.is_empty() {
616 let _: crate::pb::ProgressPayload =
617 prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
618 }
619 tracing::debug!(
620 sender = %env.sender,
621 message_id = %env.message_id,
622 message_type = %env.message_type,
623 "signal received"
624 );
625 let _ = self.signal_bus.send(env.clone());
626 Ok(ProcessResult {
627 session_state: SessionState::Open,
628 duplicate: false,
629 })
630 }
631
632 pub async fn get_session_checked(&self, session_id: &str) -> Option<Session> {
633 let mut guard = self.registry.sessions.write().await;
634 let changed = if let Some(session) = guard.get_mut(session_id) {
635 self.maybe_expire_session(session_id, session)
636 .await
637 .unwrap_or(false)
638 } else {
639 return None;
640 };
641 if changed {
642 if let Some(session) = guard.get(session_id) {
643 self.save_session_to_storage(session).await;
644 }
645 }
646 guard.get(session_id).cloned()
647 }
648
649 pub async fn cancel_session(
653 &self,
654 session_id: &str,
655 reason: &str,
656 cancelled_by: &str,
657 ) -> Result<ProcessResult, MacpError> {
658 let mut guard = self.registry.sessions.write().await;
659 let session = guard.get_mut(session_id).ok_or(MacpError::UnknownSession)?;
660
661 self.maybe_expire_session(session_id, session).await?;
662
663 if session.state.is_terminal() {
666 let result_state = session.state.clone();
667 self.save_session_to_storage(session).await;
668 return Ok(ProcessResult {
669 session_state: result_state,
670 duplicate: false,
671 });
672 }
673
674 let cancel_payload = crate::pb::SessionCancelPayload {
677 reason: reason.to_string(),
678 cancelled_by: cancelled_by.to_string(),
679 };
680 let cancel_entry = Self::make_internal_entry(
681 "SessionCancel",
682 &prost::Message::encode_to_vec(&cancel_payload),
683 session_id,
684 &session.mode,
685 );
686 self.storage
687 .append_log_entry(session_id, &cancel_entry)
688 .await
689 .map_err(|_| MacpError::StorageFailed)?;
690 self.log_store.append(session_id, cancel_entry).await;
691 let _ = session.cancel();
694 self.save_session_to_storage(session).await;
695 if !self.maybe_compact_log(session_id, session).await {
696 self.force_insert_checkpoint(session_id, session).await;
697 }
698 self.metrics.record_session_cancelled(&session.mode);
699 tracing::info!(session_id, reason, "session cancelled");
700 let _ = self
701 .session_lifecycle_bus
702 .send(SessionLifecycleEvent::Cancelled {
703 session_id: session_id.to_string(),
704 });
705
706 Ok(ProcessResult {
707 session_state: SessionState::Cancelled,
708 duplicate: false,
709 })
710 }
711
712 pub async fn suspend_session(
716 &self,
717 session_id: &str,
718 reason: &str,
719 suspended_by: &str,
720 ) -> Result<ProcessResult, MacpError> {
721 let mut guard = self.registry.sessions.write().await;
722 let session = guard.get_mut(session_id).ok_or(MacpError::UnknownSession)?;
723
724 self.maybe_expire_session(session_id, session).await?;
725 if session.state != SessionState::Open {
726 return Err(MacpError::SessionNotOpen);
727 }
728
729 let now_ms = chrono::Utc::now().timestamp_millis();
730 let payload = crate::pb::SessionSuspendPayload {
731 reason: reason.to_string(),
732 suspended_by: suspended_by.to_string(),
733 };
734 let entry = Self::make_internal_entry(
735 "SessionSuspend",
736 &prost::Message::encode_to_vec(&payload),
737 session_id,
738 &session.mode,
739 );
740 self.storage
741 .append_log_entry(session_id, &entry)
742 .await
743 .map_err(|_| MacpError::StorageFailed)?;
744 self.log_store.append(session_id, entry).await;
745 session.suspend(now_ms)?;
746 self.save_session_to_storage(session).await;
747 self.metrics.record_session_suspended(&session.mode);
748 tracing::info!(session_id, reason, "session suspended");
749 let _ = self
750 .session_lifecycle_bus
751 .send(SessionLifecycleEvent::Suspended {
752 session_id: session_id.to_string(),
753 });
754
755 Ok(ProcessResult {
756 session_state: SessionState::Suspended,
757 duplicate: false,
758 })
759 }
760
761 pub async fn resume_session(
765 &self,
766 session_id: &str,
767 reason: &str,
768 resumed_by: &str,
769 ) -> Result<ProcessResult, MacpError> {
770 let mut guard = self.registry.sessions.write().await;
771 let session = guard.get_mut(session_id).ok_or(MacpError::UnknownSession)?;
772
773 if session.state != SessionState::Suspended {
774 return Err(MacpError::SessionNotOpen);
775 }
776
777 let now_ms = chrono::Utc::now().timestamp_millis();
778 let banked_before = session
779 .suspended_at_ms
780 .map(|at| (now_ms - at).max(0))
781 .unwrap_or(0);
782 let payload = crate::pb::SessionResumePayload {
783 reason: reason.to_string(),
784 resumed_by: resumed_by.to_string(),
785 banked_ms: banked_before,
786 };
787 let entry = Self::make_internal_entry(
788 "SessionResume",
789 &prost::Message::encode_to_vec(&payload),
790 session_id,
791 &session.mode,
792 );
793 self.storage
794 .append_log_entry(session_id, &entry)
795 .await
796 .map_err(|_| MacpError::StorageFailed)?;
797 self.log_store.append(session_id, entry).await;
798
799 match session.resume(now_ms) {
801 Ok(()) => {
802 self.save_session_to_storage(session).await;
803 self.metrics.record_session_resumed(&session.mode);
804 tracing::info!(session_id, reason, "session resumed");
805 let _ = self
806 .session_lifecycle_bus
807 .send(SessionLifecycleEvent::Resumed {
808 session_id: session_id.to_string(),
809 });
810 Ok(ProcessResult {
811 session_state: SessionState::Open,
812 duplicate: false,
813 })
814 }
815 Err(_) => {
816 self.save_session_to_storage(session).await;
818 self.metrics.record_session_expired(&session.mode);
819 let _ = self
820 .session_lifecycle_bus
821 .send(SessionLifecycleEvent::Expired {
822 session_id: session_id.to_string(),
823 });
824 Err(MacpError::TtlExpired)
825 }
826 }
827 }
828
829 async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
832 match crate::storage::compaction::compact_session_log(&*self.storage, session_id, session)
833 .await
834 {
835 Ok(()) => true,
836 Err(e) => {
837 tracing::debug!(
838 session_id,
839 error = %e,
840 "log compaction skipped (backend may not support it)"
841 );
842 false
843 }
844 }
845 }
846
847 async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
850 let persisted = crate::registry::PersistedSession::from(session);
851 let raw_payload = match serde_json::to_vec(&persisted) {
852 Ok(bytes) => bytes,
853 Err(e) => {
854 tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
855 return;
856 }
857 };
858 let now = Utc::now().timestamp_millis();
859 let checkpoint = LogEntry {
860 message_id: String::new(),
861 received_at_ms: now,
862 sender: "_runtime".into(),
863 message_type: "Checkpoint".into(),
864 raw_payload,
865 entry_kind: EntryKind::Checkpoint,
866 session_id: session_id.into(),
867 mode: session.mode.clone(),
868 macp_version: String::new(),
869 timestamp_unix_ms: now,
870 };
871 if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
872 tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
873 return;
874 }
875 self.log_store.append(session_id, checkpoint).await;
876 tracing::debug!(
877 session_id,
878 "forced checkpoint inserted for terminal session"
879 );
880 }
881
882 async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
884 if self.checkpoint_interval == 0 {
885 return;
886 }
887 let log_len = self
888 .log_store
889 .get_log(session_id)
890 .await
891 .map(|l| l.len())
892 .unwrap_or(0);
893 if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
895 return;
896 }
897 self.force_insert_checkpoint(session_id, session).await;
898 tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
899 }
900
901 pub async fn cleanup_expired_sessions(&self) {
905 let now = Utc::now().timestamp_millis();
906 let mut guard = self.registry.sessions.write().await;
907 let expired_ids: Vec<String> = guard
908 .iter()
909 .filter(|(_, s)| s.state == SessionState::Open && now > s.ttl_expiry)
910 .map(|(id, _)| id.clone())
911 .collect();
912
913 for session_id in &expired_ids {
914 if let Some(session) = guard.get_mut(session_id) {
915 if session.state != SessionState::Open || now <= session.ttl_expiry {
916 continue;
917 }
918 let entry = Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode);
919 if let Err(e) = self.storage.append_log_entry(session_id, &entry).await {
920 tracing::warn!(
921 session_id,
922 error = %e,
923 "failed to write TTL expiry during cleanup"
924 );
925 continue;
926 }
927 self.log_store.append(session_id, entry).await;
928 session.state = SessionState::Expired;
929 self.metrics.record_session_expired(&session.mode);
930 self.save_session_to_storage(session).await;
931 if !self.maybe_compact_log(session_id, session).await {
932 self.force_insert_checkpoint(session_id, session).await;
933 }
934 tracing::info!(session_id, "session expired via background cleanup");
935 let _ = self
936 .session_lifecycle_bus
937 .send(SessionLifecycleEvent::Expired {
938 session_id: session_id.clone(),
939 });
940 }
941 }
942
943 if !expired_ids.is_empty() {
944 tracing::info!(
945 count = expired_ids.len(),
946 "background cleanup expired sessions"
947 );
948 }
949 }
950
951 pub async fn evict_stale_sessions(&self, retention_secs: u64) {
954 let now = Utc::now().timestamp_millis();
955 let cutoff = now - (retention_secs as i64 * 1000);
956 let mut guard = self.registry.sessions.write().await;
957 let evict_ids: Vec<String> = guard
958 .iter()
959 .filter(|(_, s)| {
960 matches!(s.state, SessionState::Resolved | SessionState::Expired)
961 && s.started_at_unix_ms < cutoff
962 })
963 .map(|(id, _)| id.clone())
964 .collect();
965
966 for id in &evict_ids {
967 guard.remove(id);
968 }
969
970 if !evict_ids.is_empty() {
971 tracing::info!(
972 count = evict_ids.len(),
973 "evicted stale sessions from memory"
974 );
975 }
976 }
977}
978
979#[cfg(test)]
980mod tests {
981 use super::*;
982 use crate::decision_pb::ProposalPayload;
983 use crate::pb::{CommitmentPayload, SessionStartPayload};
984 use prost::Message;
985
986 fn new_sid() -> String {
987 uuid::Uuid::new_v4().as_hyphenated().to_string()
988 }
989
990 fn make_runtime() -> Runtime {
991 let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
992 let registry = Arc::new(SessionRegistry::new());
993 let log_store = Arc::new(LogStore::new());
994 Runtime::new(storage, registry, log_store)
995 }
996
997 fn session_start(participants: Vec<String>) -> Vec<u8> {
998 SessionStartPayload {
999 intent: "intent".into(),
1000 participants,
1001 mode_version: "1.0.0".into(),
1002 configuration_version: "cfg-1".into(),
1003 policy_version: String::new(),
1004 ttl_ms: 1_000,
1005 context_id: String::new(),
1006 extensions: std::collections::HashMap::new(),
1007 roots: vec![],
1008 }
1009 .encode_to_vec()
1010 }
1011
1012 fn env(
1013 mode: &str,
1014 message_type: &str,
1015 message_id: &str,
1016 session_id: &str,
1017 sender: &str,
1018 payload: Vec<u8>,
1019 ) -> Envelope {
1020 Envelope {
1021 macp_version: "1.0".into(),
1022 mode: mode.into(),
1023 message_type: message_type.into(),
1024 message_id: message_id.into(),
1025 session_id: session_id.into(),
1026 sender: sender.into(),
1027 timestamp_unix_ms: Utc::now().timestamp_millis(),
1028 payload,
1029 }
1030 }
1031
1032 #[tokio::test]
1033 async fn standard_session_start_is_strict() {
1034 let rt = make_runtime();
1035 let sid = new_sid();
1036 let bad = SessionStartPayload {
1037 ttl_ms: 0,
1038 ..Default::default()
1039 }
1040 .encode_to_vec();
1041 let err = rt
1042 .process(
1043 &env(
1044 "macp.mode.decision.v1",
1045 "SessionStart",
1046 "m1",
1047 &sid,
1048 "agent://orchestrator",
1049 bad,
1050 ),
1051 None,
1052 )
1053 .await
1054 .unwrap_err();
1055 assert!(matches!(
1056 err,
1057 MacpError::InvalidPayload | MacpError::InvalidTtl
1058 ));
1059 }
1060
1061 #[tokio::test]
1062 async fn empty_mode_is_rejected() {
1063 let rt = make_runtime();
1064 let sid = new_sid();
1065 let err = rt
1066 .process(
1067 &env(
1068 "",
1069 "SessionStart",
1070 "m1",
1071 &sid,
1072 "agent://orchestrator",
1073 session_start(vec!["agent://fraud".into()]),
1074 ),
1075 None,
1076 )
1077 .await
1078 .unwrap_err();
1079 assert_eq!(err.to_string(), "InvalidEnvelope");
1080 }
1081
1082 #[tokio::test]
1083 async fn rejected_messages_do_not_enter_dedup_state() {
1084 let rt = make_runtime();
1085 let sid = new_sid();
1086 rt.process(
1087 &env(
1088 "macp.mode.decision.v1",
1089 "SessionStart",
1090 "m1",
1091 &sid,
1092 "agent://orchestrator",
1093 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1094 ),
1095 None,
1096 )
1097 .await
1098 .unwrap();
1099
1100 let bad = rt
1101 .process(
1102 &env(
1103 "macp.mode.decision.v1",
1104 "Proposal",
1105 "m2",
1106 &sid,
1107 "agent://fraud",
1108 b"not-protobuf".to_vec(),
1109 ),
1110 None,
1111 )
1112 .await
1113 .unwrap_err();
1114 assert_eq!(bad.to_string(), "InvalidPayload");
1115
1116 let good = ProposalPayload {
1117 proposal_id: "p1".into(),
1118 option: "step-up".into(),
1119 rationale: "risk".into(),
1120 supporting_data: vec![],
1121 }
1122 .encode_to_vec();
1123 let result = rt
1124 .process(
1125 &env(
1126 "macp.mode.decision.v1",
1127 "Proposal",
1128 "m2",
1129 &sid,
1130 "agent://orchestrator",
1131 good,
1132 ),
1133 None,
1134 )
1135 .await
1136 .unwrap();
1137 assert!(!result.duplicate);
1138 }
1139
1140 #[tokio::test]
1141 async fn get_session_transitions_expired_sessions() {
1142 let rt = make_runtime();
1143 let sid = new_sid();
1144 let payload = SessionStartPayload {
1145 intent: "intent".into(),
1146 participants: vec!["agent://fraud".into()],
1147 mode_version: "1.0.0".into(),
1148 configuration_version: "cfg-1".into(),
1149 policy_version: String::new(),
1150 ttl_ms: 1,
1151 context_id: String::new(),
1152 extensions: std::collections::HashMap::new(),
1153 roots: vec![],
1154 }
1155 .encode_to_vec();
1156 rt.process(
1157 &env(
1158 "macp.mode.decision.v1",
1159 "SessionStart",
1160 "m1",
1161 &sid,
1162 "agent://orchestrator",
1163 payload,
1164 ),
1165 None,
1166 )
1167 .await
1168 .unwrap();
1169 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1170 let session = rt.get_session_checked(&sid).await.unwrap();
1171 assert_eq!(session.state, SessionState::Expired);
1172 }
1173
1174 #[tokio::test]
1175 async fn multi_round_requires_standard_session_start() {
1176 let rt = make_runtime();
1177 let sid = new_sid();
1178 let payload = SessionStartPayload {
1180 participants: vec!["creator".into(), "other".into()],
1181 ..Default::default()
1182 }
1183 .encode_to_vec();
1184 let err = rt
1185 .process(
1186 &env(
1187 "ext.multi_round.v1",
1188 "SessionStart",
1189 "m1",
1190 &sid,
1191 "creator",
1192 payload,
1193 ),
1194 None,
1195 )
1196 .await
1197 .unwrap_err();
1198 assert!(matches!(
1199 err,
1200 MacpError::InvalidPayload | MacpError::InvalidTtl
1201 ));
1202 }
1203
1204 #[tokio::test]
1205 async fn multi_round_valid_session_start() {
1206 let rt = make_runtime();
1207 let sid = new_sid();
1208 let payload = session_start(vec!["alice".into(), "bob".into()]);
1209 rt.process(
1210 &env(
1211 "ext.multi_round.v1",
1212 "SessionStart",
1213 "m1",
1214 &sid,
1215 "coordinator",
1216 payload,
1217 ),
1218 None,
1219 )
1220 .await
1221 .unwrap();
1222 let session = rt.get_session_checked(&sid).await.unwrap();
1223 assert_eq!(session.mode, "ext.multi_round.v1");
1224 assert_eq!(session.participants, vec!["alice", "bob"]);
1225 }
1226
1227 #[tokio::test]
1228 async fn duplicate_session_start_message_id_returns_duplicate() {
1229 let rt = make_runtime();
1230 let sid = new_sid();
1231 let payload = session_start(vec!["agent://fraud".into()]);
1232 rt.process(
1233 &env(
1234 "macp.mode.decision.v1",
1235 "SessionStart",
1236 "m1",
1237 &sid,
1238 "agent://orchestrator",
1239 payload.clone(),
1240 ),
1241 None,
1242 )
1243 .await
1244 .unwrap();
1245
1246 let result = rt
1247 .process(
1248 &env(
1249 "macp.mode.decision.v1",
1250 "SessionStart",
1251 "m1",
1252 &sid,
1253 "agent://orchestrator",
1254 payload,
1255 ),
1256 None,
1257 )
1258 .await
1259 .unwrap();
1260 assert!(result.duplicate);
1261 }
1262
1263 #[tokio::test]
1264 async fn non_start_mode_mismatch_rejected() {
1265 let rt = make_runtime();
1266 let sid = new_sid();
1267 rt.process(
1268 &env(
1269 "macp.mode.decision.v1",
1270 "SessionStart",
1271 "m1",
1272 &sid,
1273 "agent://orchestrator",
1274 session_start(vec!["agent://fraud".into()]),
1275 ),
1276 None,
1277 )
1278 .await
1279 .unwrap();
1280
1281 let proposal = ProposalPayload {
1282 proposal_id: "p1".into(),
1283 option: "step-up".into(),
1284 rationale: "risk".into(),
1285 supporting_data: vec![],
1286 }
1287 .encode_to_vec();
1288 let err = rt
1289 .process(
1290 &env(
1291 "macp.mode.task.v1",
1292 "Proposal",
1293 "m2",
1294 &sid,
1295 "agent://orchestrator",
1296 proposal,
1297 ),
1298 None,
1299 )
1300 .await
1301 .unwrap_err();
1302 assert_eq!(err.to_string(), "InvalidEnvelope");
1303 }
1304
1305 #[tokio::test]
1306 async fn cancel_idempotent_on_already_expired() {
1307 let rt = make_runtime();
1308 let sid = new_sid();
1309 let payload = SessionStartPayload {
1310 intent: "intent".into(),
1311 participants: vec!["agent://fraud".into()],
1312 mode_version: "1.0.0".into(),
1313 configuration_version: "cfg-1".into(),
1314 policy_version: String::new(),
1315 ttl_ms: 1,
1316 context_id: String::new(),
1317 extensions: std::collections::HashMap::new(),
1318 roots: vec![],
1319 }
1320 .encode_to_vec();
1321 rt.process(
1322 &env(
1323 "macp.mode.decision.v1",
1324 "SessionStart",
1325 "m1",
1326 &sid,
1327 "agent://orchestrator",
1328 payload,
1329 ),
1330 None,
1331 )
1332 .await
1333 .unwrap();
1334 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1335 let result = rt
1336 .cancel_session(&sid, "cleanup", "agent://orchestrator")
1337 .await
1338 .unwrap();
1339 assert_eq!(result.session_state, SessionState::Expired);
1340 }
1341
1342 #[tokio::test]
1343 async fn accepted_envelopes_are_published_in_order() {
1344 let rt = make_runtime();
1345 let sid = new_sid();
1346 let mut events = rt.subscribe_session_stream(&sid);
1347
1348 let start = env(
1349 "macp.mode.decision.v1",
1350 "SessionStart",
1351 "m1",
1352 &sid,
1353 "agent://orchestrator",
1354 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1355 );
1356 rt.process(&start, None).await.unwrap();
1357 let first = events.recv().await.unwrap();
1358 assert_eq!(first.message_id, "m1");
1359 assert_eq!(first.message_type, "SessionStart");
1360
1361 let proposal = ProposalPayload {
1362 proposal_id: "p1".into(),
1363 option: "step-up".into(),
1364 rationale: "risk".into(),
1365 supporting_data: vec![],
1366 }
1367 .encode_to_vec();
1368 let proposal_env = env(
1369 "macp.mode.decision.v1",
1370 "Proposal",
1371 "m2",
1372 &sid,
1373 "agent://orchestrator",
1374 proposal,
1375 );
1376 rt.process(&proposal_env, None).await.unwrap();
1377 let second = events.recv().await.unwrap();
1378 assert_eq!(second.message_id, "m2");
1379 assert_eq!(second.message_type, "Proposal");
1380 }
1381
1382 #[tokio::test]
1383 async fn commitment_versions_are_carried_into_resolution() {
1384 let rt = make_runtime();
1385 let sid = new_sid();
1386 rt.process(
1387 &env(
1388 "macp.mode.proposal.v1",
1389 "SessionStart",
1390 "m1",
1391 &sid,
1392 "agent://buyer",
1393 session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
1394 ),
1395 None,
1396 )
1397 .await
1398 .unwrap();
1399
1400 let proposal = crate::proposal_pb::ProposalPayload {
1401 proposal_id: "p1".into(),
1402 title: "offer".into(),
1403 summary: "summary".into(),
1404 details: vec![],
1405 tags: vec![],
1406 }
1407 .encode_to_vec();
1408 rt.process(
1409 &env(
1410 "macp.mode.proposal.v1",
1411 "Proposal",
1412 "m2",
1413 &sid,
1414 "agent://seller",
1415 proposal,
1416 ),
1417 None,
1418 )
1419 .await
1420 .unwrap();
1421 let accept = crate::proposal_pb::AcceptPayload {
1422 proposal_id: "p1".into(),
1423 reason: String::new(),
1424 }
1425 .encode_to_vec();
1426 rt.process(
1427 &env(
1428 "macp.mode.proposal.v1",
1429 "Accept",
1430 "m3",
1431 &sid,
1432 "agent://seller",
1433 accept.clone(),
1434 ),
1435 None,
1436 )
1437 .await
1438 .unwrap();
1439 rt.process(
1440 &env(
1441 "macp.mode.proposal.v1",
1442 "Accept",
1443 "m4",
1444 &sid,
1445 "agent://buyer",
1446 accept,
1447 ),
1448 None,
1449 )
1450 .await
1451 .unwrap();
1452 let commitment = CommitmentPayload {
1453 commitment_id: "c1".into(),
1454 action: "proposal.accepted".into(),
1455 authority_scope: "commercial".into(),
1456 reason: "bound".into(),
1457 mode_version: "1.0.0".into(),
1458 policy_version: "policy.default".into(),
1459 configuration_version: "cfg-1".into(),
1460 outcome_positive: true,
1461 supersedes: None,
1462 }
1463 .encode_to_vec();
1464 let result = rt
1465 .process(
1466 &env(
1467 "macp.mode.proposal.v1",
1468 "Commitment",
1469 "m5",
1470 &sid,
1471 "agent://buyer",
1472 commitment,
1473 ),
1474 None,
1475 )
1476 .await
1477 .unwrap();
1478 assert_eq!(result.session_state, SessionState::Resolved);
1479 }
1480
1481 #[tokio::test]
1482 async fn max_open_sessions_enforced_under_write_lock() {
1483 let rt = make_runtime();
1484 let sid1 = new_sid();
1485 let sid2 = new_sid();
1486 let sid3 = new_sid();
1487 rt.process(
1488 &env(
1489 "macp.mode.decision.v1",
1490 "SessionStart",
1491 "m1",
1492 &sid1,
1493 "agent://orchestrator",
1494 session_start(vec!["agent://fraud".into()]),
1495 ),
1496 Some(1),
1497 )
1498 .await
1499 .unwrap();
1500
1501 let err = rt
1502 .process(
1503 &env(
1504 "macp.mode.decision.v1",
1505 "SessionStart",
1506 "m2",
1507 &sid2,
1508 "agent://orchestrator",
1509 session_start(vec!["agent://fraud".into()]),
1510 ),
1511 Some(1),
1512 )
1513 .await
1514 .unwrap_err();
1515 assert!(matches!(err, MacpError::RateLimited));
1516
1517 rt.process(
1518 &env(
1519 "macp.mode.decision.v1",
1520 "SessionStart",
1521 "m3",
1522 &sid3,
1523 "agent://other",
1524 session_start(vec!["agent://fraud".into()]),
1525 ),
1526 Some(1),
1527 )
1528 .await
1529 .unwrap();
1530 }
1531
1532 #[tokio::test]
1533 async fn weak_session_id_rejected() {
1534 let rt = make_runtime();
1535 let err = rt
1536 .process(
1537 &env(
1538 "macp.mode.decision.v1",
1539 "SessionStart",
1540 "m1",
1541 "s1",
1542 "agent://orchestrator",
1543 session_start(vec!["agent://fraud".into()]),
1544 ),
1545 None,
1546 )
1547 .await
1548 .unwrap_err();
1549 assert_eq!(err.to_string(), "InvalidSessionId");
1550 }
1551
1552 #[tokio::test]
1553 async fn log_append_failure_rejects_session_start() {
1554 use std::io;
1555 struct FailingBackend;
1556 #[async_trait::async_trait]
1557 impl StorageBackend for FailingBackend {
1558 async fn save_session(&self, _: &Session) -> io::Result<()> {
1559 Ok(())
1560 }
1561 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
1562 Ok(None)
1563 }
1564 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
1565 Ok(vec![])
1566 }
1567 async fn delete_session(&self, _: &str) -> io::Result<()> {
1568 Ok(())
1569 }
1570 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
1571 Ok(vec![])
1572 }
1573 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
1574 Err(io::Error::other("disk full"))
1575 }
1576 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
1577 Ok(vec![])
1578 }
1579 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
1580 Ok(())
1581 }
1582 }
1583
1584 let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
1585 let registry = Arc::new(SessionRegistry::new());
1586 let log_store = Arc::new(LogStore::new());
1587 let rt = Runtime::new(storage, registry, log_store);
1588 let sid = new_sid();
1589
1590 let err = rt
1591 .process(
1592 &env(
1593 "macp.mode.decision.v1",
1594 "SessionStart",
1595 "m1",
1596 &sid,
1597 "agent://orchestrator",
1598 session_start(vec!["agent://fraud".into()]),
1599 ),
1600 None,
1601 )
1602 .await
1603 .unwrap_err();
1604 assert_eq!(err.to_string(), "StorageFailed");
1605 }
1606
1607 #[tokio::test]
1608 async fn log_append_failure_rejects_in_session_message() {
1609 use std::io;
1610 use std::sync::atomic::{AtomicUsize, Ordering};
1611
1612 struct FailOnSecondAppend {
1613 count: AtomicUsize,
1614 }
1615 #[async_trait::async_trait]
1616 impl StorageBackend for FailOnSecondAppend {
1617 async fn save_session(&self, _: &Session) -> io::Result<()> {
1618 Ok(())
1619 }
1620 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
1621 Ok(None)
1622 }
1623 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
1624 Ok(vec![])
1625 }
1626 async fn delete_session(&self, _: &str) -> io::Result<()> {
1627 Ok(())
1628 }
1629 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
1630 Ok(vec![])
1631 }
1632 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
1633 let n = self.count.fetch_add(1, Ordering::SeqCst);
1634 if n >= 1 {
1635 Err(io::Error::other("disk full"))
1636 } else {
1637 Ok(())
1638 }
1639 }
1640 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
1641 Ok(vec![])
1642 }
1643 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
1644 Ok(())
1645 }
1646 }
1647
1648 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
1649 count: AtomicUsize::new(0),
1650 });
1651 let registry = Arc::new(SessionRegistry::new());
1652 let log_store = Arc::new(LogStore::new());
1653 let rt = Runtime::new(storage, registry, log_store);
1654 let sid = new_sid();
1655
1656 rt.process(
1658 &env(
1659 "macp.mode.decision.v1",
1660 "SessionStart",
1661 "m1",
1662 &sid,
1663 "agent://orchestrator",
1664 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1665 ),
1666 None,
1667 )
1668 .await
1669 .unwrap();
1670
1671 let proposal = ProposalPayload {
1673 proposal_id: "p1".into(),
1674 option: "step-up".into(),
1675 rationale: "risk".into(),
1676 supporting_data: vec![],
1677 }
1678 .encode_to_vec();
1679 let err = rt
1680 .process(
1681 &env(
1682 "macp.mode.decision.v1",
1683 "Proposal",
1684 "m2",
1685 &sid,
1686 "agent://orchestrator",
1687 proposal,
1688 ),
1689 None,
1690 )
1691 .await
1692 .unwrap_err();
1693 assert_eq!(err.to_string(), "StorageFailed");
1694
1695 let session = rt.get_session_checked(&sid).await.unwrap();
1697 assert!(!session.seen_message_ids.contains("m2"));
1698 }
1699
1700 #[tokio::test]
1701 async fn cancel_session_fails_if_log_append_fails() {
1702 use std::io;
1703 use std::sync::atomic::{AtomicUsize, Ordering};
1704
1705 struct FailOnSecondAppend {
1706 count: AtomicUsize,
1707 }
1708 #[async_trait::async_trait]
1709 impl StorageBackend for FailOnSecondAppend {
1710 async fn save_session(&self, _: &Session) -> io::Result<()> {
1711 Ok(())
1712 }
1713 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
1714 Ok(None)
1715 }
1716 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
1717 Ok(vec![])
1718 }
1719 async fn delete_session(&self, _: &str) -> io::Result<()> {
1720 Ok(())
1721 }
1722 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
1723 Ok(vec![])
1724 }
1725 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
1726 let n = self.count.fetch_add(1, Ordering::SeqCst);
1727 if n >= 1 {
1728 Err(io::Error::other("disk full"))
1729 } else {
1730 Ok(())
1731 }
1732 }
1733 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
1734 Ok(vec![])
1735 }
1736 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
1737 Ok(())
1738 }
1739 }
1740
1741 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
1742 count: AtomicUsize::new(0),
1743 });
1744 let registry = Arc::new(SessionRegistry::new());
1745 let log_store = Arc::new(LogStore::new());
1746 let rt = Runtime::new(storage, registry, log_store);
1747 let sid = new_sid();
1748
1749 rt.process(
1750 &env(
1751 "macp.mode.decision.v1",
1752 "SessionStart",
1753 "m1",
1754 &sid,
1755 "agent://orchestrator",
1756 session_start(vec!["agent://fraud".into()]),
1757 ),
1758 None,
1759 )
1760 .await
1761 .unwrap();
1762
1763 let err = rt
1764 .cancel_session(&sid, "test cancel", "agent://orchestrator")
1765 .await
1766 .unwrap_err();
1767 assert_eq!(err.to_string(), "StorageFailed");
1768 }
1769
1770 #[tokio::test]
1771 async fn ttl_expiration_rejects_message() {
1772 let rt = make_runtime();
1773 let sid = new_sid();
1774 let payload = SessionStartPayload {
1775 intent: "intent".into(),
1776 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
1777 mode_version: "1.0.0".into(),
1778 configuration_version: "cfg-1".into(),
1779 policy_version: String::new(),
1780 ttl_ms: 1,
1781 context_id: String::new(),
1782 extensions: std::collections::HashMap::new(),
1783 roots: vec![],
1784 }
1785 .encode_to_vec();
1786 rt.process(
1787 &env(
1788 "macp.mode.decision.v1",
1789 "SessionStart",
1790 "m1",
1791 &sid,
1792 "agent://orchestrator",
1793 payload,
1794 ),
1795 None,
1796 )
1797 .await
1798 .unwrap();
1799 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1800 let proposal = ProposalPayload {
1801 proposal_id: "p1".into(),
1802 option: "step-up".into(),
1803 rationale: "risk".into(),
1804 supporting_data: vec![],
1805 }
1806 .encode_to_vec();
1807 let err = rt
1808 .process(
1809 &env(
1810 "macp.mode.decision.v1",
1811 "Proposal",
1812 "m2",
1813 &sid,
1814 "agent://orchestrator",
1815 proposal,
1816 ),
1817 None,
1818 )
1819 .await
1820 .unwrap_err();
1821 assert_eq!(err.to_string(), "TtlExpired");
1822 }
1823
1824 #[tokio::test]
1825 async fn cleanup_expired_sessions_marks_expired() {
1826 let rt = make_runtime();
1827 let sid = new_sid();
1828 let payload = SessionStartPayload {
1829 intent: "intent".into(),
1830 participants: vec!["agent://fraud".into()],
1831 mode_version: "1.0.0".into(),
1832 configuration_version: "cfg-1".into(),
1833 policy_version: String::new(),
1834 ttl_ms: 1,
1835 context_id: String::new(),
1836 extensions: std::collections::HashMap::new(),
1837 roots: vec![],
1838 }
1839 .encode_to_vec();
1840 rt.process(
1841 &env(
1842 "macp.mode.decision.v1",
1843 "SessionStart",
1844 "m1",
1845 &sid,
1846 "agent://orchestrator",
1847 payload,
1848 ),
1849 None,
1850 )
1851 .await
1852 .unwrap();
1853 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1854 rt.cleanup_expired_sessions().await;
1855 let session = rt.get_session_checked(&sid).await.unwrap();
1856 assert_eq!(session.state, SessionState::Expired);
1857 }
1858
1859 #[tokio::test]
1860 async fn evict_stale_sessions_removes_resolved() {
1861 let rt = make_runtime();
1862 let sid = new_sid();
1863 rt.process(
1865 &env(
1866 "macp.mode.decision.v1",
1867 "SessionStart",
1868 "m1",
1869 &sid,
1870 "agent://orchestrator",
1871 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1872 ),
1873 None,
1874 )
1875 .await
1876 .unwrap();
1877 let proposal = ProposalPayload {
1879 proposal_id: "p1".into(),
1880 option: "step-up".into(),
1881 rationale: "risk".into(),
1882 supporting_data: vec![],
1883 }
1884 .encode_to_vec();
1885 rt.process(
1886 &env(
1887 "macp.mode.decision.v1",
1888 "Proposal",
1889 "m2",
1890 &sid,
1891 "agent://orchestrator",
1892 proposal,
1893 ),
1894 None,
1895 )
1896 .await
1897 .unwrap();
1898 let commitment = CommitmentPayload {
1900 commitment_id: "c1".into(),
1901 action: "decision.selected".into(),
1902 authority_scope: "payments".into(),
1903 reason: "bound".into(),
1904 mode_version: "1.0.0".into(),
1905 policy_version: "policy.default".into(),
1906 configuration_version: "cfg-1".into(),
1907 outcome_positive: true,
1908 supersedes: None,
1909 }
1910 .encode_to_vec();
1911 let result = rt
1912 .process(
1913 &env(
1914 "macp.mode.decision.v1",
1915 "Commitment",
1916 "m3",
1917 &sid,
1918 "agent://orchestrator",
1919 commitment,
1920 ),
1921 None,
1922 )
1923 .await
1924 .unwrap();
1925 assert_eq!(result.session_state, SessionState::Resolved);
1926 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1928 rt.evict_stale_sessions(0).await;
1930 assert!(rt.registry.get_session(&sid).await.is_none());
1932 }
1933
1934 #[tokio::test]
1935 async fn session_start_with_wrong_mode_version_rejected() {
1936 let rt = make_runtime();
1937 let sid = new_sid();
1938 let payload = SessionStartPayload {
1939 intent: "test".into(),
1940 participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
1941 mode_version: "99.0.0".into(), configuration_version: "cfg-1".into(),
1943 policy_version: String::new(),
1944 ttl_ms: 60_000,
1945 context_id: String::new(),
1946 extensions: std::collections::HashMap::new(),
1947 roots: vec![],
1948 }
1949 .encode_to_vec();
1950
1951 let err = rt
1952 .process(
1953 &env(
1954 "macp.mode.decision.v1",
1955 "SessionStart",
1956 "m1",
1957 &sid,
1958 "agent://orchestrator",
1959 payload,
1960 ),
1961 None,
1962 )
1963 .await
1964 .unwrap_err();
1965 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
1966 }
1967
1968 #[tokio::test]
1969 async fn signal_empty_signal_type_rejected() {
1970 let rt = make_runtime();
1971 let signal_payload = crate::pb::SignalPayload {
1973 signal_type: String::new(),
1974 data: b"some data".to_vec(),
1975 confidence: 0.0,
1976 correlation_session_id: String::new(),
1977 }
1978 .encode_to_vec();
1979 let signal = Envelope {
1980 macp_version: "1.0".into(),
1981 mode: String::new(),
1982 message_type: "Signal".into(),
1983 message_id: "sig-1".into(),
1984 session_id: String::new(),
1985 sender: "agent://a".into(),
1986 timestamp_unix_ms: 0,
1987 payload: signal_payload,
1988 };
1989 let err = rt.process_signal(&signal).await.unwrap_err();
1990 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
1991 }
1992
1993 #[tokio::test]
1994 async fn signal_valid_payload_accepted() {
1995 let rt = make_runtime();
1996 let signal_payload = crate::pb::SignalPayload {
1997 signal_type: "heartbeat".into(),
1998 data: vec![],
1999 confidence: 0.8,
2000 correlation_session_id: String::new(),
2001 }
2002 .encode_to_vec();
2003 let signal = Envelope {
2004 macp_version: "1.0".into(),
2005 mode: String::new(),
2006 message_type: "Signal".into(),
2007 message_id: "sig-2".into(),
2008 session_id: String::new(),
2009 sender: "agent://a".into(),
2010 timestamp_unix_ms: 0,
2011 payload: signal_payload,
2012 };
2013 rt.process_signal(&signal).await.unwrap();
2014 }
2015
2016 #[tokio::test]
2017 async fn signal_empty_payload_accepted() {
2018 let rt = make_runtime();
2019 let signal = Envelope {
2020 macp_version: "1.0".into(),
2021 mode: String::new(),
2022 message_type: "Signal".into(),
2023 message_id: "sig-3".into(),
2024 session_id: String::new(),
2025 sender: "agent://a".into(),
2026 timestamp_unix_ms: 0,
2027 payload: vec![],
2028 };
2029 rt.process_signal(&signal).await.unwrap();
2030 }
2031}