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, MAX_SUSPENSION_CYCLES,
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 macp_core::MACP_VERSION.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(
298 message_type: &str,
299 payload: &[u8],
300 session_id: &str,
301 mode: &str,
302 at_ms: i64,
303 ) -> LogEntry {
304 LogEntry {
305 message_id: String::new(),
306 received_at_ms: at_ms,
307 sender: "_runtime".into(),
308 message_type: message_type.into(),
309 raw_payload: payload.to_vec(),
310 entry_kind: EntryKind::Internal,
311 session_id: session_id.into(),
312 mode: mode.into(),
313 macp_version: macp_core::MACP_VERSION.into(),
314 timestamp_unix_ms: at_ms,
315 bound_mode_version: None,
316 semantics_rev: 0,
317 bound_max_suspend_ms: None,
318 compacted_incoming_ordinals: 0,
319 }
320 }
321
322 async fn save_session_to_storage(&self, session: &Session) {
323 if let Err(err) = self.storage.save_session(session).await {
324 tracing::warn!(
325 session_id = %session.session_id,
326 error = %err,
327 "failed to persist session snapshot"
328 );
329 }
330 }
331
332 async fn maybe_expire_session(
333 &self,
334 session_id: &str,
335 session: &mut Session,
336 ) -> Result<bool, MacpError> {
337 let now = Utc::now().timestamp_millis();
338 let expires = (session.state == SessionState::Open && now > session.ttl_expiry)
341 || (session.state == SessionState::Suspended && session.suspend_cap_exceeded(now));
342 if expires {
343 let entry =
344 Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode, now);
345 self.storage
346 .append_log_entry(session_id, &entry)
347 .await
348 .map_err(|_| MacpError::StorageFailed)?;
349 self.log_store.append(session_id, entry).await;
350 session.state = SessionState::Expired;
351 session.suspended_at_ms = None;
352 self.metrics.record_session_expired(&session.mode);
353 tracing::info!(session_id, "session expired via TTL");
354 let _ = self
355 .session_lifecycle_bus
356 .send(SessionLifecycleEvent::Expired {
357 session_id: session_id.to_string(),
358 });
359 return Ok(true);
360 }
361 Ok(false)
362 }
363
364 pub async fn process(
365 &self,
366 env: &Envelope,
367 max_open_sessions: Option<usize>,
368 ) -> Result<ProcessResult, MacpError> {
369 match env.message_type.as_str() {
370 "SessionStart" => self.process_session_start(env, max_open_sessions).await,
371 "Signal" | "Progress" => self.process_signal(env).await,
372 _ => self.process_message(env).await,
373 }
374 }
375
376 async fn process_session_start(
377 &self,
378 env: &Envelope,
379 max_open_sessions: Option<usize>,
380 ) -> Result<ProcessResult, MacpError> {
381 if env.mode.trim().is_empty() {
382 return Err(MacpError::InvalidEnvelope);
383 }
384 validate_session_id_for_acceptance(&env.session_id)?;
385 let mode_name = env.mode.as_str();
386 let mode = self
387 .mode_registry
388 .get_mode(mode_name)
389 .ok_or(MacpError::UnknownMode)?;
390
391 let start_payload = parse_session_start_payload(&env.payload)?;
392 let require_complete_start = self.mode_registry.requires_strict_session_start(mode_name);
400 if require_complete_start {
401 validate_canonical_session_start_payload_for_mode(mode_name, &start_payload)?;
402 }
403
404 let descriptor_version = self.mode_registry.get_mode_version(mode_name);
412 if let Some(descriptor_version) = &descriptor_version {
413 if !start_payload.mode_version.is_empty()
414 && &start_payload.mode_version != descriptor_version
415 {
416 tracing::warn!(
417 mode = mode_name,
418 payload_version = %start_payload.mode_version,
419 descriptor_version = %descriptor_version,
420 "mode_version mismatch"
421 );
422 return Err(MacpError::InvalidEnvelope);
423 }
424 }
425 let bound_mode_version: Option<String> = if start_payload.mode_version.is_empty() {
426 descriptor_version
427 } else {
428 None
429 };
430 let effective_mode_version = bound_mode_version
431 .clone()
432 .unwrap_or_else(|| start_payload.mode_version.clone());
433
434 let ttl_ms = extract_ttl_ms(&start_payload)?;
435
436 if let Some(existing) = self.registry.get_shared(&env.session_id).await {
441 let existing = existing.lock().await;
442 if existing.seen_message_ids.contains(&env.message_id) {
443 return Ok(ProcessResult {
444 session_state: existing.state.clone(),
445 duplicate: true,
446 });
447 }
448 return Err(MacpError::SessionAlreadyExists);
449 }
450
451 let effective_policy_version = if start_payload.policy_version.is_empty() {
456 crate::policy::defaults::DEFAULT_POLICY_ID.to_string()
457 } else {
458 start_payload.policy_version.clone()
459 };
460 let policy_definition = match self.policy_registry.resolve(&effective_policy_version) {
461 Ok(policy) => {
462 if policy.mode != "*" && policy.mode != mode_name {
464 return Err(MacpError::InvalidPolicyDefinition);
465 }
466 Some(policy)
467 }
468 Err(_) => {
469 return Err(MacpError::UnknownPolicyVersion);
470 }
471 };
472
473 let accepted_at = Utc::now().timestamp_millis();
474 let ttl_base = if env.timestamp_unix_ms > 0 {
478 env.timestamp_unix_ms
479 } else {
480 accepted_at
481 };
482 let ttl_expiry = ttl_base.saturating_add(ttl_ms);
483 let bound_max_suspend_ms = if start_payload.max_suspend_ms > 0 {
488 start_payload.max_suspend_ms
489 } else {
490 macp_core::session::MAX_SUSPEND_MS
491 };
492 let session = Session::builder(env.session_id.clone(), mode_name, env.sender.clone())
493 .ttl_expiry(ttl_expiry)
494 .ttl_ms(ttl_ms)
495 .max_suspend_ms(bound_max_suspend_ms)
496 .started_at_unix_ms(accepted_at)
497 .participants(start_payload.participants.clone())
498 .intent(start_payload.intent.clone())
499 .mode_version(effective_mode_version)
500 .configuration_version(start_payload.configuration_version.clone())
501 .policy_version(effective_policy_version)
502 .context_id(start_payload.context_id.clone())
503 .extensions(start_payload.extensions.clone())
504 .roots(start_payload.roots.clone())
505 .policy_definition(policy_definition)
506 .build();
507
508 let response = mode.on_session_start(&session, env)?;
509 mode.validate_client_envelope(&session, env)?;
516 let semantics_rev = session.semantics_rev;
517
518 let shared = std::sync::Arc::new(tokio::sync::Mutex::new(session));
523 let mut session_guard = shared
526 .clone()
527 .try_lock_owned()
528 .expect("freshly created mutex is uncontended");
529 {
530 let mut map = self.registry.sessions.write().await;
531 if map.contains_key(&env.session_id) {
532 return Err(MacpError::SessionAlreadyExists);
534 }
535 if let Some(max_open) = max_open_sessions {
536 let now = Utc::now().timestamp_millis();
537 let mut count = 0usize;
538 for arc in map.values() {
539 let counts = match arc.try_lock() {
544 Ok(s) => {
545 s.initiator_sender == env.sender
546 && s.state == SessionState::Open
547 && now <= s.ttl_expiry
548 }
549 Err(_) => true,
550 };
551 if counts {
552 count += 1;
553 }
554 }
555 if count >= max_open {
556 return Err(MacpError::RateLimited);
557 }
558 }
559 map.insert(env.session_id.clone(), std::sync::Arc::clone(&shared));
560 }
561
562 let rollback = |runtime: &Self, session_guard: &mut Session| {
567 session_guard.state = SessionState::Expired;
568 let registry = std::sync::Arc::clone(&runtime.registry);
569 let sid = env.session_id.clone();
570 async move {
571 let mut map = registry.sessions.write().await;
572 map.remove(&sid);
573 }
574 };
575
576 if self
578 .storage
579 .create_session_storage(&env.session_id)
580 .await
581 .is_err()
582 {
583 rollback(self, &mut session_guard).await;
584 return Err(MacpError::StorageFailed);
585 }
586 let mut incoming_entry = Self::make_incoming_entry(env, accepted_at);
587 incoming_entry.bound_mode_version = bound_mode_version;
588 incoming_entry.semantics_rev = semantics_rev;
589 incoming_entry.bound_max_suspend_ms = Some(bound_max_suspend_ms);
590 if self
591 .storage
592 .append_log_entry(&env.session_id, &incoming_entry)
593 .await
594 .is_err()
595 {
596 rollback(self, &mut session_guard).await;
597 return Err(MacpError::StorageFailed);
598 }
599
600 self.log_store.create_session_log(&env.session_id).await;
602 self.log_store.append(&env.session_id, incoming_entry).await;
603
604 session_guard
605 .seen_message_ids
606 .insert(env.message_id.clone());
607 session_guard.apply_mode_response(response);
608
609 let result_state = session_guard.state.clone();
610 if let Err(err) = self.storage.save_session(&session_guard).await {
619 tracing::warn!(
620 session_id = %session_guard.session_id,
621 error = %err,
622 "failed to persist session snapshot at SessionStart (recoverable via replay)"
623 );
624 }
625 self.metrics.record_session_start(mode_name);
626 tracing::info!(
627 session_id = %env.session_id,
628 mode = mode_name,
629 sender = %env.sender,
630 "session started"
631 );
632 self.publish_accepted_envelope(env);
638 drop(session_guard);
639 let _ = self
640 .session_lifecycle_bus
641 .send(SessionLifecycleEvent::Created {
642 session_id: env.session_id.clone(),
643 });
644
645 Ok(ProcessResult {
646 session_state: result_state,
647 duplicate: false,
648 })
649 }
650
651 async fn synthesize_due_accept(
749 &self,
750 session_id: &str,
751 session: &mut Session,
752 now_ms: i64,
753 ) -> Result<bool, MacpError> {
754 if session.state != SessionState::Open {
755 return Ok(false);
756 }
757 let Some(mode) = self.mode_registry.get_mode(&session.mode) else {
758 return Ok(false);
759 };
760 let Some(syn) = mode.due_synthetic_envelope(session, now_ms) else {
761 return Ok(false);
762 };
763 if session.seen_message_ids.contains(&syn.message_id) {
767 return Ok(false);
768 }
769 mode.authorize_sender(session, &syn)?;
773 let response = mode.on_message_at(
774 session,
775 &syn,
776 &macp_core::mode::MessageContext::new(syn.timestamp_unix_ms),
777 )?;
778
779 let entry = Self::make_incoming_entry(&syn, syn.timestamp_unix_ms);
785 self.storage
786 .append_log_entry(session_id, &entry)
787 .await
788 .map_err(|_| MacpError::StorageFailed)?;
789 self.log_store.append(session_id, entry).await;
790
791 session.seen_message_ids.insert(syn.message_id.clone());
792 session.apply_mode_response(response);
793 self.metrics.record_message_accepted(&session.mode);
794
795 tracing::info!(
796 session_id = %session_id,
797 message_type = %syn.message_type,
798 message_id = %syn.message_id,
799 sender = %syn.sender,
800 deadline_ms = syn.timestamp_unix_ms,
801 "synthetic envelope appended to accepted history"
802 );
803
804 self.save_session_to_storage(session).await;
805 self.maybe_insert_checkpoint(session_id, session).await;
822 self.publish_accepted_envelope(&syn);
823 Ok(true)
824 }
825
826 async fn process_message(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
834 let shared = self
841 .registry
842 .get_shared(&env.session_id)
843 .await
844 .ok_or(MacpError::UnknownSession)?;
845 let mut session_guard = shared.lock().await;
846 let session = &mut *session_guard;
847
848 let now_ms = chrono::Utc::now().timestamp_millis();
855 match macp_modes::step::check_preconditions(session, env, now_ms)? {
856 macp_modes::step::Precheck::Duplicate => {
857 return Ok(ProcessResult {
858 session_state: session.state.clone(),
859 duplicate: true,
860 });
861 }
862 macp_modes::step::Precheck::Expired => {
863 let expired = self.maybe_expire_session(&env.session_id, session).await?;
869 debug_assert!(expired, "check_preconditions reported Expired");
870 self.save_session_to_storage(session).await;
871 return Err(MacpError::TtlExpired);
872 }
873 macp_modes::step::Precheck::Proceed => {}
874 }
875
876 let mode = self
877 .mode_registry
878 .get_mode(&session.mode)
879 .ok_or(MacpError::UnknownMode)?;
880 mode.authorize_sender(session, env)?;
881 mode.validate_client_envelope(session, env)?;
893 let accepted_at_ms = Utc::now().timestamp_millis();
896 self.synthesize_due_accept(&env.session_id, session, accepted_at_ms)
904 .await?;
905 let response = mode.on_message_at(
906 session,
907 env,
908 &macp_core::mode::MessageContext::new(accepted_at_ms),
909 )?;
910
911 let incoming_entry = Self::make_incoming_entry(env, accepted_at_ms);
913 self.storage
914 .append_log_entry(&env.session_id, &incoming_entry)
915 .await
916 .map_err(|_| MacpError::StorageFailed)?;
917
918 self.log_store.append(&env.session_id, incoming_entry).await;
922 let result_state = macp_modes::step::commit(session, env, response, now_ms);
923
924 self.metrics.record_message_accepted(&session.mode);
925 if env.message_type == "Commitment" {
926 self.metrics.record_commitment_accepted(&session.mode);
927 }
928
929 if Self::audit_verbose(session) {
934 tracing::info!(
935 session_id = %env.session_id,
936 message_type = %env.message_type,
937 sender = %env.sender,
938 state = ?result_state,
939 "message accepted (audit)"
940 );
941 } else {
942 tracing::debug!(
943 session_id = %env.session_id,
944 message_type = %env.message_type,
945 sender = %env.sender,
946 state = ?result_state,
947 "message accepted"
948 );
949 }
950
951 if result_state == SessionState::Resolved {
952 self.metrics.record_session_resolved(&session.mode);
953 tracing::info!(session_id = %env.session_id, mode = %session.mode, "session resolved");
954 let _ = self
955 .session_lifecycle_bus
956 .send(SessionLifecycleEvent::Resolved {
957 session_id: env.session_id.clone(),
958 });
959 }
960
961 self.save_session_to_storage(session).await;
963 if result_state == SessionState::Resolved {
964 if !self.maybe_compact_log(&env.session_id, session).await {
965 self.force_insert_checkpoint(&env.session_id, session).await;
966 }
967 } else {
968 self.maybe_insert_checkpoint(&env.session_id, session).await;
969 }
970 self.publish_accepted_envelope(env);
971
972 Ok(ProcessResult {
973 session_state: result_state,
974 duplicate: false,
975 })
976 }
977
978 async fn process_signal(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
982 if env.message_type == "Signal" && !env.payload.is_empty() {
985 let signal: crate::pb::SignalPayload =
986 prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
987 if signal.signal_type.trim().is_empty() {
988 return Err(MacpError::InvalidPayload);
989 }
990 }
991 if env.message_type == "Progress" && !env.payload.is_empty() {
993 let _: crate::pb::ProgressPayload =
994 prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
995 }
996 tracing::debug!(
997 sender = %env.sender,
998 message_id = %env.message_id,
999 message_type = %env.message_type,
1000 "signal received"
1001 );
1002 let _ = self.signal_bus.send(env.clone());
1003 Ok(ProcessResult {
1004 session_state: SessionState::Open,
1005 duplicate: false,
1006 })
1007 }
1008
1009 pub async fn get_session_checked(&self, session_id: &str) -> Option<Session> {
1010 let shared = self.registry.get_shared(session_id).await?;
1011 let mut session = shared.lock().await;
1012 let changed = self
1013 .maybe_expire_session(session_id, &mut session)
1014 .await
1015 .unwrap_or(false);
1016 if changed {
1017 self.save_session_to_storage(&session).await;
1018 }
1019 Some(session.clone())
1020 }
1021
1022 pub async fn cancel_session(
1026 &self,
1027 session_id: &str,
1028 reason: &str,
1029 cancelled_by: &str,
1030 ) -> Result<ProcessResult, MacpError> {
1031 let shared = self
1032 .registry
1033 .get_shared(session_id)
1034 .await
1035 .ok_or(MacpError::UnknownSession)?;
1036 let mut session_guard = shared.lock().await;
1037 let session = &mut *session_guard;
1038
1039 self.maybe_expire_session(session_id, session).await?;
1040
1041 if session.state.is_terminal() {
1044 let result_state = session.state.clone();
1045 self.save_session_to_storage(session).await;
1046 return Ok(ProcessResult {
1047 session_state: result_state,
1048 duplicate: false,
1049 });
1050 }
1051
1052 let now_ms = Utc::now().timestamp_millis();
1055 let cancel_payload = crate::pb::SessionCancelPayload {
1056 reason: reason.to_string(),
1057 cancelled_by: cancelled_by.to_string(),
1058 };
1059 let cancel_entry = Self::make_internal_entry(
1062 "SessionCancel",
1063 &prost::Message::encode_to_vec(&cancel_payload),
1064 session_id,
1065 &session.mode,
1066 now_ms,
1067 );
1068 self.storage
1069 .append_log_entry(session_id, &cancel_entry)
1070 .await
1071 .map_err(|_| MacpError::StorageFailed)?;
1072 self.log_store.append(session_id, cancel_entry).await;
1073 let _ = session.cancel();
1076 self.save_session_to_storage(session).await;
1077 if !self.maybe_compact_log(session_id, session).await {
1078 self.force_insert_checkpoint(session_id, session).await;
1079 }
1080 self.metrics.record_session_cancelled(&session.mode);
1081 tracing::info!(session_id, reason, "session cancelled");
1082 let _ = self
1083 .session_lifecycle_bus
1084 .send(SessionLifecycleEvent::Cancelled {
1085 session_id: session_id.to_string(),
1086 });
1087
1088 Ok(ProcessResult {
1089 session_state: SessionState::Cancelled,
1090 duplicate: false,
1091 })
1092 }
1093
1094 pub async fn suspend_session(
1098 &self,
1099 session_id: &str,
1100 reason: &str,
1101 suspended_by: &str,
1102 ) -> Result<ProcessResult, MacpError> {
1103 let shared = self
1104 .registry
1105 .get_shared(session_id)
1106 .await
1107 .ok_or(MacpError::UnknownSession)?;
1108 let mut session_guard = shared.lock().await;
1109 let session = &mut *session_guard;
1110
1111 self.maybe_expire_session(session_id, session).await?;
1112 if session.state != SessionState::Open {
1113 return Err(MacpError::SessionNotOpen);
1114 }
1115
1116 let now_ms = chrono::Utc::now().timestamp_millis();
1117 let payload = crate::pb::SessionSuspendPayload {
1118 reason: reason.to_string(),
1119 suspended_by: suspended_by.to_string(),
1120 };
1121 let entry = Self::make_internal_entry(
1124 "SessionSuspend",
1125 &prost::Message::encode_to_vec(&payload),
1126 session_id,
1127 &session.mode,
1128 now_ms,
1129 );
1130 self.storage
1131 .append_log_entry(session_id, &entry)
1132 .await
1133 .map_err(|_| MacpError::StorageFailed)?;
1134 self.log_store.append(session_id, entry).await;
1135 session.suspend(now_ms)?;
1136 self.save_session_to_storage(session).await;
1137 self.metrics.record_session_suspended(&session.mode);
1138 tracing::info!(session_id, reason, "session suspended");
1139 let _ = self
1140 .session_lifecycle_bus
1141 .send(SessionLifecycleEvent::Suspended {
1142 session_id: session_id.to_string(),
1143 });
1144
1145 Ok(ProcessResult {
1146 session_state: SessionState::Suspended,
1147 duplicate: false,
1148 })
1149 }
1150
1151 pub async fn resume_session(
1155 &self,
1156 session_id: &str,
1157 reason: &str,
1158 resumed_by: &str,
1159 ) -> Result<ProcessResult, MacpError> {
1160 let shared = self
1161 .registry
1162 .get_shared(session_id)
1163 .await
1164 .ok_or(MacpError::UnknownSession)?;
1165 let mut session_guard = shared.lock().await;
1166 let session = &mut *session_guard;
1167
1168 if session.state != SessionState::Suspended {
1169 return Err(MacpError::SessionNotOpen);
1170 }
1171
1172 let now_ms = chrono::Utc::now().timestamp_millis();
1173 let banked_ms = session
1206 .suspended_at_ms
1207 .map(|suspended_at| session.ttl_expiry.saturating_sub(suspended_at).max(0))
1208 .unwrap_or(0);
1209 let payload = crate::pb::SessionResumePayload {
1210 reason: reason.to_string(),
1211 resumed_by: resumed_by.to_string(),
1212 banked_ms,
1213 };
1214 let entry = Self::make_internal_entry(
1217 "SessionResume",
1218 &prost::Message::encode_to_vec(&payload),
1219 session_id,
1220 &session.mode,
1221 now_ms,
1222 );
1223 self.storage
1224 .append_log_entry(session_id, &entry)
1225 .await
1226 .map_err(|_| MacpError::StorageFailed)?;
1227 self.log_store.append(session_id, entry).await;
1228
1229 match session.resume(now_ms) {
1231 Ok(()) => {
1232 self.save_session_to_storage(session).await;
1233 self.metrics.record_session_resumed(&session.mode);
1234 tracing::info!(session_id, reason, "session resumed");
1235 let _ = self
1236 .session_lifecycle_bus
1237 .send(SessionLifecycleEvent::Resumed {
1238 session_id: session_id.to_string(),
1239 });
1240 Ok(ProcessResult {
1241 session_state: SessionState::Open,
1242 duplicate: false,
1243 })
1244 }
1245 Err(_) => {
1246 let cycle_cap_exceeded = session.suspension_intervals.len() > MAX_SUSPENSION_CYCLES;
1255 let duration_cap_exceeded =
1256 session.accumulated_suspended_ms > session.effective_max_suspend_ms();
1257 tracing::warn!(
1258 session_id,
1259 cycle_cap_exceeded,
1260 duration_cap_exceeded,
1261 suspension_cycles = session.suspension_intervals.len(),
1262 accumulated_suspended_ms = session.accumulated_suspended_ms,
1263 "session force-expired: suspension cap exceeded"
1264 );
1265 self.save_session_to_storage(session).await;
1266 self.metrics.record_session_expired(&session.mode);
1267 let _ = self
1268 .session_lifecycle_bus
1269 .send(SessionLifecycleEvent::Expired {
1270 session_id: session_id.to_string(),
1271 });
1272 Err(MacpError::TtlExpired)
1273 }
1274 }
1275 }
1276
1277 async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
1280 let discarded = match self.log_store.get_log(session_id).await {
1284 Some(entries) => {
1285 let prior_base: u64 = entries
1286 .iter()
1287 .filter(|e| e.entry_kind == EntryKind::Checkpoint)
1288 .map(|e| e.compacted_incoming_ordinals)
1289 .max()
1290 .unwrap_or(0);
1291 prior_base
1292 + entries
1293 .iter()
1294 .filter(|e| e.entry_kind == EntryKind::Incoming)
1295 .count() as u64
1296 }
1297 None => 0,
1298 };
1299 match crate::storage::compaction::compact_session_log(
1300 &*self.storage,
1301 session_id,
1302 session,
1303 discarded,
1304 )
1305 .await
1306 {
1307 Ok(checkpoint) => {
1308 self.log_store
1312 .replace_session_log(session_id, vec![checkpoint])
1313 .await;
1314 true
1315 }
1316 Err(e) => {
1317 tracing::debug!(
1318 session_id,
1319 error = %e,
1320 "log compaction skipped (backend may not support it)"
1321 );
1322 false
1323 }
1324 }
1325 }
1326
1327 async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
1330 let persisted = crate::registry::PersistedSession::from(session);
1331 let raw_payload = match serde_json::to_vec(&persisted) {
1332 Ok(bytes) => bytes,
1333 Err(e) => {
1334 tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
1335 return;
1336 }
1337 };
1338 let now = Utc::now().timestamp_millis();
1339 let checkpoint = LogEntry {
1340 message_id: String::new(),
1341 received_at_ms: now,
1342 sender: "_runtime".into(),
1343 message_type: "Checkpoint".into(),
1344 raw_payload,
1345 entry_kind: EntryKind::Checkpoint,
1346 session_id: session_id.into(),
1347 mode: session.mode.clone(),
1348 macp_version: String::new(),
1349 timestamp_unix_ms: now,
1350 bound_mode_version: None,
1351 semantics_rev: 0,
1352 bound_max_suspend_ms: None,
1353 compacted_incoming_ordinals: 0,
1354 };
1355 if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
1356 tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
1357 return;
1358 }
1359 self.log_store.append(session_id, checkpoint).await;
1360 tracing::debug!(
1361 session_id,
1362 "forced checkpoint inserted for terminal session"
1363 );
1364 }
1365
1366 async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
1376 if self.checkpoint_interval == 0 {
1377 return;
1378 }
1379 let log_len = self
1380 .log_store
1381 .get_log(session_id)
1382 .await
1383 .map(|l| l.len())
1384 .unwrap_or(0);
1385 if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
1387 return;
1388 }
1389 self.force_insert_checkpoint(session_id, session).await;
1390 tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
1391 }
1392
1393 pub async fn cleanup_expired_sessions(&self) {
1397 let now = Utc::now().timestamp_millis();
1398 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1403 let guard = self.registry.sessions.read().await;
1404 guard
1405 .iter()
1406 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1407 .collect()
1408 };
1409
1410 let mut expired_count = 0usize;
1411 for (session_id, shared) in candidates {
1412 let mut session = shared.lock().await;
1413 if session.state != SessionState::Open || now <= session.ttl_expiry {
1414 continue;
1415 }
1416 let entry =
1417 Self::make_internal_entry("TtlExpired", b"", &session_id, &session.mode, now);
1418 if let Err(e) = self.storage.append_log_entry(&session_id, &entry).await {
1419 tracing::warn!(
1420 session_id,
1421 error = %e,
1422 "failed to write TTL expiry during cleanup"
1423 );
1424 continue;
1425 }
1426 self.log_store.append(&session_id, entry).await;
1427 session.state = SessionState::Expired;
1428 self.metrics.record_session_expired(&session.mode);
1429 self.save_session_to_storage(&session).await;
1430 if !self.maybe_compact_log(&session_id, &session).await {
1431 self.force_insert_checkpoint(&session_id, &session).await;
1432 }
1433 expired_count += 1;
1434 tracing::info!(session_id = %session_id, "session expired via background cleanup");
1435 let _ = self
1436 .session_lifecycle_bus
1437 .send(SessionLifecycleEvent::Expired {
1438 session_id: session_id.clone(),
1439 });
1440 }
1441
1442 if expired_count > 0 {
1443 tracing::info!(count = expired_count, "background cleanup expired sessions");
1444 }
1445 }
1446
1447 pub async fn sweep_due_synthetic_accepts(&self) -> usize {
1507 let now = Utc::now().timestamp_millis();
1508 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1512 let guard = self.registry.sessions.read().await;
1513 guard
1514 .iter()
1515 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1516 .collect()
1517 };
1518
1519 let mut emitted = 0usize;
1520 for (session_id, shared) in candidates {
1521 let mut session = shared.lock().await;
1522 if session.state != SessionState::Open {
1533 continue;
1534 }
1535 match self
1536 .synthesize_due_accept(&session_id, &mut session, now)
1537 .await
1538 {
1539 Ok(true) => emitted += 1,
1540 Ok(false) => {}
1541 Err(e) => {
1547 tracing::warn!(
1548 session_id = %session_id,
1549 error = %e,
1550 "eager sweep could not emit a due synthetic envelope"
1551 );
1552 }
1553 }
1554 }
1555
1556 if emitted > 0 {
1557 tracing::info!(
1558 count = emitted,
1559 "eager sweep appended due synthetic envelopes"
1560 );
1561 }
1562 emitted
1563 }
1564
1565 pub async fn gc_disk_sessions(&self, retention_secs: u64) -> usize {
1573 let now = Utc::now().timestamp_millis();
1574 let cutoff = now - (retention_secs as i64 * 1000);
1575 let ids = match self.storage.list_session_ids().await {
1576 Ok(ids) => ids,
1577 Err(e) => {
1578 tracing::warn!(error = %e, "disk GC: cannot list sessions");
1579 return 0;
1580 }
1581 };
1582 let mut removed = 0usize;
1583 for id in ids {
1584 let eligible = if let Some(shared) = self.registry.get_shared(&id).await {
1587 let s = shared.lock().await;
1588 s.state.is_terminal() && s.started_at_unix_ms < cutoff
1589 } else {
1590 match self.storage.load_session(&id).await {
1591 Ok(Some(s)) => s.state.is_terminal() && s.started_at_unix_ms < cutoff,
1592 _ => false,
1595 }
1596 };
1597 if !eligible {
1598 continue;
1599 }
1600 match self.storage.delete_session(&id).await {
1601 Ok(()) => {
1602 {
1603 let mut guard = self.registry.sessions.write().await;
1604 guard.remove(&id);
1605 }
1606 self.log_store.remove_session_log(&id).await;
1607 let _ = self.stream_bus.remove_if_unused(&id);
1608 removed += 1;
1609 }
1610 Err(e) => {
1611 tracing::warn!(session_id = %id, error = %e, "disk GC: delete failed");
1612 }
1613 }
1614 }
1615 if removed > 0 {
1616 tracing::info!(count = removed, "disk GC removed terminal sessions");
1617 }
1618 removed
1619 }
1620
1621 pub async fn evict_stale_sessions(&self, retention_secs: u64) {
1627 let now = Utc::now().timestamp_millis();
1628 let cutoff = now - (retention_secs as i64 * 1000);
1629
1630 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1631 let guard = self.registry.sessions.read().await;
1632 guard
1633 .iter()
1634 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1635 .collect()
1636 };
1637 let mut evict_ids = Vec::new();
1638 for (id, shared) in candidates {
1639 let session = shared.lock().await;
1640 if matches!(
1641 session.state,
1642 SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
1643 ) && session.started_at_unix_ms < cutoff
1644 {
1645 evict_ids.push(id);
1646 }
1647 }
1648
1649 if evict_ids.is_empty() {
1650 return;
1651 }
1652 {
1653 let mut guard = self.registry.sessions.write().await;
1654 for id in &evict_ids {
1655 guard.remove(id);
1656 }
1657 }
1658 for id in &evict_ids {
1659 self.log_store.remove_session_log(id).await;
1660 let _ = self.stream_bus.remove_if_unused(id);
1663 }
1664 tracing::info!(
1665 count = evict_ids.len(),
1666 "evicted stale sessions from memory (registry + log cache + stream bus)"
1667 );
1668 }
1669}
1670
1671#[cfg(test)]
1672mod tests {
1673 use super::*;
1674 use crate::decision_pb::ProposalPayload;
1675 use crate::pb::{CommitmentPayload, SessionStartPayload};
1676 use prost::Message;
1677
1678 fn new_sid() -> String {
1679 uuid::Uuid::new_v4().as_hyphenated().to_string()
1680 }
1681
1682 fn make_runtime() -> Runtime {
1683 let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
1684 let registry = Arc::new(SessionRegistry::new());
1685 let log_store = Arc::new(LogStore::new());
1686 Runtime::new(storage, registry, log_store)
1687 }
1688
1689 fn session_start(participants: Vec<String>) -> Vec<u8> {
1690 SessionStartPayload {
1691 intent: "intent".into(),
1692 participants,
1693 mode_version: "1.0.0".into(),
1694 configuration_version: "cfg-1".into(),
1695 policy_version: String::new(),
1696 ttl_ms: 1_000,
1697 context_id: String::new(),
1698 extensions: std::collections::HashMap::new(),
1699 roots: vec![],
1700 max_suspend_ms: 0,
1701 }
1702 .encode_to_vec()
1703 }
1704
1705 fn env(
1706 mode: &str,
1707 message_type: &str,
1708 message_id: &str,
1709 session_id: &str,
1710 sender: &str,
1711 payload: Vec<u8>,
1712 ) -> Envelope {
1713 Envelope {
1714 macp_version: "1.0".into(),
1715 mode: mode.into(),
1716 message_type: message_type.into(),
1717 message_id: message_id.into(),
1718 session_id: session_id.into(),
1719 sender: sender.into(),
1720 timestamp_unix_ms: Utc::now().timestamp_millis(),
1721 payload,
1722 }
1723 }
1724
1725 #[tokio::test]
1726 async fn standard_session_start_is_strict() {
1727 let rt = make_runtime();
1728 let sid = new_sid();
1729 let bad = SessionStartPayload {
1730 ttl_ms: 0,
1731 ..Default::default()
1732 }
1733 .encode_to_vec();
1734 let err = rt
1735 .process(
1736 &env(
1737 "macp.mode.decision.v1",
1738 "SessionStart",
1739 "m1",
1740 &sid,
1741 "agent://orchestrator",
1742 bad,
1743 ),
1744 None,
1745 )
1746 .await
1747 .unwrap_err();
1748 assert!(matches!(
1749 err,
1750 MacpError::InvalidPayload | MacpError::InvalidTtl
1751 ));
1752 }
1753
1754 #[tokio::test]
1767 async fn a_promoted_mode_still_gets_canonical_session_start_validation() {
1768 let mode_registry = Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
1769 macp_policy::DefaultPolicyEvaluator,
1770 )));
1771 mode_registry
1772 .register_extension(crate::pb::ModeDescriptor {
1773 mode: "ext.promoted.v1".into(),
1774 mode_version: "1.0.0".into(),
1775 title: "Promoted".into(),
1776 description: "promotion target".into(),
1777 determinism_class: "semantic-deterministic".into(),
1778 participant_model: "declared".into(),
1779 message_types: vec!["SessionStart".into(), "Commitment".into()],
1780 terminal_message_types: vec!["Commitment".into()],
1781 ..Default::default()
1782 })
1783 .expect("register extension");
1784 assert_eq!(
1785 mode_registry.promote_mode("ext.promoted.v1", None).unwrap(),
1786 "ext.promoted.v1"
1787 );
1788 assert!(
1789 mode_registry.requires_strict_session_start("ext.promoted.v1"),
1790 "promotion must mark the entry strict"
1791 );
1792 assert!(
1793 !crate::session::requires_strict_session_start("ext.promoted.v1"),
1794 "the core's static list must NOT know this name — that disagreement is the point"
1795 );
1796
1797 let rt = Runtime::with_mode_registry(
1798 Arc::new(crate::storage::MemoryBackend),
1799 Arc::new(SessionRegistry::new()),
1800 Arc::new(LogStore::new()),
1801 mode_registry,
1802 );
1803
1804 let err = rt
1806 .process(
1807 &env(
1808 "ext.promoted.v1",
1809 "SessionStart",
1810 "m1",
1811 &new_sid(),
1812 "agent://orchestrator",
1813 session_start(vec![]),
1814 ),
1815 None,
1816 )
1817 .await
1818 .unwrap_err();
1819 assert_eq!(err.to_string(), "InvalidPayload");
1820
1821 let no_versions = SessionStartPayload {
1824 participants: vec!["agent://fraud".into()],
1825 ttl_ms: 1_000,
1826 ..Default::default()
1827 }
1828 .encode_to_vec();
1829 let err = rt
1830 .process(
1831 &env(
1832 "ext.promoted.v1",
1833 "SessionStart",
1834 "m2",
1835 &new_sid(),
1836 "agent://orchestrator",
1837 no_versions,
1838 ),
1839 None,
1840 )
1841 .await
1842 .unwrap_err();
1843 assert_eq!(err.to_string(), "InvalidPayload");
1844
1845 rt.process(
1848 &env(
1849 "ext.promoted.v1",
1850 "SessionStart",
1851 "m3",
1852 &new_sid(),
1853 "agent://orchestrator",
1854 session_start(vec!["agent://fraud".into()]),
1855 ),
1856 None,
1857 )
1858 .await
1859 .expect("a complete SessionStart must still be accepted for a promoted mode");
1860 }
1861
1862 #[tokio::test]
1863 async fn empty_mode_is_rejected() {
1864 let rt = make_runtime();
1865 let sid = new_sid();
1866 let err = rt
1867 .process(
1868 &env(
1869 "",
1870 "SessionStart",
1871 "m1",
1872 &sid,
1873 "agent://orchestrator",
1874 session_start(vec!["agent://fraud".into()]),
1875 ),
1876 None,
1877 )
1878 .await
1879 .unwrap_err();
1880 assert_eq!(err.to_string(), "InvalidEnvelope");
1881 }
1882
1883 #[tokio::test]
1884 async fn rejected_messages_do_not_enter_dedup_state() {
1885 let rt = make_runtime();
1886 let sid = new_sid();
1887 rt.process(
1888 &env(
1889 "macp.mode.decision.v1",
1890 "SessionStart",
1891 "m1",
1892 &sid,
1893 "agent://orchestrator",
1894 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1895 ),
1896 None,
1897 )
1898 .await
1899 .unwrap();
1900
1901 let bad = rt
1902 .process(
1903 &env(
1904 "macp.mode.decision.v1",
1905 "Proposal",
1906 "m2",
1907 &sid,
1908 "agent://fraud",
1909 b"not-protobuf".to_vec(),
1910 ),
1911 None,
1912 )
1913 .await
1914 .unwrap_err();
1915 assert_eq!(bad.to_string(), "InvalidPayload");
1916
1917 let good = ProposalPayload {
1918 proposal_id: "p1".into(),
1919 option: "step-up".into(),
1920 rationale: "risk".into(),
1921 supporting_data: vec![],
1922 }
1923 .encode_to_vec();
1924 let result = rt
1925 .process(
1926 &env(
1927 "macp.mode.decision.v1",
1928 "Proposal",
1929 "m2",
1930 &sid,
1931 "agent://orchestrator",
1932 good,
1933 ),
1934 None,
1935 )
1936 .await
1937 .unwrap();
1938 assert!(!result.duplicate);
1939 }
1940
1941 #[tokio::test]
1942 async fn get_session_transitions_expired_sessions() {
1943 let rt = make_runtime();
1944 let sid = new_sid();
1945 let payload = SessionStartPayload {
1946 intent: "intent".into(),
1947 participants: vec!["agent://fraud".into()],
1948 mode_version: "1.0.0".into(),
1949 configuration_version: "cfg-1".into(),
1950 policy_version: String::new(),
1951 ttl_ms: 1,
1952 context_id: String::new(),
1953 extensions: std::collections::HashMap::new(),
1954 roots: vec![],
1955 max_suspend_ms: 0,
1956 }
1957 .encode_to_vec();
1958 rt.process(
1959 &env(
1960 "macp.mode.decision.v1",
1961 "SessionStart",
1962 "m1",
1963 &sid,
1964 "agent://orchestrator",
1965 payload,
1966 ),
1967 None,
1968 )
1969 .await
1970 .unwrap();
1971 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1972 let session = rt.get_session_checked(&sid).await.unwrap();
1973 assert_eq!(session.state, SessionState::Expired);
1974 }
1975
1976 #[tokio::test]
1977 async fn multi_round_requires_standard_session_start() {
1978 let rt = make_runtime();
1979 let sid = new_sid();
1980 let payload = SessionStartPayload {
1982 participants: vec!["creator".into(), "other".into()],
1983 ..Default::default()
1984 }
1985 .encode_to_vec();
1986 let err = rt
1987 .process(
1988 &env(
1989 "ext.multi_round.v1",
1990 "SessionStart",
1991 "m1",
1992 &sid,
1993 "creator",
1994 payload,
1995 ),
1996 None,
1997 )
1998 .await
1999 .unwrap_err();
2000 assert!(matches!(
2001 err,
2002 MacpError::InvalidPayload | MacpError::InvalidTtl
2003 ));
2004 }
2005
2006 #[tokio::test]
2007 async fn multi_round_valid_session_start() {
2008 let rt = make_runtime();
2009 let sid = new_sid();
2010 let payload = session_start(vec!["alice".into(), "bob".into()]);
2011 rt.process(
2012 &env(
2013 "ext.multi_round.v1",
2014 "SessionStart",
2015 "m1",
2016 &sid,
2017 "coordinator",
2018 payload,
2019 ),
2020 None,
2021 )
2022 .await
2023 .unwrap();
2024 let session = rt.get_session_checked(&sid).await.unwrap();
2025 assert_eq!(session.mode, "ext.multi_round.v1");
2026 assert_eq!(session.participants, vec!["alice", "bob"]);
2027 }
2028
2029 #[tokio::test]
2030 async fn duplicate_session_start_message_id_returns_duplicate() {
2031 let rt = make_runtime();
2032 let sid = new_sid();
2033 let payload = session_start(vec!["agent://fraud".into()]);
2034 rt.process(
2035 &env(
2036 "macp.mode.decision.v1",
2037 "SessionStart",
2038 "m1",
2039 &sid,
2040 "agent://orchestrator",
2041 payload.clone(),
2042 ),
2043 None,
2044 )
2045 .await
2046 .unwrap();
2047
2048 let result = rt
2049 .process(
2050 &env(
2051 "macp.mode.decision.v1",
2052 "SessionStart",
2053 "m1",
2054 &sid,
2055 "agent://orchestrator",
2056 payload,
2057 ),
2058 None,
2059 )
2060 .await
2061 .unwrap();
2062 assert!(result.duplicate);
2063 }
2064
2065 #[tokio::test]
2066 async fn non_start_mode_mismatch_rejected() {
2067 let rt = make_runtime();
2068 let sid = new_sid();
2069 rt.process(
2070 &env(
2071 "macp.mode.decision.v1",
2072 "SessionStart",
2073 "m1",
2074 &sid,
2075 "agent://orchestrator",
2076 session_start(vec!["agent://fraud".into()]),
2077 ),
2078 None,
2079 )
2080 .await
2081 .unwrap();
2082
2083 let proposal = ProposalPayload {
2084 proposal_id: "p1".into(),
2085 option: "step-up".into(),
2086 rationale: "risk".into(),
2087 supporting_data: vec![],
2088 }
2089 .encode_to_vec();
2090 let err = rt
2091 .process(
2092 &env(
2093 "macp.mode.task.v1",
2094 "Proposal",
2095 "m2",
2096 &sid,
2097 "agent://orchestrator",
2098 proposal,
2099 ),
2100 None,
2101 )
2102 .await
2103 .unwrap_err();
2104 assert_eq!(err.to_string(), "InvalidEnvelope");
2105 }
2106
2107 #[tokio::test]
2108 async fn cancel_idempotent_on_already_expired() {
2109 let rt = make_runtime();
2110 let sid = new_sid();
2111 let payload = SessionStartPayload {
2112 intent: "intent".into(),
2113 participants: vec!["agent://fraud".into()],
2114 mode_version: "1.0.0".into(),
2115 configuration_version: "cfg-1".into(),
2116 policy_version: String::new(),
2117 ttl_ms: 1,
2118 context_id: String::new(),
2119 extensions: std::collections::HashMap::new(),
2120 roots: vec![],
2121 max_suspend_ms: 0,
2122 }
2123 .encode_to_vec();
2124 rt.process(
2125 &env(
2126 "macp.mode.decision.v1",
2127 "SessionStart",
2128 "m1",
2129 &sid,
2130 "agent://orchestrator",
2131 payload,
2132 ),
2133 None,
2134 )
2135 .await
2136 .unwrap();
2137 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2138 let result = rt
2139 .cancel_session(&sid, "cleanup", "agent://orchestrator")
2140 .await
2141 .unwrap();
2142 assert_eq!(result.session_state, SessionState::Expired);
2143 }
2144
2145 #[tokio::test]
2146 async fn accepted_envelopes_are_published_in_order() {
2147 let rt = make_runtime();
2148 let sid = new_sid();
2149 let mut events = rt.subscribe_session_stream(&sid);
2150
2151 let start = env(
2152 "macp.mode.decision.v1",
2153 "SessionStart",
2154 "m1",
2155 &sid,
2156 "agent://orchestrator",
2157 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2158 );
2159 rt.process(&start, None).await.unwrap();
2160 let first = events.recv().await.unwrap();
2161 assert_eq!(first.message_id, "m1");
2162 assert_eq!(first.message_type, "SessionStart");
2163
2164 let proposal = ProposalPayload {
2165 proposal_id: "p1".into(),
2166 option: "step-up".into(),
2167 rationale: "risk".into(),
2168 supporting_data: vec![],
2169 }
2170 .encode_to_vec();
2171 let proposal_env = env(
2172 "macp.mode.decision.v1",
2173 "Proposal",
2174 "m2",
2175 &sid,
2176 "agent://orchestrator",
2177 proposal,
2178 );
2179 rt.process(&proposal_env, None).await.unwrap();
2180 let second = events.recv().await.unwrap();
2181 assert_eq!(second.message_id, "m2");
2182 assert_eq!(second.message_type, "Proposal");
2183 }
2184
2185 #[tokio::test]
2186 async fn commitment_versions_are_carried_into_resolution() {
2187 let rt = make_runtime();
2188 let sid = new_sid();
2189 rt.process(
2190 &env(
2191 "macp.mode.proposal.v1",
2192 "SessionStart",
2193 "m1",
2194 &sid,
2195 "agent://buyer",
2196 session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
2197 ),
2198 None,
2199 )
2200 .await
2201 .unwrap();
2202
2203 let proposal = crate::proposal_pb::ProposalPayload {
2204 proposal_id: "p1".into(),
2205 title: "offer".into(),
2206 summary: "summary".into(),
2207 details: vec![],
2208 tags: vec![],
2209 }
2210 .encode_to_vec();
2211 rt.process(
2212 &env(
2213 "macp.mode.proposal.v1",
2214 "Proposal",
2215 "m2",
2216 &sid,
2217 "agent://seller",
2218 proposal,
2219 ),
2220 None,
2221 )
2222 .await
2223 .unwrap();
2224 let accept = crate::proposal_pb::AcceptPayload {
2225 proposal_id: "p1".into(),
2226 reason: String::new(),
2227 }
2228 .encode_to_vec();
2229 rt.process(
2230 &env(
2231 "macp.mode.proposal.v1",
2232 "Accept",
2233 "m3",
2234 &sid,
2235 "agent://seller",
2236 accept.clone(),
2237 ),
2238 None,
2239 )
2240 .await
2241 .unwrap();
2242 rt.process(
2243 &env(
2244 "macp.mode.proposal.v1",
2245 "Accept",
2246 "m4",
2247 &sid,
2248 "agent://buyer",
2249 accept,
2250 ),
2251 None,
2252 )
2253 .await
2254 .unwrap();
2255 let commitment = CommitmentPayload {
2256 commitment_id: "c1".into(),
2257 action: "proposal.accepted".into(),
2258 authority_scope: "commercial".into(),
2259 reason: "bound".into(),
2260 mode_version: "1.0.0".into(),
2261 policy_version: "policy.default".into(),
2262 configuration_version: "cfg-1".into(),
2263 outcome_positive: true,
2264 supersedes: None,
2265 }
2266 .encode_to_vec();
2267 let result = rt
2268 .process(
2269 &env(
2270 "macp.mode.proposal.v1",
2271 "Commitment",
2272 "m5",
2273 &sid,
2274 "agent://buyer",
2275 commitment,
2276 ),
2277 None,
2278 )
2279 .await
2280 .unwrap();
2281 assert_eq!(result.session_state, SessionState::Resolved);
2282 }
2283
2284 #[tokio::test]
2285 async fn max_open_sessions_enforced_under_write_lock() {
2286 let rt = make_runtime();
2287 let sid1 = new_sid();
2288 let sid2 = new_sid();
2289 let sid3 = new_sid();
2290 rt.process(
2291 &env(
2292 "macp.mode.decision.v1",
2293 "SessionStart",
2294 "m1",
2295 &sid1,
2296 "agent://orchestrator",
2297 session_start(vec!["agent://fraud".into()]),
2298 ),
2299 Some(1),
2300 )
2301 .await
2302 .unwrap();
2303
2304 let err = rt
2305 .process(
2306 &env(
2307 "macp.mode.decision.v1",
2308 "SessionStart",
2309 "m2",
2310 &sid2,
2311 "agent://orchestrator",
2312 session_start(vec!["agent://fraud".into()]),
2313 ),
2314 Some(1),
2315 )
2316 .await
2317 .unwrap_err();
2318 assert!(matches!(err, MacpError::RateLimited));
2319
2320 rt.process(
2321 &env(
2322 "macp.mode.decision.v1",
2323 "SessionStart",
2324 "m3",
2325 &sid3,
2326 "agent://other",
2327 session_start(vec!["agent://fraud".into()]),
2328 ),
2329 Some(1),
2330 )
2331 .await
2332 .unwrap();
2333 }
2334
2335 #[tokio::test]
2336 async fn weak_session_id_rejected() {
2337 let rt = make_runtime();
2338 let err = rt
2339 .process(
2340 &env(
2341 "macp.mode.decision.v1",
2342 "SessionStart",
2343 "m1",
2344 "s1",
2345 "agent://orchestrator",
2346 session_start(vec!["agent://fraud".into()]),
2347 ),
2348 None,
2349 )
2350 .await
2351 .unwrap_err();
2352 assert_eq!(err.to_string(), "InvalidSessionId");
2353 }
2354
2355 #[tokio::test]
2356 async fn log_append_failure_rejects_session_start() {
2357 use std::io;
2358 struct FailingBackend;
2359 #[async_trait::async_trait]
2360 impl StorageBackend for FailingBackend {
2361 async fn save_session(&self, _: &Session) -> io::Result<()> {
2362 Ok(())
2363 }
2364 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2365 Ok(None)
2366 }
2367 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2368 Ok(vec![])
2369 }
2370 async fn delete_session(&self, _: &str) -> io::Result<()> {
2371 Ok(())
2372 }
2373 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2374 Ok(vec![])
2375 }
2376 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2377 Err(io::Error::other("disk full"))
2378 }
2379 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2380 Ok(vec![])
2381 }
2382 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2383 Ok(())
2384 }
2385 }
2386
2387 let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
2388 let registry = Arc::new(SessionRegistry::new());
2389 let log_store = Arc::new(LogStore::new());
2390 let rt = Runtime::new(storage, registry, log_store);
2391 let sid = new_sid();
2392
2393 let err = rt
2394 .process(
2395 &env(
2396 "macp.mode.decision.v1",
2397 "SessionStart",
2398 "m1",
2399 &sid,
2400 "agent://orchestrator",
2401 session_start(vec!["agent://fraud".into()]),
2402 ),
2403 None,
2404 )
2405 .await
2406 .unwrap_err();
2407 assert_eq!(err.to_string(), "StorageFailed");
2408 }
2409
2410 #[tokio::test]
2411 async fn log_append_failure_rejects_in_session_message() {
2412 use std::io;
2413 use std::sync::atomic::{AtomicUsize, Ordering};
2414
2415 struct FailOnSecondAppend {
2416 count: AtomicUsize,
2417 }
2418 #[async_trait::async_trait]
2419 impl StorageBackend for FailOnSecondAppend {
2420 async fn save_session(&self, _: &Session) -> io::Result<()> {
2421 Ok(())
2422 }
2423 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2424 Ok(None)
2425 }
2426 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2427 Ok(vec![])
2428 }
2429 async fn delete_session(&self, _: &str) -> io::Result<()> {
2430 Ok(())
2431 }
2432 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2433 Ok(vec![])
2434 }
2435 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2436 let n = self.count.fetch_add(1, Ordering::SeqCst);
2437 if n >= 1 {
2438 Err(io::Error::other("disk full"))
2439 } else {
2440 Ok(())
2441 }
2442 }
2443 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2444 Ok(vec![])
2445 }
2446 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2447 Ok(())
2448 }
2449 }
2450
2451 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2452 count: AtomicUsize::new(0),
2453 });
2454 let registry = Arc::new(SessionRegistry::new());
2455 let log_store = Arc::new(LogStore::new());
2456 let rt = Runtime::new(storage, registry, log_store);
2457 let sid = new_sid();
2458
2459 rt.process(
2461 &env(
2462 "macp.mode.decision.v1",
2463 "SessionStart",
2464 "m1",
2465 &sid,
2466 "agent://orchestrator",
2467 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2468 ),
2469 None,
2470 )
2471 .await
2472 .unwrap();
2473
2474 let proposal = ProposalPayload {
2476 proposal_id: "p1".into(),
2477 option: "step-up".into(),
2478 rationale: "risk".into(),
2479 supporting_data: vec![],
2480 }
2481 .encode_to_vec();
2482 let err = rt
2483 .process(
2484 &env(
2485 "macp.mode.decision.v1",
2486 "Proposal",
2487 "m2",
2488 &sid,
2489 "agent://orchestrator",
2490 proposal,
2491 ),
2492 None,
2493 )
2494 .await
2495 .unwrap_err();
2496 assert_eq!(err.to_string(), "StorageFailed");
2497
2498 let session = rt.get_session_checked(&sid).await.unwrap();
2500 assert!(!session.seen_message_ids.contains("m2"));
2501 }
2502
2503 #[tokio::test]
2504 async fn cancel_session_fails_if_log_append_fails() {
2505 use std::io;
2506 use std::sync::atomic::{AtomicUsize, Ordering};
2507
2508 struct FailOnSecondAppend {
2509 count: AtomicUsize,
2510 }
2511 #[async_trait::async_trait]
2512 impl StorageBackend for FailOnSecondAppend {
2513 async fn save_session(&self, _: &Session) -> io::Result<()> {
2514 Ok(())
2515 }
2516 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2517 Ok(None)
2518 }
2519 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2520 Ok(vec![])
2521 }
2522 async fn delete_session(&self, _: &str) -> io::Result<()> {
2523 Ok(())
2524 }
2525 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2526 Ok(vec![])
2527 }
2528 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2529 let n = self.count.fetch_add(1, Ordering::SeqCst);
2530 if n >= 1 {
2531 Err(io::Error::other("disk full"))
2532 } else {
2533 Ok(())
2534 }
2535 }
2536 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2537 Ok(vec![])
2538 }
2539 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2540 Ok(())
2541 }
2542 }
2543
2544 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2545 count: AtomicUsize::new(0),
2546 });
2547 let registry = Arc::new(SessionRegistry::new());
2548 let log_store = Arc::new(LogStore::new());
2549 let rt = Runtime::new(storage, registry, log_store);
2550 let sid = new_sid();
2551
2552 rt.process(
2553 &env(
2554 "macp.mode.decision.v1",
2555 "SessionStart",
2556 "m1",
2557 &sid,
2558 "agent://orchestrator",
2559 session_start(vec!["agent://fraud".into()]),
2560 ),
2561 None,
2562 )
2563 .await
2564 .unwrap();
2565
2566 let err = rt
2567 .cancel_session(&sid, "test cancel", "agent://orchestrator")
2568 .await
2569 .unwrap_err();
2570 assert_eq!(err.to_string(), "StorageFailed");
2571 }
2572
2573 #[tokio::test]
2574 async fn ttl_expiration_rejects_message() {
2575 let rt = make_runtime();
2576 let sid = new_sid();
2577 let payload = SessionStartPayload {
2578 intent: "intent".into(),
2579 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
2580 mode_version: "1.0.0".into(),
2581 configuration_version: "cfg-1".into(),
2582 policy_version: String::new(),
2583 ttl_ms: 1,
2584 context_id: String::new(),
2585 extensions: std::collections::HashMap::new(),
2586 roots: vec![],
2587 max_suspend_ms: 0,
2588 }
2589 .encode_to_vec();
2590 rt.process(
2591 &env(
2592 "macp.mode.decision.v1",
2593 "SessionStart",
2594 "m1",
2595 &sid,
2596 "agent://orchestrator",
2597 payload,
2598 ),
2599 None,
2600 )
2601 .await
2602 .unwrap();
2603 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2604 let proposal = ProposalPayload {
2605 proposal_id: "p1".into(),
2606 option: "step-up".into(),
2607 rationale: "risk".into(),
2608 supporting_data: vec![],
2609 }
2610 .encode_to_vec();
2611 let err = rt
2612 .process(
2613 &env(
2614 "macp.mode.decision.v1",
2615 "Proposal",
2616 "m2",
2617 &sid,
2618 "agent://orchestrator",
2619 proposal,
2620 ),
2621 None,
2622 )
2623 .await
2624 .unwrap_err();
2625 assert_eq!(err.to_string(), "TtlExpired");
2626 }
2627
2628 #[tokio::test]
2629 async fn cleanup_expired_sessions_marks_expired() {
2630 let rt = make_runtime();
2631 let sid = new_sid();
2632 let payload = SessionStartPayload {
2633 intent: "intent".into(),
2634 participants: vec!["agent://fraud".into()],
2635 mode_version: "1.0.0".into(),
2636 configuration_version: "cfg-1".into(),
2637 policy_version: String::new(),
2638 ttl_ms: 1,
2639 context_id: String::new(),
2640 extensions: std::collections::HashMap::new(),
2641 roots: vec![],
2642 max_suspend_ms: 0,
2643 }
2644 .encode_to_vec();
2645 rt.process(
2646 &env(
2647 "macp.mode.decision.v1",
2648 "SessionStart",
2649 "m1",
2650 &sid,
2651 "agent://orchestrator",
2652 payload,
2653 ),
2654 None,
2655 )
2656 .await
2657 .unwrap();
2658 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2659 rt.cleanup_expired_sessions().await;
2660 let session = rt.get_session_checked(&sid).await.unwrap();
2661 assert_eq!(session.state, SessionState::Expired);
2662 }
2663
2664 #[tokio::test]
2665 async fn evict_stale_sessions_removes_resolved() {
2666 let rt = make_runtime();
2667 let sid = new_sid();
2668 rt.process(
2670 &env(
2671 "macp.mode.decision.v1",
2672 "SessionStart",
2673 "m1",
2674 &sid,
2675 "agent://orchestrator",
2676 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2677 ),
2678 None,
2679 )
2680 .await
2681 .unwrap();
2682 let proposal = ProposalPayload {
2684 proposal_id: "p1".into(),
2685 option: "step-up".into(),
2686 rationale: "risk".into(),
2687 supporting_data: vec![],
2688 }
2689 .encode_to_vec();
2690 rt.process(
2691 &env(
2692 "macp.mode.decision.v1",
2693 "Proposal",
2694 "m2",
2695 &sid,
2696 "agent://orchestrator",
2697 proposal,
2698 ),
2699 None,
2700 )
2701 .await
2702 .unwrap();
2703 let commitment = CommitmentPayload {
2705 commitment_id: "c1".into(),
2706 action: "decision.selected".into(),
2707 authority_scope: "payments".into(),
2708 reason: "bound".into(),
2709 mode_version: "1.0.0".into(),
2710 policy_version: "policy.default".into(),
2711 configuration_version: "cfg-1".into(),
2712 outcome_positive: true,
2713 supersedes: None,
2714 }
2715 .encode_to_vec();
2716 let result = rt
2717 .process(
2718 &env(
2719 "macp.mode.decision.v1",
2720 "Commitment",
2721 "m3",
2722 &sid,
2723 "agent://orchestrator",
2724 commitment,
2725 ),
2726 None,
2727 )
2728 .await
2729 .unwrap();
2730 assert_eq!(result.session_state, SessionState::Resolved);
2731 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2733 rt.evict_stale_sessions(0).await;
2735 assert!(rt.registry.get_session(&sid).await.is_none());
2737 }
2738
2739 #[tokio::test]
2740 async fn session_start_with_wrong_mode_version_rejected() {
2741 let rt = make_runtime();
2742 let sid = new_sid();
2743 let payload = SessionStartPayload {
2744 intent: "test".into(),
2745 participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
2746 mode_version: "99.0.0".into(), configuration_version: "cfg-1".into(),
2748 policy_version: String::new(),
2749 ttl_ms: 60_000,
2750 context_id: String::new(),
2751 extensions: std::collections::HashMap::new(),
2752 roots: vec![],
2753 max_suspend_ms: 0,
2754 }
2755 .encode_to_vec();
2756
2757 let err = rt
2758 .process(
2759 &env(
2760 "macp.mode.decision.v1",
2761 "SessionStart",
2762 "m1",
2763 &sid,
2764 "agent://orchestrator",
2765 payload,
2766 ),
2767 None,
2768 )
2769 .await
2770 .unwrap_err();
2771 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2772 }
2773
2774 #[tokio::test]
2775 async fn signal_empty_signal_type_rejected() {
2776 let rt = make_runtime();
2777 let signal_payload = crate::pb::SignalPayload {
2779 signal_type: String::new(),
2780 data: b"some data".to_vec(),
2781 confidence: 0.0,
2782 correlation_session_id: String::new(),
2783 }
2784 .encode_to_vec();
2785 let signal = Envelope {
2786 macp_version: "1.0".into(),
2787 mode: String::new(),
2788 message_type: "Signal".into(),
2789 message_id: "sig-1".into(),
2790 session_id: String::new(),
2791 sender: "agent://a".into(),
2792 timestamp_unix_ms: 0,
2793 payload: signal_payload,
2794 };
2795 let err = rt.process_signal(&signal).await.unwrap_err();
2796 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2797 }
2798
2799 #[tokio::test]
2800 async fn signal_valid_payload_accepted() {
2801 let rt = make_runtime();
2802 let signal_payload = crate::pb::SignalPayload {
2803 signal_type: "heartbeat".into(),
2804 data: vec![],
2805 confidence: 0.8,
2806 correlation_session_id: String::new(),
2807 }
2808 .encode_to_vec();
2809 let signal = Envelope {
2810 macp_version: "1.0".into(),
2811 mode: String::new(),
2812 message_type: "Signal".into(),
2813 message_id: "sig-2".into(),
2814 session_id: String::new(),
2815 sender: "agent://a".into(),
2816 timestamp_unix_ms: 0,
2817 payload: signal_payload,
2818 };
2819 rt.process_signal(&signal).await.unwrap();
2820 }
2821
2822 #[tokio::test]
2823 async fn signal_empty_payload_accepted() {
2824 let rt = make_runtime();
2825 let signal = Envelope {
2826 macp_version: "1.0".into(),
2827 mode: String::new(),
2828 message_type: "Signal".into(),
2829 message_id: "sig-3".into(),
2830 session_id: String::new(),
2831 sender: "agent://a".into(),
2832 timestamp_unix_ms: 0,
2833 payload: vec![],
2834 };
2835 rt.process_signal(&signal).await.unwrap();
2836 }
2837
2838 #[tokio::test]
2844 async fn ext_mode_empty_version_binds_descriptor_version() {
2845 let rt = make_runtime();
2846 rt.register_extension(ModeDescriptor {
2847 mode: "ext.dyn.v1".into(),
2848 mode_version: "2.5.0".into(),
2849 message_types: vec!["SessionStart".into(), "Note".into(), "Commitment".into()],
2850 terminal_message_types: vec!["Commitment".into()],
2851 ..Default::default()
2852 })
2853 .unwrap();
2854
2855 let sid = new_sid();
2856 let payload = SessionStartPayload {
2857 participants: vec!["alice".into()],
2858 configuration_version: "cfg-1".into(),
2859 ttl_ms: 60_000,
2860 ..Default::default()
2861 }
2862 .encode_to_vec();
2863 rt.process(
2864 &env("ext.dyn.v1", "SessionStart", "m1", &sid, "alice", payload),
2865 None,
2866 )
2867 .await
2868 .unwrap();
2869
2870 let session = rt.get_session_checked(&sid).await.unwrap();
2872 assert_eq!(session.mode_version, "2.5.0");
2873
2874 let bad = CommitmentPayload {
2876 commitment_id: "c1".into(),
2877 action: "work.completed".into(),
2878 authority_scope: "test".into(),
2879 reason: "done".into(),
2880 mode_version: String::new(),
2881 policy_version: "policy.default".into(),
2882 configuration_version: "cfg-1".into(),
2883 outcome_positive: true,
2884 supersedes: None,
2885 }
2886 .encode_to_vec();
2887 let err = rt
2888 .process(
2889 &env("ext.dyn.v1", "Commitment", "m2", &sid, "alice", bad),
2890 None,
2891 )
2892 .await
2893 .unwrap_err();
2894 assert_eq!(err.to_string(), "InvalidPayload");
2895
2896 let good = CommitmentPayload {
2898 commitment_id: "c1".into(),
2899 action: "work.completed".into(),
2900 authority_scope: "test".into(),
2901 reason: "done".into(),
2902 mode_version: "2.5.0".into(),
2903 policy_version: "policy.default".into(),
2904 configuration_version: "cfg-1".into(),
2905 outcome_positive: true,
2906 supersedes: None,
2907 }
2908 .encode_to_vec();
2909 let result = rt
2910 .process(
2911 &env("ext.dyn.v1", "Commitment", "m3", &sid, "alice", good),
2912 None,
2913 )
2914 .await
2915 .unwrap();
2916 assert_eq!(result.session_state, SessionState::Resolved);
2917 }
2918
2919 #[tokio::test]
2922 async fn ext_mode_binding_recorded_on_session_start_log_entry() {
2923 let rt = make_runtime();
2924 rt.register_extension(ModeDescriptor {
2925 mode: "ext.dyn2.v1".into(),
2926 mode_version: "3.0.0".into(),
2927 message_types: vec!["SessionStart".into(), "Commitment".into()],
2928 terminal_message_types: vec!["Commitment".into()],
2929 ..Default::default()
2930 })
2931 .unwrap();
2932
2933 let sid = new_sid();
2934 let payload = SessionStartPayload {
2935 participants: vec!["alice".into()],
2936 configuration_version: "cfg-1".into(),
2937 ttl_ms: 60_000,
2938 ..Default::default()
2939 }
2940 .encode_to_vec();
2941 rt.process(
2942 &env("ext.dyn2.v1", "SessionStart", "m1", &sid, "alice", payload),
2943 None,
2944 )
2945 .await
2946 .unwrap();
2947
2948 let log = rt.log_store.get_log(&sid).await.unwrap();
2949 assert_eq!(log[0].message_type, "SessionStart");
2950 assert_eq!(log[0].bound_mode_version.as_deref(), Some("3.0.0"));
2951
2952 let sid2 = new_sid();
2954 let payload2 = SessionStartPayload {
2955 participants: vec!["alice".into()],
2956 mode_version: "3.0.0".into(),
2957 configuration_version: "cfg-1".into(),
2958 ttl_ms: 60_000,
2959 ..Default::default()
2960 }
2961 .encode_to_vec();
2962 rt.process(
2963 &env(
2964 "ext.dyn2.v1",
2965 "SessionStart",
2966 "m1",
2967 &sid2,
2968 "alice",
2969 payload2,
2970 ),
2971 None,
2972 )
2973 .await
2974 .unwrap();
2975 let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2976 assert_eq!(log2[0].bound_mode_version, None);
2977 }
2978
2979 #[tokio::test]
2984 async fn session_start_binds_and_records_max_suspend_cap() {
2985 let rt = make_runtime();
2986
2987 let sid = new_sid();
2989 let payload = SessionStartPayload {
2990 participants: vec!["alice".into(), "bob".into()],
2991 mode_version: "1.0.0".into(),
2992 configuration_version: "cfg-1".into(),
2993 ttl_ms: 60_000,
2994 max_suspend_ms: 12_345,
2995 ..Default::default()
2996 }
2997 .encode_to_vec();
2998 rt.process(
2999 &env(
3000 "macp.mode.decision.v1",
3001 "SessionStart",
3002 "m1",
3003 &sid,
3004 "alice",
3005 payload,
3006 ),
3007 None,
3008 )
3009 .await
3010 .unwrap();
3011 let log = rt.log_store.get_log(&sid).await.unwrap();
3012 assert_eq!(log[0].bound_max_suspend_ms, Some(12_345));
3013
3014 let sid2 = new_sid();
3016 let payload2 = SessionStartPayload {
3017 participants: vec!["alice".into(), "bob".into()],
3018 mode_version: "1.0.0".into(),
3019 configuration_version: "cfg-1".into(),
3020 ttl_ms: 60_000,
3021 max_suspend_ms: 0,
3022 ..Default::default()
3023 }
3024 .encode_to_vec();
3025 rt.process(
3026 &env(
3027 "macp.mode.decision.v1",
3028 "SessionStart",
3029 "m2",
3030 &sid2,
3031 "alice",
3032 payload2,
3033 ),
3034 None,
3035 )
3036 .await
3037 .unwrap();
3038 let log2 = rt.log_store.get_log(&sid2).await.unwrap();
3039 assert_eq!(
3040 log2[0].bound_max_suspend_ms,
3041 Some(macp_core::session::MAX_SUSPEND_MS)
3042 );
3043 }
3044
3045 #[test]
3046 fn audit_verbosity_reads_policy_rules() {
3047 let mut session = Session::builder("s1", "macp.mode.decision.v1", "a").build();
3048 assert!(!Runtime::audit_verbose(&session));
3049
3050 session.policy_definition = Some(macp_core::policy::PolicyDefinition {
3051 policy_id: "policy.test.audit".into(),
3052 mode: "*".into(),
3053 description: "audited".into(),
3054 rules: serde_json::json!({ "audit": { "level": "info" } }),
3055 schema_version: 1,
3056 });
3057 assert!(Runtime::audit_verbose(&session));
3058
3059 session.policy_definition.as_mut().unwrap().rules =
3060 serde_json::json!({ "audit": { "level": "debug" } });
3061 assert!(!Runtime::audit_verbose(&session));
3062 }
3063
3064 #[tokio::test]
3070 async fn session_start_snapshot_failure_is_nonfatal_after_commit_point() {
3071 use std::io;
3072
3073 struct FailSnapshotBackend;
3074 #[async_trait::async_trait]
3075 impl StorageBackend for FailSnapshotBackend {
3076 async fn create_session_storage(&self, _s: &str) -> io::Result<()> {
3077 Ok(())
3078 }
3079 async fn save_session(&self, _s: &Session) -> io::Result<()> {
3080 Err(io::Error::other("snapshot disk full"))
3081 }
3082 async fn load_session(&self, _s: &str) -> io::Result<Option<Session>> {
3083 Ok(None)
3084 }
3085 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
3086 Ok(vec![])
3087 }
3088 async fn delete_session(&self, _s: &str) -> io::Result<()> {
3089 Ok(())
3090 }
3091 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
3092 Ok(vec![])
3093 }
3094 async fn append_log_entry(
3095 &self,
3096 _s: &str,
3097 _e: &crate::log_store::LogEntry,
3098 ) -> io::Result<()> {
3099 Ok(())
3100 }
3101 async fn load_log(&self, _s: &str) -> io::Result<Vec<crate::log_store::LogEntry>> {
3102 Ok(vec![])
3103 }
3104 }
3105
3106 let rt = Runtime::new(
3107 Arc::new(FailSnapshotBackend),
3108 Arc::new(SessionRegistry::new()),
3109 Arc::new(LogStore::new()),
3110 );
3111 let sid = new_sid();
3112 let result = rt
3113 .process(
3114 &env(
3115 "macp.mode.decision.v1",
3116 "SessionStart",
3117 "m1",
3118 &sid,
3119 "agent://orchestrator",
3120 session_start(vec!["agent://orchestrator".into()]),
3121 ),
3122 None,
3123 )
3124 .await
3125 .expect("start must succeed: the log append (commit point) succeeded");
3126 assert!(!result.duplicate);
3127 assert!(rt.get_session_checked(&sid).await.is_some());
3129 }
3130
3131 const HANDOFF_MODE: &str = "macp.mode.handoff.v1";
3147 const OWNER: &str = "agent://owner";
3148 const TARGET: &str = "agent://target";
3149
3150 fn reserved_id(handoff_id: &str) -> String {
3151 format!(
3152 "{}{handoff_id}",
3153 crate::mode::handoff::IMPLICIT_ACCEPT_MESSAGE_ID_PREFIX
3154 )
3155 }
3156
3157 fn handoff_start_payload() -> Vec<u8> {
3158 session_start(vec![OWNER.into(), TARGET.into()])
3159 }
3160
3161 fn handoff_commitment_payload() -> Vec<u8> {
3164 CommitmentPayload {
3165 commitment_id: "c1".into(),
3166 action: "handoff.accepted".into(),
3167 authority_scope: "support".into(),
3168 reason: "bound".into(),
3169 mode_version: "1.0.0".into(),
3170 policy_version: "policy.default".into(),
3171 configuration_version: "cfg-1".into(),
3172 outcome_positive: true,
3173 supersedes: None,
3174 }
3175 .encode_to_vec()
3176 }
3177
3178 fn handoff_offer(handoff_id: &str) -> Vec<u8> {
3179 crate::handoff_pb::HandoffOfferPayload {
3180 handoff_id: handoff_id.into(),
3181 target_participant: TARGET.into(),
3182 scope: "support".into(),
3183 reason: "escalate".into(),
3184 }
3185 .encode_to_vec()
3186 }
3187
3188 fn handoff_context(handoff_id: &str) -> Vec<u8> {
3189 crate::handoff_pb::HandoffContextPayload {
3190 handoff_id: handoff_id.into(),
3191 content_type: "text/plain".into(),
3192 context: b"background".to_vec(),
3193 }
3194 .encode_to_vec()
3195 }
3196
3197 fn handoff_accept(handoff_id: &str, implicit: bool) -> Vec<u8> {
3198 crate::handoff_pb::HandoffAcceptPayload {
3199 handoff_id: handoff_id.into(),
3200 accepted_by: TARGET.into(),
3201 reason: "ready".into(),
3202 implicit,
3203 }
3204 .encode_to_vec()
3205 }
3206
3207 async fn handoff_session_with_offer(rt: &Runtime) -> String {
3210 let sid = new_sid();
3211 rt.process(
3212 &env(
3213 HANDOFF_MODE,
3214 "SessionStart",
3215 "start-1",
3216 &sid,
3217 OWNER,
3218 handoff_start_payload(),
3219 ),
3220 None,
3221 )
3222 .await
3223 .expect("handoff session start");
3224 rt.process(
3225 &env(
3226 HANDOFF_MODE,
3227 "HandoffOffer",
3228 "offer-1",
3229 &sid,
3230 OWNER,
3231 handoff_offer("h1"),
3232 ),
3233 None,
3234 )
3235 .await
3236 .expect("handoff offer");
3237 assert_eq!(
3238 rt.get_session_checked(&sid).await.unwrap().semantics_rev,
3239 macp_core::session::CURRENT_SEMANTICS_REV
3240 );
3241 sid
3242 }
3243
3244 #[tokio::test]
3257 async fn reserved_message_id_namespace_is_rejected_at_rev2() {
3258 let rt = make_runtime();
3259 let sid = handoff_session_with_offer(&rt).await;
3260
3261 let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3262 let dedup_before = rt
3263 .get_session_checked(&sid)
3264 .await
3265 .unwrap()
3266 .seen_message_ids
3267 .clone();
3268
3269 for (message_type, sender, payload) in [
3270 ("HandoffContext", OWNER, handoff_context("h1")),
3271 ("Commitment", OWNER, handoff_commitment_payload()),
3272 ("HandoffAccept", TARGET, handoff_accept("h1", false)),
3273 ] {
3274 let err = rt
3275 .process(
3276 &env(
3277 HANDOFF_MODE,
3278 message_type,
3279 &reserved_id("h1"),
3280 &sid,
3281 sender,
3282 payload,
3283 ),
3284 None,
3285 )
3286 .await
3287 .unwrap_err();
3288 assert!(
3289 matches!(err, MacpError::InvalidEnvelope),
3290 "{message_type} with a reserved id must be InvalidEnvelope, got {err}"
3291 );
3292 }
3293
3294 let session = rt.get_session_checked(&sid).await.unwrap();
3296 assert_eq!(
3297 rt.log_store.get_log(&sid).await.unwrap().len(),
3298 history_before
3299 );
3300 assert_eq!(session.seen_message_ids, dedup_before);
3301 assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3302 assert_eq!(session.state, SessionState::Open);
3303
3304 rt.process(
3307 &env(
3308 HANDOFF_MODE,
3309 "HandoffContext",
3310 "ctx-1",
3311 &sid,
3312 OWNER,
3313 handoff_context("h1"),
3314 ),
3315 None,
3316 )
3317 .await
3318 .expect("an ordinary id is accepted");
3319 assert_eq!(
3320 rt.log_store.get_log(&sid).await.unwrap().len(),
3321 history_before + 1
3322 );
3323 }
3324
3325 #[tokio::test]
3336 async fn reserved_message_id_is_rejected_on_the_session_start_path() {
3337 let rt = make_runtime();
3338 let sid = new_sid();
3339
3340 let err = rt
3341 .process(
3342 &env(
3343 HANDOFF_MODE,
3344 "SessionStart",
3345 &reserved_id("h1"),
3346 &sid,
3347 OWNER,
3348 handoff_start_payload(),
3349 ),
3350 None,
3351 )
3352 .await
3353 .unwrap_err();
3354 assert!(
3355 matches!(err, MacpError::InvalidEnvelope),
3356 "reserved id on SessionStart must be InvalidEnvelope, got {err}"
3357 );
3358
3359 assert!(rt.get_session_checked(&sid).await.is_none());
3361 assert!(rt.log_store.get_log(&sid).await.is_none());
3362 assert!(!rt.registry.sessions.read().await.contains_key(&sid));
3363
3364 rt.process(
3367 &env(
3368 HANDOFF_MODE,
3369 "SessionStart",
3370 "start-1",
3371 &sid,
3372 OWNER,
3373 handoff_start_payload(),
3374 ),
3375 None,
3376 )
3377 .await
3378 .expect("a rejected SessionStart must not reserve the session id");
3379 let session = rt.get_session_checked(&sid).await.unwrap();
3380 assert!(session.seen_message_ids.contains("start-1"));
3381 assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3382 }
3383
3384 #[tokio::test]
3414 async fn client_implicit_accept_rejected_through_the_runtime() {
3415 let rt = make_runtime();
3416 let sid = handoff_session_with_offer(&rt).await;
3417 let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3418
3419 let err = rt
3421 .process(
3422 &env(
3423 HANDOFF_MODE,
3424 "HandoffAccept",
3425 "accept-1",
3426 &sid,
3427 TARGET,
3428 handoff_accept("h1", true),
3429 ),
3430 None,
3431 )
3432 .await
3433 .unwrap_err();
3434 assert!(matches!(err, MacpError::InvalidPayload), "got {err}");
3435
3436 let err = rt
3438 .process(
3439 &env(
3440 HANDOFF_MODE,
3441 "HandoffAccept",
3442 &reserved_id("h1"),
3443 &sid,
3444 TARGET,
3445 handoff_accept("h1", true),
3446 ),
3447 None,
3448 )
3449 .await
3450 .unwrap_err();
3451 assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3452
3453 let session = rt.get_session_checked(&sid).await.unwrap();
3455 assert_eq!(
3456 rt.log_store.get_log(&sid).await.unwrap().len(),
3457 history_before
3458 );
3459 assert!(session.seen_message_ids.is_disjoint(
3460 &["accept-1".to_string(), reserved_id("h1")]
3461 .into_iter()
3462 .collect()
3463 ));
3464 let mode_state: serde_json::Value = serde_json::from_slice(&session.mode_state).unwrap();
3465 assert_eq!(mode_state["offers"]["h1"]["disposition"], "Offered");
3466
3467 rt.process(
3471 &env(
3472 HANDOFF_MODE,
3473 "HandoffAccept",
3474 "accept-2",
3475 &sid,
3476 TARGET,
3477 handoff_accept("h1", false),
3478 ),
3479 None,
3480 )
3481 .await
3482 .expect("an explicit accept is still accepted");
3483 }
3484
3485 #[tokio::test]
3491 async fn client_boundary_error_ordering_is_unchanged_at_rev2() {
3492 let rt = make_runtime();
3493 let sid = handoff_session_with_offer(&rt).await;
3494
3495 let err = rt
3497 .process(
3498 &env(
3499 HANDOFF_MODE,
3500 "HandoffAccept",
3501 &reserved_id("h1"),
3502 &sid,
3503 "agent://stranger",
3504 handoff_accept("h1", true),
3505 ),
3506 None,
3507 )
3508 .await
3509 .unwrap_err();
3510 assert!(
3511 matches!(err, MacpError::Forbidden),
3512 "authorization must be reported before the client boundary, got {err}"
3513 );
3514
3515 let err = rt
3517 .process(
3518 &env(
3519 HANDOFF_MODE,
3520 "HandoffAccept",
3521 &reserved_id("h1"),
3522 &sid,
3523 TARGET,
3524 handoff_accept("h1", true),
3525 ),
3526 None,
3527 )
3528 .await
3529 .unwrap_err();
3530 assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3531 }
3532
3533 async fn handoff_session_with_timed_offer(rt: &Runtime, timeout_ms: i64) -> String {
3539 rt.register_policy(macp_core::policy::PolicyDefinition {
3540 policy_id: "handoff-timed".into(),
3541 mode: HANDOFF_MODE.into(),
3542 description: "implicit accept".into(),
3543 rules: serde_json::json!({
3544 "acceptance": { "implicit_accept_timeout_ms": timeout_ms },
3545 "commitment": { "authority": "initiator_only" }
3546 }),
3547 schema_version: 1,
3548 })
3549 .expect("policy registers");
3550
3551 let sid = new_sid();
3552 let start = SessionStartPayload {
3553 intent: "escalate".into(),
3554 participants: vec![OWNER.into(), TARGET.into()],
3555 mode_version: "1.0.0".into(),
3556 configuration_version: "cfg-1".into(),
3557 policy_version: "handoff-timed".into(),
3558 ttl_ms: 60_000,
3559 context_id: String::new(),
3560 extensions: std::collections::HashMap::new(),
3561 roots: vec![],
3562 max_suspend_ms: 0,
3563 }
3564 .encode_to_vec();
3565 rt.process(
3566 &env(HANDOFF_MODE, "SessionStart", "start-1", &sid, OWNER, start),
3567 None,
3568 )
3569 .await
3570 .expect("session start");
3571 rt.process(
3572 &env(
3573 HANDOFF_MODE,
3574 "HandoffOffer",
3575 "offer-1",
3576 &sid,
3577 OWNER,
3578 handoff_offer("h1"),
3579 ),
3580 None,
3581 )
3582 .await
3583 .expect("offer");
3584 sid
3585 }
3586
3587 #[tokio::test]
3607 async fn synthesis_is_skipped_for_a_non_open_session() {
3608 let rt = make_runtime();
3609 let sid = handoff_session_with_timed_offer(&rt, 20).await;
3610 rt.suspend_session(&sid, "hold", OWNER)
3611 .await
3612 .expect("suspend");
3613
3614 let shared = rt.registry.get_shared(&sid).await.unwrap();
3615 let mut guard = shared.lock().await;
3616 let session = &mut *guard;
3617 assert_eq!(session.state, SessionState::Suspended);
3618 assert!(session.suspended_at_ms.is_some());
3619
3620 let long_after = session.suspended_at_ms.unwrap() + 10_000;
3623 let log_before = rt.log_store.get_log(&sid).await.unwrap_or_default().len();
3624 let dedup_before = session.seen_message_ids.len();
3625 let mode_state_before = session.mode_state.clone();
3626
3627 rt.synthesize_due_accept(&sid, session, long_after)
3628 .await
3629 .expect("the filter is a skip, not an error");
3630
3631 assert_eq!(
3632 rt.log_store.get_log(&sid).await.unwrap_or_default().len(),
3633 log_before,
3634 "a suspended session must not gain a synthetic entry"
3635 );
3636 assert_eq!(session.seen_message_ids.len(), dedup_before);
3637 assert_eq!(session.mode_state, mode_state_before);
3638
3639 let mut without_the_tripwire = session.clone();
3651 without_the_tripwire.state = SessionState::Open;
3652 without_the_tripwire.suspended_at_ms = None;
3653 let mode = rt.mode_registry.get_mode(&session.mode).unwrap();
3654 let would_have_emitted = mode
3655 .due_synthetic_envelope(&without_the_tripwire, long_after)
3656 .expect("the mode would have synthesized; only the kernel filter stopped it");
3657 let suspended_at = session.suspended_at_ms.unwrap();
3663 assert!(
3664 would_have_emitted.timestamp_unix_ms >= suspended_at
3665 && would_have_emitted.timestamp_unix_ms < long_after,
3666 "D {} must fall inside the still-open pause starting at {suspended_at}",
3667 would_have_emitted.timestamp_unix_ms
3668 );
3669
3670 drop(guard);
3673 rt.resume_session(&sid, "go", OWNER).await.expect("resume");
3674 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3675 let shared = rt.registry.get_shared(&sid).await.unwrap();
3676 let mut guard = shared.lock().await;
3677 let session = &mut *guard;
3678 let now = Utc::now().timestamp_millis();
3679 rt.synthesize_due_accept(&sid, session, now).await.unwrap();
3680 assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3681 }
3682
3683 #[tokio::test]
3693 async fn synthetic_entry_stamps_received_at_with_the_deadline() {
3694 let rt = make_runtime();
3695 let sid = handoff_session_with_timed_offer(&rt, 20).await;
3696 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3697
3698 let offer_received_at = rt
3699 .log_store
3700 .get_log(&sid)
3701 .await
3702 .unwrap()
3703 .iter()
3704 .find(|e| e.message_type == "HandoffOffer")
3705 .expect("offer entry")
3706 .received_at_ms;
3707 let expected_d = offer_received_at + 20;
3708
3709 let shared = rt.registry.get_shared(&sid).await.unwrap();
3710 let mut guard = shared.lock().await;
3711 let session = &mut *guard;
3712 let observed = Utc::now().timestamp_millis();
3714 assert!(observed > expected_d);
3715 rt.synthesize_due_accept(&sid, session, observed)
3716 .await
3717 .unwrap();
3718 drop(guard);
3719
3720 let entry = rt
3721 .log_store
3722 .get_log(&sid)
3723 .await
3724 .unwrap()
3725 .into_iter()
3726 .find(|e| e.message_id == reserved_id("h1"))
3727 .expect("the synthetic entry");
3728 assert_eq!(entry.timestamp_unix_ms, expected_d, "envelope clock is D");
3729 assert_eq!(entry.received_at_ms, expected_d, "entry clock is D");
3730 assert_ne!(
3731 entry.received_at_ms, observed,
3732 "received_at_ms must not be the observation time"
3733 );
3734 assert_eq!(entry.entry_kind, EntryKind::Incoming);
3735 }
3736
3737 #[tokio::test]
3747 async fn synthetic_accept_is_not_credited_as_participant_activity() {
3748 let rt = make_runtime();
3749 let sid = handoff_session_with_timed_offer(&rt, 20).await;
3750 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3751
3752 let shared = rt.registry.get_shared(&sid).await.unwrap();
3753 let mut guard = shared.lock().await;
3754 let session = &mut *guard;
3755 let before = session.participant_message_counts.get(TARGET).copied();
3756 rt.synthesize_due_accept(&sid, session, Utc::now().timestamp_millis())
3757 .await
3758 .unwrap();
3759 assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3760 assert_eq!(
3761 session.participant_message_counts.get(TARGET).copied(),
3762 before,
3763 "the target must not be credited with a message they did not send"
3764 );
3765 }
3766 #[tokio::test]
3793 async fn a_synthetic_entry_on_the_checkpoint_boundary_checkpoints_either_path() {
3794 async fn shape(rt: &Runtime, sid: &str) -> Vec<(EntryKind, String)> {
3795 rt.log_store
3796 .get_log(sid)
3797 .await
3798 .expect("log")
3799 .iter()
3800 .map(|e| (e.entry_kind.clone(), e.message_type.clone()))
3801 .collect()
3802 }
3803
3804 let mut eager = make_runtime();
3806 eager.checkpoint_interval = 3;
3807 let eager_sid = handoff_session_with_timed_offer(&eager, 20).await;
3808 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3809 assert_eq!(
3810 eager.sweep_due_synthetic_accepts().await,
3811 1,
3812 "sweep emitted"
3813 );
3814
3815 let mut lazy = make_runtime();
3819 lazy.checkpoint_interval = 3;
3820 let lazy_sid = handoff_session_with_timed_offer(&lazy, 20).await;
3821 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3822 lazy.process(
3823 &env(
3824 HANDOFF_MODE,
3825 "HandoffContext",
3826 "ctx-1",
3827 &lazy_sid,
3828 OWNER,
3829 handoff_context("h1"),
3830 ),
3831 None,
3832 )
3833 .await
3834 .expect("context accepted");
3835
3836 let expected = vec![
3837 (EntryKind::Incoming, "SessionStart".to_string()),
3838 (EntryKind::Incoming, "HandoffOffer".to_string()),
3839 (EntryKind::Incoming, "HandoffAccept".to_string()),
3840 (EntryKind::Checkpoint, "Checkpoint".to_string()),
3841 ];
3842 assert_eq!(shape(&eager, &eager_sid).await, expected, "eager sweep");
3843
3844 let lazy_shape = shape(&lazy, &lazy_sid).await;
3845 assert_eq!(
3846 lazy_shape[..4],
3847 expected[..],
3848 "the lazy path must checkpoint in the same place as the eager one"
3849 );
3850 assert_eq!(
3853 lazy_shape[4..],
3854 [(EntryKind::Incoming, "HandoffContext".to_string())],
3855 "no second checkpoint for the trigger's own append"
3856 );
3857 }
3858}