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(
283 message_type: &str,
284 payload: &[u8],
285 session_id: &str,
286 mode: &str,
287 at_ms: i64,
288 ) -> LogEntry {
289 LogEntry {
290 message_id: String::new(),
291 received_at_ms: at_ms,
292 sender: "_runtime".into(),
293 message_type: message_type.into(),
294 raw_payload: payload.to_vec(),
295 entry_kind: EntryKind::Internal,
296 session_id: session_id.into(),
297 mode: mode.into(),
298 macp_version: macp_core::MACP_VERSION.into(),
299 timestamp_unix_ms: at_ms,
300 bound_mode_version: None,
301 semantics_rev: 0,
302 bound_max_suspend_ms: None,
303 compacted_incoming_ordinals: 0,
304 }
305 }
306
307 async fn save_session_to_storage(&self, session: &Session) {
308 if let Err(err) = self.storage.save_session(session).await {
309 tracing::warn!(
310 session_id = %session.session_id,
311 error = %err,
312 "failed to persist session snapshot"
313 );
314 }
315 }
316
317 async fn maybe_expire_session(
318 &self,
319 session_id: &str,
320 session: &mut Session,
321 ) -> Result<bool, MacpError> {
322 let now = Utc::now().timestamp_millis();
323 let expires = (session.state == SessionState::Open && now > session.ttl_expiry)
326 || (session.state == SessionState::Suspended && session.suspend_cap_exceeded(now));
327 if expires {
328 let entry =
329 Self::make_internal_entry("TtlExpired", b"", session_id, &session.mode, now);
330 self.storage
331 .append_log_entry(session_id, &entry)
332 .await
333 .map_err(|_| MacpError::StorageFailed)?;
334 self.log_store.append(session_id, entry).await;
335 session.state = SessionState::Expired;
336 session.suspended_at_ms = None;
337 self.metrics.record_session_expired(&session.mode);
338 tracing::info!(session_id, "session expired via TTL");
339 let _ = self
340 .session_lifecycle_bus
341 .send(SessionLifecycleEvent::Expired {
342 session_id: session_id.to_string(),
343 });
344 return Ok(true);
345 }
346 Ok(false)
347 }
348
349 pub async fn process(
350 &self,
351 env: &Envelope,
352 max_open_sessions: Option<usize>,
353 ) -> Result<ProcessResult, MacpError> {
354 match env.message_type.as_str() {
355 "SessionStart" => self.process_session_start(env, max_open_sessions).await,
356 "Signal" | "Progress" => self.process_signal(env).await,
357 _ => self.process_message(env).await,
358 }
359 }
360
361 async fn process_session_start(
362 &self,
363 env: &Envelope,
364 max_open_sessions: Option<usize>,
365 ) -> Result<ProcessResult, MacpError> {
366 if env.mode.trim().is_empty() {
367 return Err(MacpError::InvalidEnvelope);
368 }
369 validate_session_id_for_acceptance(&env.session_id)?;
370 let mode_name = env.mode.as_str();
371 let mode = self
372 .mode_registry
373 .get_mode(mode_name)
374 .ok_or(MacpError::UnknownMode)?;
375
376 let start_payload = parse_session_start_payload(&env.payload)?;
377 let require_complete_start = self.mode_registry.requires_strict_session_start(mode_name);
385 if require_complete_start {
386 validate_canonical_session_start_payload_for_mode(mode_name, &start_payload)?;
387 }
388
389 let descriptor_version = self.mode_registry.get_mode_version(mode_name);
397 if let Some(descriptor_version) = &descriptor_version {
398 if !start_payload.mode_version.is_empty()
399 && &start_payload.mode_version != descriptor_version
400 {
401 tracing::warn!(
402 mode = mode_name,
403 payload_version = %start_payload.mode_version,
404 descriptor_version = %descriptor_version,
405 "mode_version mismatch"
406 );
407 return Err(MacpError::InvalidEnvelope);
408 }
409 }
410 let bound_mode_version: Option<String> = if start_payload.mode_version.is_empty() {
411 descriptor_version
412 } else {
413 None
414 };
415 let effective_mode_version = bound_mode_version
416 .clone()
417 .unwrap_or_else(|| start_payload.mode_version.clone());
418
419 let ttl_ms = extract_ttl_ms(&start_payload)?;
420
421 if let Some(existing) = self.registry.get_shared(&env.session_id).await {
426 let existing = existing.lock().await;
427 if existing.seen_message_ids.contains(&env.message_id) {
428 return Ok(ProcessResult {
429 session_state: existing.state.clone(),
430 duplicate: true,
431 });
432 }
433 return Err(MacpError::SessionAlreadyExists);
434 }
435
436 let effective_policy_version = if start_payload.policy_version.is_empty() {
441 crate::policy::defaults::DEFAULT_POLICY_ID.to_string()
442 } else {
443 start_payload.policy_version.clone()
444 };
445 let policy_definition = match self.policy_registry.resolve(&effective_policy_version) {
446 Ok(policy) => {
447 if policy.mode != "*" && policy.mode != mode_name {
449 return Err(MacpError::InvalidPolicyDefinition);
450 }
451 Some(policy)
452 }
453 Err(_) => {
454 return Err(MacpError::UnknownPolicyVersion);
455 }
456 };
457
458 let accepted_at = Utc::now().timestamp_millis();
459 let ttl_base = if env.timestamp_unix_ms > 0 {
463 env.timestamp_unix_ms
464 } else {
465 accepted_at
466 };
467 let ttl_expiry = ttl_base.saturating_add(ttl_ms);
468 let bound_max_suspend_ms = if start_payload.max_suspend_ms > 0 {
473 start_payload.max_suspend_ms
474 } else {
475 macp_core::session::MAX_SUSPEND_MS
476 };
477 let session = Session::builder(env.session_id.clone(), mode_name, env.sender.clone())
478 .ttl_expiry(ttl_expiry)
479 .ttl_ms(ttl_ms)
480 .max_suspend_ms(bound_max_suspend_ms)
481 .started_at_unix_ms(accepted_at)
482 .participants(start_payload.participants.clone())
483 .intent(start_payload.intent.clone())
484 .mode_version(effective_mode_version)
485 .configuration_version(start_payload.configuration_version.clone())
486 .policy_version(effective_policy_version)
487 .context_id(start_payload.context_id.clone())
488 .extensions(start_payload.extensions.clone())
489 .roots(start_payload.roots.clone())
490 .policy_definition(policy_definition)
491 .build();
492
493 let response = mode.on_session_start(&session, env)?;
494 mode.validate_client_envelope(&session, env)?;
501 let semantics_rev = session.semantics_rev;
502
503 let shared = std::sync::Arc::new(tokio::sync::Mutex::new(session));
508 let mut session_guard = shared
511 .clone()
512 .try_lock_owned()
513 .expect("freshly created mutex is uncontended");
514 {
515 let mut map = self.registry.sessions.write().await;
516 if map.contains_key(&env.session_id) {
517 return Err(MacpError::SessionAlreadyExists);
519 }
520 if let Some(max_open) = max_open_sessions {
521 let now = Utc::now().timestamp_millis();
522 let mut count = 0usize;
523 for arc in map.values() {
524 let counts = match arc.try_lock() {
529 Ok(s) => {
530 s.initiator_sender == env.sender
531 && s.state == SessionState::Open
532 && now <= s.ttl_expiry
533 }
534 Err(_) => true,
535 };
536 if counts {
537 count += 1;
538 }
539 }
540 if count >= max_open {
541 return Err(MacpError::RateLimited);
542 }
543 }
544 map.insert(env.session_id.clone(), std::sync::Arc::clone(&shared));
545 }
546
547 let rollback = |runtime: &Self, session_guard: &mut Session| {
552 session_guard.state = SessionState::Expired;
553 let registry = std::sync::Arc::clone(&runtime.registry);
554 let sid = env.session_id.clone();
555 async move {
556 let mut map = registry.sessions.write().await;
557 map.remove(&sid);
558 }
559 };
560
561 if self
563 .storage
564 .create_session_storage(&env.session_id)
565 .await
566 .is_err()
567 {
568 rollback(self, &mut session_guard).await;
569 return Err(MacpError::StorageFailed);
570 }
571 let mut incoming_entry = Self::make_incoming_entry(env, accepted_at);
572 incoming_entry.bound_mode_version = bound_mode_version;
573 incoming_entry.semantics_rev = semantics_rev;
574 incoming_entry.bound_max_suspend_ms = Some(bound_max_suspend_ms);
575 if self
576 .storage
577 .append_log_entry(&env.session_id, &incoming_entry)
578 .await
579 .is_err()
580 {
581 rollback(self, &mut session_guard).await;
582 return Err(MacpError::StorageFailed);
583 }
584
585 self.log_store.create_session_log(&env.session_id).await;
587 self.log_store.append(&env.session_id, incoming_entry).await;
588
589 session_guard
590 .seen_message_ids
591 .insert(env.message_id.clone());
592 session_guard.apply_mode_response(response);
593
594 let result_state = session_guard.state.clone();
595 if let Err(err) = self.storage.save_session(&session_guard).await {
604 tracing::warn!(
605 session_id = %session_guard.session_id,
606 error = %err,
607 "failed to persist session snapshot at SessionStart (recoverable via replay)"
608 );
609 }
610 self.metrics.record_session_start(mode_name);
611 tracing::info!(
612 session_id = %env.session_id,
613 mode = mode_name,
614 sender = %env.sender,
615 "session started"
616 );
617 self.publish_accepted_envelope(env);
623 drop(session_guard);
624 let _ = self
625 .session_lifecycle_bus
626 .send(SessionLifecycleEvent::Created {
627 session_id: env.session_id.clone(),
628 });
629
630 Ok(ProcessResult {
631 session_state: result_state,
632 duplicate: false,
633 })
634 }
635
636 async fn synthesize_due_accept(
722 &self,
723 session_id: &str,
724 session: &mut Session,
725 now_ms: i64,
726 ) -> Result<bool, MacpError> {
727 if session.state != SessionState::Open {
728 return Ok(false);
729 }
730 let Some(mode) = self.mode_registry.get_mode(&session.mode) else {
731 return Ok(false);
732 };
733 let Some(syn) = mode.due_synthetic_envelope(session, now_ms) else {
734 return Ok(false);
735 };
736 if session.seen_message_ids.contains(&syn.message_id) {
740 return Ok(false);
741 }
742 mode.authorize_sender(session, &syn)?;
746 let response = mode.on_message_at(
747 session,
748 &syn,
749 &macp_core::mode::MessageContext::new(syn.timestamp_unix_ms),
750 )?;
751
752 let entry = Self::make_incoming_entry(&syn, syn.timestamp_unix_ms);
758 self.storage
759 .append_log_entry(session_id, &entry)
760 .await
761 .map_err(|_| MacpError::StorageFailed)?;
762 self.log_store.append(session_id, entry).await;
763
764 session.seen_message_ids.insert(syn.message_id.clone());
765 session.apply_mode_response(response);
766 self.metrics.record_message_accepted(&session.mode);
767
768 tracing::info!(
769 session_id = %session_id,
770 message_type = %syn.message_type,
771 message_id = %syn.message_id,
772 sender = %syn.sender,
773 deadline_ms = syn.timestamp_unix_ms,
774 "synthetic envelope appended to accepted history"
775 );
776
777 self.save_session_to_storage(session).await;
778 self.maybe_insert_checkpoint(session_id, session).await;
795 self.publish_accepted_envelope(&syn);
796 Ok(true)
797 }
798
799 async fn process_message(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
807 let shared = self
814 .registry
815 .get_shared(&env.session_id)
816 .await
817 .ok_or(MacpError::UnknownSession)?;
818 let mut session_guard = shared.lock().await;
819 let session = &mut *session_guard;
820
821 let now_ms = chrono::Utc::now().timestamp_millis();
828 match macp_modes::step::check_preconditions(session, env, now_ms)? {
829 macp_modes::step::Precheck::Duplicate => {
830 return Ok(ProcessResult {
831 session_state: session.state.clone(),
832 duplicate: true,
833 });
834 }
835 macp_modes::step::Precheck::Expired => {
836 let expired = self.maybe_expire_session(&env.session_id, session).await?;
842 debug_assert!(expired, "check_preconditions reported Expired");
843 self.save_session_to_storage(session).await;
844 return Err(MacpError::TtlExpired);
845 }
846 macp_modes::step::Precheck::Proceed => {}
847 }
848
849 let mode = self
850 .mode_registry
851 .get_mode(&session.mode)
852 .ok_or(MacpError::UnknownMode)?;
853 mode.authorize_sender(session, env)?;
854 mode.validate_client_envelope(session, env)?;
866 let accepted_at_ms = Utc::now().timestamp_millis();
869 self.synthesize_due_accept(&env.session_id, session, accepted_at_ms)
877 .await?;
878 let response = mode.on_message_at(
879 session,
880 env,
881 &macp_core::mode::MessageContext::new(accepted_at_ms),
882 )?;
883
884 let incoming_entry = Self::make_incoming_entry(env, accepted_at_ms);
886 self.storage
887 .append_log_entry(&env.session_id, &incoming_entry)
888 .await
889 .map_err(|_| MacpError::StorageFailed)?;
890
891 self.log_store.append(&env.session_id, incoming_entry).await;
895 let result_state = macp_modes::step::commit(session, env, response, now_ms);
896
897 self.metrics.record_message_accepted(&session.mode);
898 if env.message_type == "Commitment" {
899 self.metrics.record_commitment_accepted(&session.mode);
900 }
901
902 if Self::audit_verbose(session) {
907 tracing::info!(
908 session_id = %env.session_id,
909 message_type = %env.message_type,
910 sender = %env.sender,
911 state = ?result_state,
912 "message accepted (audit)"
913 );
914 } else {
915 tracing::debug!(
916 session_id = %env.session_id,
917 message_type = %env.message_type,
918 sender = %env.sender,
919 state = ?result_state,
920 "message accepted"
921 );
922 }
923
924 if result_state == SessionState::Resolved {
925 self.metrics.record_session_resolved(&session.mode);
926 tracing::info!(session_id = %env.session_id, mode = %session.mode, "session resolved");
927 let _ = self
928 .session_lifecycle_bus
929 .send(SessionLifecycleEvent::Resolved {
930 session_id: env.session_id.clone(),
931 });
932 }
933
934 self.save_session_to_storage(session).await;
936 if result_state == SessionState::Resolved {
937 if !self.maybe_compact_log(&env.session_id, session).await {
938 self.force_insert_checkpoint(&env.session_id, session).await;
939 }
940 } else {
941 self.maybe_insert_checkpoint(&env.session_id, session).await;
942 }
943 self.publish_accepted_envelope(env);
944
945 Ok(ProcessResult {
946 session_state: result_state,
947 duplicate: false,
948 })
949 }
950
951 async fn process_signal(&self, env: &Envelope) -> Result<ProcessResult, MacpError> {
955 if env.message_type == "Signal" && !env.payload.is_empty() {
958 let signal: crate::pb::SignalPayload =
959 prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
960 if signal.signal_type.trim().is_empty() {
961 return Err(MacpError::InvalidPayload);
962 }
963 }
964 if env.message_type == "Progress" && !env.payload.is_empty() {
966 let _: crate::pb::ProgressPayload =
967 prost::Message::decode(&*env.payload).map_err(|_| MacpError::InvalidPayload)?;
968 }
969 tracing::debug!(
970 sender = %env.sender,
971 message_id = %env.message_id,
972 message_type = %env.message_type,
973 "signal received"
974 );
975 let _ = self.signal_bus.send(env.clone());
976 Ok(ProcessResult {
977 session_state: SessionState::Open,
978 duplicate: false,
979 })
980 }
981
982 pub async fn get_session_checked(&self, session_id: &str) -> Option<Session> {
983 let shared = self.registry.get_shared(session_id).await?;
984 let mut session = shared.lock().await;
985 let changed = self
986 .maybe_expire_session(session_id, &mut session)
987 .await
988 .unwrap_or(false);
989 if changed {
990 self.save_session_to_storage(&session).await;
991 }
992 Some(session.clone())
993 }
994
995 pub async fn cancel_session(
999 &self,
1000 session_id: &str,
1001 reason: &str,
1002 cancelled_by: &str,
1003 ) -> Result<ProcessResult, MacpError> {
1004 let shared = self
1005 .registry
1006 .get_shared(session_id)
1007 .await
1008 .ok_or(MacpError::UnknownSession)?;
1009 let mut session_guard = shared.lock().await;
1010 let session = &mut *session_guard;
1011
1012 self.maybe_expire_session(session_id, session).await?;
1013
1014 if session.state.is_terminal() {
1017 let result_state = session.state.clone();
1018 self.save_session_to_storage(session).await;
1019 return Ok(ProcessResult {
1020 session_state: result_state,
1021 duplicate: false,
1022 });
1023 }
1024
1025 let now_ms = Utc::now().timestamp_millis();
1028 let cancel_payload = crate::pb::SessionCancelPayload {
1029 reason: reason.to_string(),
1030 cancelled_by: cancelled_by.to_string(),
1031 };
1032 let cancel_entry = Self::make_internal_entry(
1033 "SessionCancel",
1034 &prost::Message::encode_to_vec(&cancel_payload),
1035 session_id,
1036 &session.mode,
1037 now_ms,
1038 );
1039 self.storage
1040 .append_log_entry(session_id, &cancel_entry)
1041 .await
1042 .map_err(|_| MacpError::StorageFailed)?;
1043 self.log_store.append(session_id, cancel_entry).await;
1044 let _ = session.cancel();
1047 self.save_session_to_storage(session).await;
1048 if !self.maybe_compact_log(session_id, session).await {
1049 self.force_insert_checkpoint(session_id, session).await;
1050 }
1051 self.metrics.record_session_cancelled(&session.mode);
1052 tracing::info!(session_id, reason, "session cancelled");
1053 let _ = self
1054 .session_lifecycle_bus
1055 .send(SessionLifecycleEvent::Cancelled {
1056 session_id: session_id.to_string(),
1057 });
1058
1059 Ok(ProcessResult {
1060 session_state: SessionState::Cancelled,
1061 duplicate: false,
1062 })
1063 }
1064
1065 pub async fn suspend_session(
1069 &self,
1070 session_id: &str,
1071 reason: &str,
1072 suspended_by: &str,
1073 ) -> Result<ProcessResult, MacpError> {
1074 let shared = self
1075 .registry
1076 .get_shared(session_id)
1077 .await
1078 .ok_or(MacpError::UnknownSession)?;
1079 let mut session_guard = shared.lock().await;
1080 let session = &mut *session_guard;
1081
1082 self.maybe_expire_session(session_id, session).await?;
1083 if session.state != SessionState::Open {
1084 return Err(MacpError::SessionNotOpen);
1085 }
1086
1087 let now_ms = chrono::Utc::now().timestamp_millis();
1088 let payload = crate::pb::SessionSuspendPayload {
1089 reason: reason.to_string(),
1090 suspended_by: suspended_by.to_string(),
1091 };
1092 let entry = Self::make_internal_entry(
1093 "SessionSuspend",
1094 &prost::Message::encode_to_vec(&payload),
1095 session_id,
1096 &session.mode,
1097 now_ms,
1098 );
1099 self.storage
1100 .append_log_entry(session_id, &entry)
1101 .await
1102 .map_err(|_| MacpError::StorageFailed)?;
1103 self.log_store.append(session_id, entry).await;
1104 session.suspend(now_ms)?;
1105 self.save_session_to_storage(session).await;
1106 self.metrics.record_session_suspended(&session.mode);
1107 tracing::info!(session_id, reason, "session suspended");
1108 let _ = self
1109 .session_lifecycle_bus
1110 .send(SessionLifecycleEvent::Suspended {
1111 session_id: session_id.to_string(),
1112 });
1113
1114 Ok(ProcessResult {
1115 session_state: SessionState::Suspended,
1116 duplicate: false,
1117 })
1118 }
1119
1120 pub async fn resume_session(
1124 &self,
1125 session_id: &str,
1126 reason: &str,
1127 resumed_by: &str,
1128 ) -> Result<ProcessResult, MacpError> {
1129 let shared = self
1130 .registry
1131 .get_shared(session_id)
1132 .await
1133 .ok_or(MacpError::UnknownSession)?;
1134 let mut session_guard = shared.lock().await;
1135 let session = &mut *session_guard;
1136
1137 if session.state != SessionState::Suspended {
1138 return Err(MacpError::SessionNotOpen);
1139 }
1140
1141 let now_ms = chrono::Utc::now().timestamp_millis();
1142 let banked_before = session
1154 .suspended_at_ms
1155 .map(|at| (now_ms - at).max(0))
1156 .unwrap_or(0);
1157 let payload = crate::pb::SessionResumePayload {
1158 reason: reason.to_string(),
1159 resumed_by: resumed_by.to_string(),
1160 banked_ms: banked_before,
1161 };
1162 let entry = Self::make_internal_entry(
1163 "SessionResume",
1164 &prost::Message::encode_to_vec(&payload),
1165 session_id,
1166 &session.mode,
1167 now_ms,
1168 );
1169 self.storage
1170 .append_log_entry(session_id, &entry)
1171 .await
1172 .map_err(|_| MacpError::StorageFailed)?;
1173 self.log_store.append(session_id, entry).await;
1174
1175 match session.resume(now_ms) {
1177 Ok(()) => {
1178 self.save_session_to_storage(session).await;
1179 self.metrics.record_session_resumed(&session.mode);
1180 tracing::info!(session_id, reason, "session resumed");
1181 let _ = self
1182 .session_lifecycle_bus
1183 .send(SessionLifecycleEvent::Resumed {
1184 session_id: session_id.to_string(),
1185 });
1186 Ok(ProcessResult {
1187 session_state: SessionState::Open,
1188 duplicate: false,
1189 })
1190 }
1191 Err(_) => {
1192 let cycle_cap_exceeded = session.suspension_intervals.len() > MAX_SUSPENSION_CYCLES;
1201 let duration_cap_exceeded =
1202 session.accumulated_suspended_ms > session.effective_max_suspend_ms();
1203 tracing::warn!(
1204 session_id,
1205 cycle_cap_exceeded,
1206 duration_cap_exceeded,
1207 suspension_cycles = session.suspension_intervals.len(),
1208 accumulated_suspended_ms = session.accumulated_suspended_ms,
1209 "session force-expired: suspension cap exceeded"
1210 );
1211 self.save_session_to_storage(session).await;
1212 self.metrics.record_session_expired(&session.mode);
1213 let _ = self
1214 .session_lifecycle_bus
1215 .send(SessionLifecycleEvent::Expired {
1216 session_id: session_id.to_string(),
1217 });
1218 Err(MacpError::TtlExpired)
1219 }
1220 }
1221 }
1222
1223 async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
1226 let discarded = match self.log_store.get_log(session_id).await {
1230 Some(entries) => {
1231 let prior_base: u64 = entries
1232 .iter()
1233 .filter(|e| e.entry_kind == EntryKind::Checkpoint)
1234 .map(|e| e.compacted_incoming_ordinals)
1235 .max()
1236 .unwrap_or(0);
1237 prior_base
1238 + entries
1239 .iter()
1240 .filter(|e| e.entry_kind == EntryKind::Incoming)
1241 .count() as u64
1242 }
1243 None => 0,
1244 };
1245 match crate::storage::compaction::compact_session_log(
1246 &*self.storage,
1247 session_id,
1248 session,
1249 discarded,
1250 )
1251 .await
1252 {
1253 Ok(checkpoint) => {
1254 self.log_store
1258 .replace_session_log(session_id, vec![checkpoint])
1259 .await;
1260 true
1261 }
1262 Err(e) => {
1263 tracing::debug!(
1264 session_id,
1265 error = %e,
1266 "log compaction skipped (backend may not support it)"
1267 );
1268 false
1269 }
1270 }
1271 }
1272
1273 async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
1276 let persisted = crate::registry::PersistedSession::from(session);
1277 let raw_payload = match serde_json::to_vec(&persisted) {
1278 Ok(bytes) => bytes,
1279 Err(e) => {
1280 tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
1281 return;
1282 }
1283 };
1284 let now = Utc::now().timestamp_millis();
1285 let checkpoint = LogEntry {
1286 message_id: String::new(),
1287 received_at_ms: now,
1288 sender: "_runtime".into(),
1289 message_type: "Checkpoint".into(),
1290 raw_payload,
1291 entry_kind: EntryKind::Checkpoint,
1292 session_id: session_id.into(),
1293 mode: session.mode.clone(),
1294 macp_version: String::new(),
1295 timestamp_unix_ms: now,
1296 bound_mode_version: None,
1297 semantics_rev: 0,
1298 bound_max_suspend_ms: None,
1299 compacted_incoming_ordinals: 0,
1300 };
1301 if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
1302 tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
1303 return;
1304 }
1305 self.log_store.append(session_id, checkpoint).await;
1306 tracing::debug!(
1307 session_id,
1308 "forced checkpoint inserted for terminal session"
1309 );
1310 }
1311
1312 async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
1322 if self.checkpoint_interval == 0 {
1323 return;
1324 }
1325 let log_len = self
1326 .log_store
1327 .get_log(session_id)
1328 .await
1329 .map(|l| l.len())
1330 .unwrap_or(0);
1331 if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
1333 return;
1334 }
1335 self.force_insert_checkpoint(session_id, session).await;
1336 tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
1337 }
1338
1339 pub async fn cleanup_expired_sessions(&self) {
1343 let now = Utc::now().timestamp_millis();
1344 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1349 let guard = self.registry.sessions.read().await;
1350 guard
1351 .iter()
1352 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1353 .collect()
1354 };
1355
1356 let mut expired_count = 0usize;
1357 for (session_id, shared) in candidates {
1358 let mut session = shared.lock().await;
1359 if session.state != SessionState::Open || now <= session.ttl_expiry {
1360 continue;
1361 }
1362 let entry =
1363 Self::make_internal_entry("TtlExpired", b"", &session_id, &session.mode, now);
1364 if let Err(e) = self.storage.append_log_entry(&session_id, &entry).await {
1365 tracing::warn!(
1366 session_id,
1367 error = %e,
1368 "failed to write TTL expiry during cleanup"
1369 );
1370 continue;
1371 }
1372 self.log_store.append(&session_id, entry).await;
1373 session.state = SessionState::Expired;
1374 self.metrics.record_session_expired(&session.mode);
1375 self.save_session_to_storage(&session).await;
1376 if !self.maybe_compact_log(&session_id, &session).await {
1377 self.force_insert_checkpoint(&session_id, &session).await;
1378 }
1379 expired_count += 1;
1380 tracing::info!(session_id = %session_id, "session expired via background cleanup");
1381 let _ = self
1382 .session_lifecycle_bus
1383 .send(SessionLifecycleEvent::Expired {
1384 session_id: session_id.clone(),
1385 });
1386 }
1387
1388 if expired_count > 0 {
1389 tracing::info!(count = expired_count, "background cleanup expired sessions");
1390 }
1391 }
1392
1393 pub async fn sweep_due_synthetic_accepts(&self) -> usize {
1453 let now = Utc::now().timestamp_millis();
1454 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1458 let guard = self.registry.sessions.read().await;
1459 guard
1460 .iter()
1461 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1462 .collect()
1463 };
1464
1465 let mut emitted = 0usize;
1466 for (session_id, shared) in candidates {
1467 let mut session = shared.lock().await;
1468 if session.state != SessionState::Open {
1479 continue;
1480 }
1481 match self
1482 .synthesize_due_accept(&session_id, &mut session, now)
1483 .await
1484 {
1485 Ok(true) => emitted += 1,
1486 Ok(false) => {}
1487 Err(e) => {
1493 tracing::warn!(
1494 session_id = %session_id,
1495 error = %e,
1496 "eager sweep could not emit a due synthetic envelope"
1497 );
1498 }
1499 }
1500 }
1501
1502 if emitted > 0 {
1503 tracing::info!(
1504 count = emitted,
1505 "eager sweep appended due synthetic envelopes"
1506 );
1507 }
1508 emitted
1509 }
1510
1511 pub async fn gc_disk_sessions(&self, retention_secs: u64) -> usize {
1519 let now = Utc::now().timestamp_millis();
1520 let cutoff = now - (retention_secs as i64 * 1000);
1521 let ids = match self.storage.list_session_ids().await {
1522 Ok(ids) => ids,
1523 Err(e) => {
1524 tracing::warn!(error = %e, "disk GC: cannot list sessions");
1525 return 0;
1526 }
1527 };
1528 let mut removed = 0usize;
1529 for id in ids {
1530 let eligible = if let Some(shared) = self.registry.get_shared(&id).await {
1533 let s = shared.lock().await;
1534 s.state.is_terminal() && s.started_at_unix_ms < cutoff
1535 } else {
1536 match self.storage.load_session(&id).await {
1537 Ok(Some(s)) => s.state.is_terminal() && s.started_at_unix_ms < cutoff,
1538 _ => false,
1541 }
1542 };
1543 if !eligible {
1544 continue;
1545 }
1546 match self.storage.delete_session(&id).await {
1547 Ok(()) => {
1548 {
1549 let mut guard = self.registry.sessions.write().await;
1550 guard.remove(&id);
1551 }
1552 self.log_store.remove_session_log(&id).await;
1553 let _ = self.stream_bus.remove_if_unused(&id);
1554 removed += 1;
1555 }
1556 Err(e) => {
1557 tracing::warn!(session_id = %id, error = %e, "disk GC: delete failed");
1558 }
1559 }
1560 }
1561 if removed > 0 {
1562 tracing::info!(count = removed, "disk GC removed terminal sessions");
1563 }
1564 removed
1565 }
1566
1567 pub async fn evict_stale_sessions(&self, retention_secs: u64) {
1573 let now = Utc::now().timestamp_millis();
1574 let cutoff = now - (retention_secs as i64 * 1000);
1575
1576 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1577 let guard = self.registry.sessions.read().await;
1578 guard
1579 .iter()
1580 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1581 .collect()
1582 };
1583 let mut evict_ids = Vec::new();
1584 for (id, shared) in candidates {
1585 let session = shared.lock().await;
1586 if matches!(
1587 session.state,
1588 SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
1589 ) && session.started_at_unix_ms < cutoff
1590 {
1591 evict_ids.push(id);
1592 }
1593 }
1594
1595 if evict_ids.is_empty() {
1596 return;
1597 }
1598 {
1599 let mut guard = self.registry.sessions.write().await;
1600 for id in &evict_ids {
1601 guard.remove(id);
1602 }
1603 }
1604 for id in &evict_ids {
1605 self.log_store.remove_session_log(id).await;
1606 let _ = self.stream_bus.remove_if_unused(id);
1609 }
1610 tracing::info!(
1611 count = evict_ids.len(),
1612 "evicted stale sessions from memory (registry + log cache + stream bus)"
1613 );
1614 }
1615}
1616
1617#[cfg(test)]
1618mod tests {
1619 use super::*;
1620 use crate::decision_pb::ProposalPayload;
1621 use crate::pb::{CommitmentPayload, SessionStartPayload};
1622 use prost::Message;
1623
1624 fn new_sid() -> String {
1625 uuid::Uuid::new_v4().as_hyphenated().to_string()
1626 }
1627
1628 fn make_runtime() -> Runtime {
1629 let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
1630 let registry = Arc::new(SessionRegistry::new());
1631 let log_store = Arc::new(LogStore::new());
1632 Runtime::new(storage, registry, log_store)
1633 }
1634
1635 fn session_start(participants: Vec<String>) -> Vec<u8> {
1636 SessionStartPayload {
1637 intent: "intent".into(),
1638 participants,
1639 mode_version: "1.0.0".into(),
1640 configuration_version: "cfg-1".into(),
1641 policy_version: String::new(),
1642 ttl_ms: 1_000,
1643 context_id: String::new(),
1644 extensions: std::collections::HashMap::new(),
1645 roots: vec![],
1646 max_suspend_ms: 0,
1647 }
1648 .encode_to_vec()
1649 }
1650
1651 fn env(
1652 mode: &str,
1653 message_type: &str,
1654 message_id: &str,
1655 session_id: &str,
1656 sender: &str,
1657 payload: Vec<u8>,
1658 ) -> Envelope {
1659 Envelope {
1660 macp_version: "1.0".into(),
1661 mode: mode.into(),
1662 message_type: message_type.into(),
1663 message_id: message_id.into(),
1664 session_id: session_id.into(),
1665 sender: sender.into(),
1666 timestamp_unix_ms: Utc::now().timestamp_millis(),
1667 payload,
1668 }
1669 }
1670
1671 #[tokio::test]
1672 async fn standard_session_start_is_strict() {
1673 let rt = make_runtime();
1674 let sid = new_sid();
1675 let bad = SessionStartPayload {
1676 ttl_ms: 0,
1677 ..Default::default()
1678 }
1679 .encode_to_vec();
1680 let err = rt
1681 .process(
1682 &env(
1683 "macp.mode.decision.v1",
1684 "SessionStart",
1685 "m1",
1686 &sid,
1687 "agent://orchestrator",
1688 bad,
1689 ),
1690 None,
1691 )
1692 .await
1693 .unwrap_err();
1694 assert!(matches!(
1695 err,
1696 MacpError::InvalidPayload | MacpError::InvalidTtl
1697 ));
1698 }
1699
1700 #[tokio::test]
1713 async fn a_promoted_mode_still_gets_canonical_session_start_validation() {
1714 let mode_registry = Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
1715 macp_policy::DefaultPolicyEvaluator,
1716 )));
1717 mode_registry
1718 .register_extension(crate::pb::ModeDescriptor {
1719 mode: "ext.promoted.v1".into(),
1720 mode_version: "1.0.0".into(),
1721 title: "Promoted".into(),
1722 description: "promotion target".into(),
1723 determinism_class: "semantic-deterministic".into(),
1724 participant_model: "declared".into(),
1725 message_types: vec!["SessionStart".into(), "Commitment".into()],
1726 terminal_message_types: vec!["Commitment".into()],
1727 ..Default::default()
1728 })
1729 .expect("register extension");
1730 assert_eq!(
1731 mode_registry.promote_mode("ext.promoted.v1", None).unwrap(),
1732 "ext.promoted.v1"
1733 );
1734 assert!(
1735 mode_registry.requires_strict_session_start("ext.promoted.v1"),
1736 "promotion must mark the entry strict"
1737 );
1738 assert!(
1739 !crate::session::requires_strict_session_start("ext.promoted.v1"),
1740 "the core's static list must NOT know this name — that disagreement is the point"
1741 );
1742
1743 let rt = Runtime::with_mode_registry(
1744 Arc::new(crate::storage::MemoryBackend),
1745 Arc::new(SessionRegistry::new()),
1746 Arc::new(LogStore::new()),
1747 mode_registry,
1748 );
1749
1750 let err = rt
1752 .process(
1753 &env(
1754 "ext.promoted.v1",
1755 "SessionStart",
1756 "m1",
1757 &new_sid(),
1758 "agent://orchestrator",
1759 session_start(vec![]),
1760 ),
1761 None,
1762 )
1763 .await
1764 .unwrap_err();
1765 assert_eq!(err.to_string(), "InvalidPayload");
1766
1767 let no_versions = SessionStartPayload {
1770 participants: vec!["agent://fraud".into()],
1771 ttl_ms: 1_000,
1772 ..Default::default()
1773 }
1774 .encode_to_vec();
1775 let err = rt
1776 .process(
1777 &env(
1778 "ext.promoted.v1",
1779 "SessionStart",
1780 "m2",
1781 &new_sid(),
1782 "agent://orchestrator",
1783 no_versions,
1784 ),
1785 None,
1786 )
1787 .await
1788 .unwrap_err();
1789 assert_eq!(err.to_string(), "InvalidPayload");
1790
1791 rt.process(
1794 &env(
1795 "ext.promoted.v1",
1796 "SessionStart",
1797 "m3",
1798 &new_sid(),
1799 "agent://orchestrator",
1800 session_start(vec!["agent://fraud".into()]),
1801 ),
1802 None,
1803 )
1804 .await
1805 .expect("a complete SessionStart must still be accepted for a promoted mode");
1806 }
1807
1808 #[tokio::test]
1809 async fn empty_mode_is_rejected() {
1810 let rt = make_runtime();
1811 let sid = new_sid();
1812 let err = rt
1813 .process(
1814 &env(
1815 "",
1816 "SessionStart",
1817 "m1",
1818 &sid,
1819 "agent://orchestrator",
1820 session_start(vec!["agent://fraud".into()]),
1821 ),
1822 None,
1823 )
1824 .await
1825 .unwrap_err();
1826 assert_eq!(err.to_string(), "InvalidEnvelope");
1827 }
1828
1829 #[tokio::test]
1830 async fn rejected_messages_do_not_enter_dedup_state() {
1831 let rt = make_runtime();
1832 let sid = new_sid();
1833 rt.process(
1834 &env(
1835 "macp.mode.decision.v1",
1836 "SessionStart",
1837 "m1",
1838 &sid,
1839 "agent://orchestrator",
1840 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1841 ),
1842 None,
1843 )
1844 .await
1845 .unwrap();
1846
1847 let bad = rt
1848 .process(
1849 &env(
1850 "macp.mode.decision.v1",
1851 "Proposal",
1852 "m2",
1853 &sid,
1854 "agent://fraud",
1855 b"not-protobuf".to_vec(),
1856 ),
1857 None,
1858 )
1859 .await
1860 .unwrap_err();
1861 assert_eq!(bad.to_string(), "InvalidPayload");
1862
1863 let good = ProposalPayload {
1864 proposal_id: "p1".into(),
1865 option: "step-up".into(),
1866 rationale: "risk".into(),
1867 supporting_data: vec![],
1868 }
1869 .encode_to_vec();
1870 let result = rt
1871 .process(
1872 &env(
1873 "macp.mode.decision.v1",
1874 "Proposal",
1875 "m2",
1876 &sid,
1877 "agent://orchestrator",
1878 good,
1879 ),
1880 None,
1881 )
1882 .await
1883 .unwrap();
1884 assert!(!result.duplicate);
1885 }
1886
1887 #[tokio::test]
1888 async fn get_session_transitions_expired_sessions() {
1889 let rt = make_runtime();
1890 let sid = new_sid();
1891 let payload = SessionStartPayload {
1892 intent: "intent".into(),
1893 participants: vec!["agent://fraud".into()],
1894 mode_version: "1.0.0".into(),
1895 configuration_version: "cfg-1".into(),
1896 policy_version: String::new(),
1897 ttl_ms: 1,
1898 context_id: String::new(),
1899 extensions: std::collections::HashMap::new(),
1900 roots: vec![],
1901 max_suspend_ms: 0,
1902 }
1903 .encode_to_vec();
1904 rt.process(
1905 &env(
1906 "macp.mode.decision.v1",
1907 "SessionStart",
1908 "m1",
1909 &sid,
1910 "agent://orchestrator",
1911 payload,
1912 ),
1913 None,
1914 )
1915 .await
1916 .unwrap();
1917 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1918 let session = rt.get_session_checked(&sid).await.unwrap();
1919 assert_eq!(session.state, SessionState::Expired);
1920 }
1921
1922 #[tokio::test]
1923 async fn multi_round_requires_standard_session_start() {
1924 let rt = make_runtime();
1925 let sid = new_sid();
1926 let payload = SessionStartPayload {
1928 participants: vec!["creator".into(), "other".into()],
1929 ..Default::default()
1930 }
1931 .encode_to_vec();
1932 let err = rt
1933 .process(
1934 &env(
1935 "ext.multi_round.v1",
1936 "SessionStart",
1937 "m1",
1938 &sid,
1939 "creator",
1940 payload,
1941 ),
1942 None,
1943 )
1944 .await
1945 .unwrap_err();
1946 assert!(matches!(
1947 err,
1948 MacpError::InvalidPayload | MacpError::InvalidTtl
1949 ));
1950 }
1951
1952 #[tokio::test]
1953 async fn multi_round_valid_session_start() {
1954 let rt = make_runtime();
1955 let sid = new_sid();
1956 let payload = session_start(vec!["alice".into(), "bob".into()]);
1957 rt.process(
1958 &env(
1959 "ext.multi_round.v1",
1960 "SessionStart",
1961 "m1",
1962 &sid,
1963 "coordinator",
1964 payload,
1965 ),
1966 None,
1967 )
1968 .await
1969 .unwrap();
1970 let session = rt.get_session_checked(&sid).await.unwrap();
1971 assert_eq!(session.mode, "ext.multi_round.v1");
1972 assert_eq!(session.participants, vec!["alice", "bob"]);
1973 }
1974
1975 #[tokio::test]
1976 async fn duplicate_session_start_message_id_returns_duplicate() {
1977 let rt = make_runtime();
1978 let sid = new_sid();
1979 let payload = session_start(vec!["agent://fraud".into()]);
1980 rt.process(
1981 &env(
1982 "macp.mode.decision.v1",
1983 "SessionStart",
1984 "m1",
1985 &sid,
1986 "agent://orchestrator",
1987 payload.clone(),
1988 ),
1989 None,
1990 )
1991 .await
1992 .unwrap();
1993
1994 let result = rt
1995 .process(
1996 &env(
1997 "macp.mode.decision.v1",
1998 "SessionStart",
1999 "m1",
2000 &sid,
2001 "agent://orchestrator",
2002 payload,
2003 ),
2004 None,
2005 )
2006 .await
2007 .unwrap();
2008 assert!(result.duplicate);
2009 }
2010
2011 #[tokio::test]
2012 async fn non_start_mode_mismatch_rejected() {
2013 let rt = make_runtime();
2014 let sid = new_sid();
2015 rt.process(
2016 &env(
2017 "macp.mode.decision.v1",
2018 "SessionStart",
2019 "m1",
2020 &sid,
2021 "agent://orchestrator",
2022 session_start(vec!["agent://fraud".into()]),
2023 ),
2024 None,
2025 )
2026 .await
2027 .unwrap();
2028
2029 let proposal = ProposalPayload {
2030 proposal_id: "p1".into(),
2031 option: "step-up".into(),
2032 rationale: "risk".into(),
2033 supporting_data: vec![],
2034 }
2035 .encode_to_vec();
2036 let err = rt
2037 .process(
2038 &env(
2039 "macp.mode.task.v1",
2040 "Proposal",
2041 "m2",
2042 &sid,
2043 "agent://orchestrator",
2044 proposal,
2045 ),
2046 None,
2047 )
2048 .await
2049 .unwrap_err();
2050 assert_eq!(err.to_string(), "InvalidEnvelope");
2051 }
2052
2053 #[tokio::test]
2054 async fn cancel_idempotent_on_already_expired() {
2055 let rt = make_runtime();
2056 let sid = new_sid();
2057 let payload = SessionStartPayload {
2058 intent: "intent".into(),
2059 participants: vec!["agent://fraud".into()],
2060 mode_version: "1.0.0".into(),
2061 configuration_version: "cfg-1".into(),
2062 policy_version: String::new(),
2063 ttl_ms: 1,
2064 context_id: String::new(),
2065 extensions: std::collections::HashMap::new(),
2066 roots: vec![],
2067 max_suspend_ms: 0,
2068 }
2069 .encode_to_vec();
2070 rt.process(
2071 &env(
2072 "macp.mode.decision.v1",
2073 "SessionStart",
2074 "m1",
2075 &sid,
2076 "agent://orchestrator",
2077 payload,
2078 ),
2079 None,
2080 )
2081 .await
2082 .unwrap();
2083 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2084 let result = rt
2085 .cancel_session(&sid, "cleanup", "agent://orchestrator")
2086 .await
2087 .unwrap();
2088 assert_eq!(result.session_state, SessionState::Expired);
2089 }
2090
2091 #[tokio::test]
2092 async fn accepted_envelopes_are_published_in_order() {
2093 let rt = make_runtime();
2094 let sid = new_sid();
2095 let mut events = rt.subscribe_session_stream(&sid);
2096
2097 let start = env(
2098 "macp.mode.decision.v1",
2099 "SessionStart",
2100 "m1",
2101 &sid,
2102 "agent://orchestrator",
2103 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2104 );
2105 rt.process(&start, None).await.unwrap();
2106 let first = events.recv().await.unwrap();
2107 assert_eq!(first.message_id, "m1");
2108 assert_eq!(first.message_type, "SessionStart");
2109
2110 let proposal = ProposalPayload {
2111 proposal_id: "p1".into(),
2112 option: "step-up".into(),
2113 rationale: "risk".into(),
2114 supporting_data: vec![],
2115 }
2116 .encode_to_vec();
2117 let proposal_env = env(
2118 "macp.mode.decision.v1",
2119 "Proposal",
2120 "m2",
2121 &sid,
2122 "agent://orchestrator",
2123 proposal,
2124 );
2125 rt.process(&proposal_env, None).await.unwrap();
2126 let second = events.recv().await.unwrap();
2127 assert_eq!(second.message_id, "m2");
2128 assert_eq!(second.message_type, "Proposal");
2129 }
2130
2131 #[tokio::test]
2132 async fn commitment_versions_are_carried_into_resolution() {
2133 let rt = make_runtime();
2134 let sid = new_sid();
2135 rt.process(
2136 &env(
2137 "macp.mode.proposal.v1",
2138 "SessionStart",
2139 "m1",
2140 &sid,
2141 "agent://buyer",
2142 session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
2143 ),
2144 None,
2145 )
2146 .await
2147 .unwrap();
2148
2149 let proposal = crate::proposal_pb::ProposalPayload {
2150 proposal_id: "p1".into(),
2151 title: "offer".into(),
2152 summary: "summary".into(),
2153 details: vec![],
2154 tags: vec![],
2155 }
2156 .encode_to_vec();
2157 rt.process(
2158 &env(
2159 "macp.mode.proposal.v1",
2160 "Proposal",
2161 "m2",
2162 &sid,
2163 "agent://seller",
2164 proposal,
2165 ),
2166 None,
2167 )
2168 .await
2169 .unwrap();
2170 let accept = crate::proposal_pb::AcceptPayload {
2171 proposal_id: "p1".into(),
2172 reason: String::new(),
2173 }
2174 .encode_to_vec();
2175 rt.process(
2176 &env(
2177 "macp.mode.proposal.v1",
2178 "Accept",
2179 "m3",
2180 &sid,
2181 "agent://seller",
2182 accept.clone(),
2183 ),
2184 None,
2185 )
2186 .await
2187 .unwrap();
2188 rt.process(
2189 &env(
2190 "macp.mode.proposal.v1",
2191 "Accept",
2192 "m4",
2193 &sid,
2194 "agent://buyer",
2195 accept,
2196 ),
2197 None,
2198 )
2199 .await
2200 .unwrap();
2201 let commitment = CommitmentPayload {
2202 commitment_id: "c1".into(),
2203 action: "proposal.accepted".into(),
2204 authority_scope: "commercial".into(),
2205 reason: "bound".into(),
2206 mode_version: "1.0.0".into(),
2207 policy_version: "policy.default".into(),
2208 configuration_version: "cfg-1".into(),
2209 outcome_positive: true,
2210 supersedes: None,
2211 }
2212 .encode_to_vec();
2213 let result = rt
2214 .process(
2215 &env(
2216 "macp.mode.proposal.v1",
2217 "Commitment",
2218 "m5",
2219 &sid,
2220 "agent://buyer",
2221 commitment,
2222 ),
2223 None,
2224 )
2225 .await
2226 .unwrap();
2227 assert_eq!(result.session_state, SessionState::Resolved);
2228 }
2229
2230 #[tokio::test]
2231 async fn max_open_sessions_enforced_under_write_lock() {
2232 let rt = make_runtime();
2233 let sid1 = new_sid();
2234 let sid2 = new_sid();
2235 let sid3 = new_sid();
2236 rt.process(
2237 &env(
2238 "macp.mode.decision.v1",
2239 "SessionStart",
2240 "m1",
2241 &sid1,
2242 "agent://orchestrator",
2243 session_start(vec!["agent://fraud".into()]),
2244 ),
2245 Some(1),
2246 )
2247 .await
2248 .unwrap();
2249
2250 let err = rt
2251 .process(
2252 &env(
2253 "macp.mode.decision.v1",
2254 "SessionStart",
2255 "m2",
2256 &sid2,
2257 "agent://orchestrator",
2258 session_start(vec!["agent://fraud".into()]),
2259 ),
2260 Some(1),
2261 )
2262 .await
2263 .unwrap_err();
2264 assert!(matches!(err, MacpError::RateLimited));
2265
2266 rt.process(
2267 &env(
2268 "macp.mode.decision.v1",
2269 "SessionStart",
2270 "m3",
2271 &sid3,
2272 "agent://other",
2273 session_start(vec!["agent://fraud".into()]),
2274 ),
2275 Some(1),
2276 )
2277 .await
2278 .unwrap();
2279 }
2280
2281 #[tokio::test]
2282 async fn weak_session_id_rejected() {
2283 let rt = make_runtime();
2284 let err = rt
2285 .process(
2286 &env(
2287 "macp.mode.decision.v1",
2288 "SessionStart",
2289 "m1",
2290 "s1",
2291 "agent://orchestrator",
2292 session_start(vec!["agent://fraud".into()]),
2293 ),
2294 None,
2295 )
2296 .await
2297 .unwrap_err();
2298 assert_eq!(err.to_string(), "InvalidSessionId");
2299 }
2300
2301 #[tokio::test]
2302 async fn log_append_failure_rejects_session_start() {
2303 use std::io;
2304 struct FailingBackend;
2305 #[async_trait::async_trait]
2306 impl StorageBackend for FailingBackend {
2307 async fn save_session(&self, _: &Session) -> io::Result<()> {
2308 Ok(())
2309 }
2310 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2311 Ok(None)
2312 }
2313 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2314 Ok(vec![])
2315 }
2316 async fn delete_session(&self, _: &str) -> io::Result<()> {
2317 Ok(())
2318 }
2319 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2320 Ok(vec![])
2321 }
2322 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2323 Err(io::Error::other("disk full"))
2324 }
2325 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2326 Ok(vec![])
2327 }
2328 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2329 Ok(())
2330 }
2331 }
2332
2333 let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
2334 let registry = Arc::new(SessionRegistry::new());
2335 let log_store = Arc::new(LogStore::new());
2336 let rt = Runtime::new(storage, registry, log_store);
2337 let sid = new_sid();
2338
2339 let err = rt
2340 .process(
2341 &env(
2342 "macp.mode.decision.v1",
2343 "SessionStart",
2344 "m1",
2345 &sid,
2346 "agent://orchestrator",
2347 session_start(vec!["agent://fraud".into()]),
2348 ),
2349 None,
2350 )
2351 .await
2352 .unwrap_err();
2353 assert_eq!(err.to_string(), "StorageFailed");
2354 }
2355
2356 #[tokio::test]
2357 async fn log_append_failure_rejects_in_session_message() {
2358 use std::io;
2359 use std::sync::atomic::{AtomicUsize, Ordering};
2360
2361 struct FailOnSecondAppend {
2362 count: AtomicUsize,
2363 }
2364 #[async_trait::async_trait]
2365 impl StorageBackend for FailOnSecondAppend {
2366 async fn save_session(&self, _: &Session) -> io::Result<()> {
2367 Ok(())
2368 }
2369 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2370 Ok(None)
2371 }
2372 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2373 Ok(vec![])
2374 }
2375 async fn delete_session(&self, _: &str) -> io::Result<()> {
2376 Ok(())
2377 }
2378 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2379 Ok(vec![])
2380 }
2381 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2382 let n = self.count.fetch_add(1, Ordering::SeqCst);
2383 if n >= 1 {
2384 Err(io::Error::other("disk full"))
2385 } else {
2386 Ok(())
2387 }
2388 }
2389 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2390 Ok(vec![])
2391 }
2392 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2393 Ok(())
2394 }
2395 }
2396
2397 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2398 count: AtomicUsize::new(0),
2399 });
2400 let registry = Arc::new(SessionRegistry::new());
2401 let log_store = Arc::new(LogStore::new());
2402 let rt = Runtime::new(storage, registry, log_store);
2403 let sid = new_sid();
2404
2405 rt.process(
2407 &env(
2408 "macp.mode.decision.v1",
2409 "SessionStart",
2410 "m1",
2411 &sid,
2412 "agent://orchestrator",
2413 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2414 ),
2415 None,
2416 )
2417 .await
2418 .unwrap();
2419
2420 let proposal = ProposalPayload {
2422 proposal_id: "p1".into(),
2423 option: "step-up".into(),
2424 rationale: "risk".into(),
2425 supporting_data: vec![],
2426 }
2427 .encode_to_vec();
2428 let err = rt
2429 .process(
2430 &env(
2431 "macp.mode.decision.v1",
2432 "Proposal",
2433 "m2",
2434 &sid,
2435 "agent://orchestrator",
2436 proposal,
2437 ),
2438 None,
2439 )
2440 .await
2441 .unwrap_err();
2442 assert_eq!(err.to_string(), "StorageFailed");
2443
2444 let session = rt.get_session_checked(&sid).await.unwrap();
2446 assert!(!session.seen_message_ids.contains("m2"));
2447 }
2448
2449 #[tokio::test]
2450 async fn cancel_session_fails_if_log_append_fails() {
2451 use std::io;
2452 use std::sync::atomic::{AtomicUsize, Ordering};
2453
2454 struct FailOnSecondAppend {
2455 count: AtomicUsize,
2456 }
2457 #[async_trait::async_trait]
2458 impl StorageBackend for FailOnSecondAppend {
2459 async fn save_session(&self, _: &Session) -> io::Result<()> {
2460 Ok(())
2461 }
2462 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2463 Ok(None)
2464 }
2465 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2466 Ok(vec![])
2467 }
2468 async fn delete_session(&self, _: &str) -> io::Result<()> {
2469 Ok(())
2470 }
2471 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2472 Ok(vec![])
2473 }
2474 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2475 let n = self.count.fetch_add(1, Ordering::SeqCst);
2476 if n >= 1 {
2477 Err(io::Error::other("disk full"))
2478 } else {
2479 Ok(())
2480 }
2481 }
2482 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2483 Ok(vec![])
2484 }
2485 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2486 Ok(())
2487 }
2488 }
2489
2490 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2491 count: AtomicUsize::new(0),
2492 });
2493 let registry = Arc::new(SessionRegistry::new());
2494 let log_store = Arc::new(LogStore::new());
2495 let rt = Runtime::new(storage, registry, log_store);
2496 let sid = new_sid();
2497
2498 rt.process(
2499 &env(
2500 "macp.mode.decision.v1",
2501 "SessionStart",
2502 "m1",
2503 &sid,
2504 "agent://orchestrator",
2505 session_start(vec!["agent://fraud".into()]),
2506 ),
2507 None,
2508 )
2509 .await
2510 .unwrap();
2511
2512 let err = rt
2513 .cancel_session(&sid, "test cancel", "agent://orchestrator")
2514 .await
2515 .unwrap_err();
2516 assert_eq!(err.to_string(), "StorageFailed");
2517 }
2518
2519 #[tokio::test]
2520 async fn ttl_expiration_rejects_message() {
2521 let rt = make_runtime();
2522 let sid = new_sid();
2523 let payload = SessionStartPayload {
2524 intent: "intent".into(),
2525 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
2526 mode_version: "1.0.0".into(),
2527 configuration_version: "cfg-1".into(),
2528 policy_version: String::new(),
2529 ttl_ms: 1,
2530 context_id: String::new(),
2531 extensions: std::collections::HashMap::new(),
2532 roots: vec![],
2533 max_suspend_ms: 0,
2534 }
2535 .encode_to_vec();
2536 rt.process(
2537 &env(
2538 "macp.mode.decision.v1",
2539 "SessionStart",
2540 "m1",
2541 &sid,
2542 "agent://orchestrator",
2543 payload,
2544 ),
2545 None,
2546 )
2547 .await
2548 .unwrap();
2549 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2550 let proposal = ProposalPayload {
2551 proposal_id: "p1".into(),
2552 option: "step-up".into(),
2553 rationale: "risk".into(),
2554 supporting_data: vec![],
2555 }
2556 .encode_to_vec();
2557 let err = rt
2558 .process(
2559 &env(
2560 "macp.mode.decision.v1",
2561 "Proposal",
2562 "m2",
2563 &sid,
2564 "agent://orchestrator",
2565 proposal,
2566 ),
2567 None,
2568 )
2569 .await
2570 .unwrap_err();
2571 assert_eq!(err.to_string(), "TtlExpired");
2572 }
2573
2574 #[tokio::test]
2575 async fn cleanup_expired_sessions_marks_expired() {
2576 let rt = make_runtime();
2577 let sid = new_sid();
2578 let payload = SessionStartPayload {
2579 intent: "intent".into(),
2580 participants: vec!["agent://fraud".into()],
2581 mode_version: "1.0.0".into(),
2582 configuration_version: "cfg-1".into(),
2583 policy_version: String::new(),
2584 ttl_ms: 1,
2585 context_id: String::new(),
2586 extensions: std::collections::HashMap::new(),
2587 roots: vec![],
2588 max_suspend_ms: 0,
2589 }
2590 .encode_to_vec();
2591 rt.process(
2592 &env(
2593 "macp.mode.decision.v1",
2594 "SessionStart",
2595 "m1",
2596 &sid,
2597 "agent://orchestrator",
2598 payload,
2599 ),
2600 None,
2601 )
2602 .await
2603 .unwrap();
2604 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2605 rt.cleanup_expired_sessions().await;
2606 let session = rt.get_session_checked(&sid).await.unwrap();
2607 assert_eq!(session.state, SessionState::Expired);
2608 }
2609
2610 #[tokio::test]
2611 async fn evict_stale_sessions_removes_resolved() {
2612 let rt = make_runtime();
2613 let sid = new_sid();
2614 rt.process(
2616 &env(
2617 "macp.mode.decision.v1",
2618 "SessionStart",
2619 "m1",
2620 &sid,
2621 "agent://orchestrator",
2622 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2623 ),
2624 None,
2625 )
2626 .await
2627 .unwrap();
2628 let proposal = ProposalPayload {
2630 proposal_id: "p1".into(),
2631 option: "step-up".into(),
2632 rationale: "risk".into(),
2633 supporting_data: vec![],
2634 }
2635 .encode_to_vec();
2636 rt.process(
2637 &env(
2638 "macp.mode.decision.v1",
2639 "Proposal",
2640 "m2",
2641 &sid,
2642 "agent://orchestrator",
2643 proposal,
2644 ),
2645 None,
2646 )
2647 .await
2648 .unwrap();
2649 let commitment = CommitmentPayload {
2651 commitment_id: "c1".into(),
2652 action: "decision.selected".into(),
2653 authority_scope: "payments".into(),
2654 reason: "bound".into(),
2655 mode_version: "1.0.0".into(),
2656 policy_version: "policy.default".into(),
2657 configuration_version: "cfg-1".into(),
2658 outcome_positive: true,
2659 supersedes: None,
2660 }
2661 .encode_to_vec();
2662 let result = rt
2663 .process(
2664 &env(
2665 "macp.mode.decision.v1",
2666 "Commitment",
2667 "m3",
2668 &sid,
2669 "agent://orchestrator",
2670 commitment,
2671 ),
2672 None,
2673 )
2674 .await
2675 .unwrap();
2676 assert_eq!(result.session_state, SessionState::Resolved);
2677 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2679 rt.evict_stale_sessions(0).await;
2681 assert!(rt.registry.get_session(&sid).await.is_none());
2683 }
2684
2685 #[tokio::test]
2686 async fn session_start_with_wrong_mode_version_rejected() {
2687 let rt = make_runtime();
2688 let sid = new_sid();
2689 let payload = SessionStartPayload {
2690 intent: "test".into(),
2691 participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
2692 mode_version: "99.0.0".into(), configuration_version: "cfg-1".into(),
2694 policy_version: String::new(),
2695 ttl_ms: 60_000,
2696 context_id: String::new(),
2697 extensions: std::collections::HashMap::new(),
2698 roots: vec![],
2699 max_suspend_ms: 0,
2700 }
2701 .encode_to_vec();
2702
2703 let err = rt
2704 .process(
2705 &env(
2706 "macp.mode.decision.v1",
2707 "SessionStart",
2708 "m1",
2709 &sid,
2710 "agent://orchestrator",
2711 payload,
2712 ),
2713 None,
2714 )
2715 .await
2716 .unwrap_err();
2717 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2718 }
2719
2720 #[tokio::test]
2721 async fn signal_empty_signal_type_rejected() {
2722 let rt = make_runtime();
2723 let signal_payload = crate::pb::SignalPayload {
2725 signal_type: String::new(),
2726 data: b"some data".to_vec(),
2727 confidence: 0.0,
2728 correlation_session_id: String::new(),
2729 }
2730 .encode_to_vec();
2731 let signal = Envelope {
2732 macp_version: "1.0".into(),
2733 mode: String::new(),
2734 message_type: "Signal".into(),
2735 message_id: "sig-1".into(),
2736 session_id: String::new(),
2737 sender: "agent://a".into(),
2738 timestamp_unix_ms: 0,
2739 payload: signal_payload,
2740 };
2741 let err = rt.process_signal(&signal).await.unwrap_err();
2742 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2743 }
2744
2745 #[tokio::test]
2746 async fn signal_valid_payload_accepted() {
2747 let rt = make_runtime();
2748 let signal_payload = crate::pb::SignalPayload {
2749 signal_type: "heartbeat".into(),
2750 data: vec![],
2751 confidence: 0.8,
2752 correlation_session_id: String::new(),
2753 }
2754 .encode_to_vec();
2755 let signal = Envelope {
2756 macp_version: "1.0".into(),
2757 mode: String::new(),
2758 message_type: "Signal".into(),
2759 message_id: "sig-2".into(),
2760 session_id: String::new(),
2761 sender: "agent://a".into(),
2762 timestamp_unix_ms: 0,
2763 payload: signal_payload,
2764 };
2765 rt.process_signal(&signal).await.unwrap();
2766 }
2767
2768 #[tokio::test]
2769 async fn signal_empty_payload_accepted() {
2770 let rt = make_runtime();
2771 let signal = Envelope {
2772 macp_version: "1.0".into(),
2773 mode: String::new(),
2774 message_type: "Signal".into(),
2775 message_id: "sig-3".into(),
2776 session_id: String::new(),
2777 sender: "agent://a".into(),
2778 timestamp_unix_ms: 0,
2779 payload: vec![],
2780 };
2781 rt.process_signal(&signal).await.unwrap();
2782 }
2783
2784 #[tokio::test]
2790 async fn ext_mode_empty_version_binds_descriptor_version() {
2791 let rt = make_runtime();
2792 rt.register_extension(ModeDescriptor {
2793 mode: "ext.dyn.v1".into(),
2794 mode_version: "2.5.0".into(),
2795 message_types: vec!["SessionStart".into(), "Note".into(), "Commitment".into()],
2796 terminal_message_types: vec!["Commitment".into()],
2797 ..Default::default()
2798 })
2799 .unwrap();
2800
2801 let sid = new_sid();
2802 let payload = SessionStartPayload {
2803 participants: vec!["alice".into()],
2804 configuration_version: "cfg-1".into(),
2805 ttl_ms: 60_000,
2806 ..Default::default()
2807 }
2808 .encode_to_vec();
2809 rt.process(
2810 &env("ext.dyn.v1", "SessionStart", "m1", &sid, "alice", payload),
2811 None,
2812 )
2813 .await
2814 .unwrap();
2815
2816 let session = rt.get_session_checked(&sid).await.unwrap();
2818 assert_eq!(session.mode_version, "2.5.0");
2819
2820 let bad = CommitmentPayload {
2822 commitment_id: "c1".into(),
2823 action: "work.completed".into(),
2824 authority_scope: "test".into(),
2825 reason: "done".into(),
2826 mode_version: String::new(),
2827 policy_version: "policy.default".into(),
2828 configuration_version: "cfg-1".into(),
2829 outcome_positive: true,
2830 supersedes: None,
2831 }
2832 .encode_to_vec();
2833 let err = rt
2834 .process(
2835 &env("ext.dyn.v1", "Commitment", "m2", &sid, "alice", bad),
2836 None,
2837 )
2838 .await
2839 .unwrap_err();
2840 assert_eq!(err.to_string(), "InvalidPayload");
2841
2842 let good = CommitmentPayload {
2844 commitment_id: "c1".into(),
2845 action: "work.completed".into(),
2846 authority_scope: "test".into(),
2847 reason: "done".into(),
2848 mode_version: "2.5.0".into(),
2849 policy_version: "policy.default".into(),
2850 configuration_version: "cfg-1".into(),
2851 outcome_positive: true,
2852 supersedes: None,
2853 }
2854 .encode_to_vec();
2855 let result = rt
2856 .process(
2857 &env("ext.dyn.v1", "Commitment", "m3", &sid, "alice", good),
2858 None,
2859 )
2860 .await
2861 .unwrap();
2862 assert_eq!(result.session_state, SessionState::Resolved);
2863 }
2864
2865 #[tokio::test]
2868 async fn ext_mode_binding_recorded_on_session_start_log_entry() {
2869 let rt = make_runtime();
2870 rt.register_extension(ModeDescriptor {
2871 mode: "ext.dyn2.v1".into(),
2872 mode_version: "3.0.0".into(),
2873 message_types: vec!["SessionStart".into(), "Commitment".into()],
2874 terminal_message_types: vec!["Commitment".into()],
2875 ..Default::default()
2876 })
2877 .unwrap();
2878
2879 let sid = new_sid();
2880 let payload = SessionStartPayload {
2881 participants: vec!["alice".into()],
2882 configuration_version: "cfg-1".into(),
2883 ttl_ms: 60_000,
2884 ..Default::default()
2885 }
2886 .encode_to_vec();
2887 rt.process(
2888 &env("ext.dyn2.v1", "SessionStart", "m1", &sid, "alice", payload),
2889 None,
2890 )
2891 .await
2892 .unwrap();
2893
2894 let log = rt.log_store.get_log(&sid).await.unwrap();
2895 assert_eq!(log[0].message_type, "SessionStart");
2896 assert_eq!(log[0].bound_mode_version.as_deref(), Some("3.0.0"));
2897
2898 let sid2 = new_sid();
2900 let payload2 = SessionStartPayload {
2901 participants: vec!["alice".into()],
2902 mode_version: "3.0.0".into(),
2903 configuration_version: "cfg-1".into(),
2904 ttl_ms: 60_000,
2905 ..Default::default()
2906 }
2907 .encode_to_vec();
2908 rt.process(
2909 &env(
2910 "ext.dyn2.v1",
2911 "SessionStart",
2912 "m1",
2913 &sid2,
2914 "alice",
2915 payload2,
2916 ),
2917 None,
2918 )
2919 .await
2920 .unwrap();
2921 let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2922 assert_eq!(log2[0].bound_mode_version, None);
2923 }
2924
2925 #[tokio::test]
2930 async fn session_start_binds_and_records_max_suspend_cap() {
2931 let rt = make_runtime();
2932
2933 let sid = new_sid();
2935 let payload = SessionStartPayload {
2936 participants: vec!["alice".into(), "bob".into()],
2937 mode_version: "1.0.0".into(),
2938 configuration_version: "cfg-1".into(),
2939 ttl_ms: 60_000,
2940 max_suspend_ms: 12_345,
2941 ..Default::default()
2942 }
2943 .encode_to_vec();
2944 rt.process(
2945 &env(
2946 "macp.mode.decision.v1",
2947 "SessionStart",
2948 "m1",
2949 &sid,
2950 "alice",
2951 payload,
2952 ),
2953 None,
2954 )
2955 .await
2956 .unwrap();
2957 let log = rt.log_store.get_log(&sid).await.unwrap();
2958 assert_eq!(log[0].bound_max_suspend_ms, Some(12_345));
2959
2960 let sid2 = new_sid();
2962 let payload2 = SessionStartPayload {
2963 participants: vec!["alice".into(), "bob".into()],
2964 mode_version: "1.0.0".into(),
2965 configuration_version: "cfg-1".into(),
2966 ttl_ms: 60_000,
2967 max_suspend_ms: 0,
2968 ..Default::default()
2969 }
2970 .encode_to_vec();
2971 rt.process(
2972 &env(
2973 "macp.mode.decision.v1",
2974 "SessionStart",
2975 "m2",
2976 &sid2,
2977 "alice",
2978 payload2,
2979 ),
2980 None,
2981 )
2982 .await
2983 .unwrap();
2984 let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2985 assert_eq!(
2986 log2[0].bound_max_suspend_ms,
2987 Some(macp_core::session::MAX_SUSPEND_MS)
2988 );
2989 }
2990
2991 #[test]
2992 fn audit_verbosity_reads_policy_rules() {
2993 let mut session = Session::builder("s1", "macp.mode.decision.v1", "a").build();
2994 assert!(!Runtime::audit_verbose(&session));
2995
2996 session.policy_definition = Some(macp_core::policy::PolicyDefinition {
2997 policy_id: "policy.test.audit".into(),
2998 mode: "*".into(),
2999 description: "audited".into(),
3000 rules: serde_json::json!({ "audit": { "level": "info" } }),
3001 schema_version: 1,
3002 });
3003 assert!(Runtime::audit_verbose(&session));
3004
3005 session.policy_definition.as_mut().unwrap().rules =
3006 serde_json::json!({ "audit": { "level": "debug" } });
3007 assert!(!Runtime::audit_verbose(&session));
3008 }
3009
3010 #[tokio::test]
3016 async fn session_start_snapshot_failure_is_nonfatal_after_commit_point() {
3017 use std::io;
3018
3019 struct FailSnapshotBackend;
3020 #[async_trait::async_trait]
3021 impl StorageBackend for FailSnapshotBackend {
3022 async fn create_session_storage(&self, _s: &str) -> io::Result<()> {
3023 Ok(())
3024 }
3025 async fn save_session(&self, _s: &Session) -> io::Result<()> {
3026 Err(io::Error::other("snapshot disk full"))
3027 }
3028 async fn load_session(&self, _s: &str) -> io::Result<Option<Session>> {
3029 Ok(None)
3030 }
3031 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
3032 Ok(vec![])
3033 }
3034 async fn delete_session(&self, _s: &str) -> io::Result<()> {
3035 Ok(())
3036 }
3037 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
3038 Ok(vec![])
3039 }
3040 async fn append_log_entry(
3041 &self,
3042 _s: &str,
3043 _e: &crate::log_store::LogEntry,
3044 ) -> io::Result<()> {
3045 Ok(())
3046 }
3047 async fn load_log(&self, _s: &str) -> io::Result<Vec<crate::log_store::LogEntry>> {
3048 Ok(vec![])
3049 }
3050 }
3051
3052 let rt = Runtime::new(
3053 Arc::new(FailSnapshotBackend),
3054 Arc::new(SessionRegistry::new()),
3055 Arc::new(LogStore::new()),
3056 );
3057 let sid = new_sid();
3058 let result = rt
3059 .process(
3060 &env(
3061 "macp.mode.decision.v1",
3062 "SessionStart",
3063 "m1",
3064 &sid,
3065 "agent://orchestrator",
3066 session_start(vec!["agent://orchestrator".into()]),
3067 ),
3068 None,
3069 )
3070 .await
3071 .expect("start must succeed: the log append (commit point) succeeded");
3072 assert!(!result.duplicate);
3073 assert!(rt.get_session_checked(&sid).await.is_some());
3075 }
3076
3077 const HANDOFF_MODE: &str = "macp.mode.handoff.v1";
3093 const OWNER: &str = "agent://owner";
3094 const TARGET: &str = "agent://target";
3095
3096 fn reserved_id(handoff_id: &str) -> String {
3097 format!(
3098 "{}{handoff_id}",
3099 crate::mode::handoff::IMPLICIT_ACCEPT_MESSAGE_ID_PREFIX
3100 )
3101 }
3102
3103 fn handoff_start_payload() -> Vec<u8> {
3104 session_start(vec![OWNER.into(), TARGET.into()])
3105 }
3106
3107 fn handoff_commitment_payload() -> Vec<u8> {
3110 CommitmentPayload {
3111 commitment_id: "c1".into(),
3112 action: "handoff.accepted".into(),
3113 authority_scope: "support".into(),
3114 reason: "bound".into(),
3115 mode_version: "1.0.0".into(),
3116 policy_version: "policy.default".into(),
3117 configuration_version: "cfg-1".into(),
3118 outcome_positive: true,
3119 supersedes: None,
3120 }
3121 .encode_to_vec()
3122 }
3123
3124 fn handoff_offer(handoff_id: &str) -> Vec<u8> {
3125 crate::handoff_pb::HandoffOfferPayload {
3126 handoff_id: handoff_id.into(),
3127 target_participant: TARGET.into(),
3128 scope: "support".into(),
3129 reason: "escalate".into(),
3130 }
3131 .encode_to_vec()
3132 }
3133
3134 fn handoff_context(handoff_id: &str) -> Vec<u8> {
3135 crate::handoff_pb::HandoffContextPayload {
3136 handoff_id: handoff_id.into(),
3137 content_type: "text/plain".into(),
3138 context: b"background".to_vec(),
3139 }
3140 .encode_to_vec()
3141 }
3142
3143 fn handoff_accept(handoff_id: &str, implicit: bool) -> Vec<u8> {
3144 crate::handoff_pb::HandoffAcceptPayload {
3145 handoff_id: handoff_id.into(),
3146 accepted_by: TARGET.into(),
3147 reason: "ready".into(),
3148 implicit,
3149 }
3150 .encode_to_vec()
3151 }
3152
3153 async fn handoff_session_with_offer(rt: &Runtime) -> String {
3156 let sid = new_sid();
3157 rt.process(
3158 &env(
3159 HANDOFF_MODE,
3160 "SessionStart",
3161 "start-1",
3162 &sid,
3163 OWNER,
3164 handoff_start_payload(),
3165 ),
3166 None,
3167 )
3168 .await
3169 .expect("handoff session start");
3170 rt.process(
3171 &env(
3172 HANDOFF_MODE,
3173 "HandoffOffer",
3174 "offer-1",
3175 &sid,
3176 OWNER,
3177 handoff_offer("h1"),
3178 ),
3179 None,
3180 )
3181 .await
3182 .expect("handoff offer");
3183 assert_eq!(
3184 rt.get_session_checked(&sid).await.unwrap().semantics_rev,
3185 macp_core::session::CURRENT_SEMANTICS_REV
3186 );
3187 sid
3188 }
3189
3190 #[tokio::test]
3203 async fn reserved_message_id_namespace_is_rejected_at_rev2() {
3204 let rt = make_runtime();
3205 let sid = handoff_session_with_offer(&rt).await;
3206
3207 let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3208 let dedup_before = rt
3209 .get_session_checked(&sid)
3210 .await
3211 .unwrap()
3212 .seen_message_ids
3213 .clone();
3214
3215 for (message_type, sender, payload) in [
3216 ("HandoffContext", OWNER, handoff_context("h1")),
3217 ("Commitment", OWNER, handoff_commitment_payload()),
3218 ("HandoffAccept", TARGET, handoff_accept("h1", false)),
3219 ] {
3220 let err = rt
3221 .process(
3222 &env(
3223 HANDOFF_MODE,
3224 message_type,
3225 &reserved_id("h1"),
3226 &sid,
3227 sender,
3228 payload,
3229 ),
3230 None,
3231 )
3232 .await
3233 .unwrap_err();
3234 assert!(
3235 matches!(err, MacpError::InvalidEnvelope),
3236 "{message_type} with a reserved id must be InvalidEnvelope, got {err}"
3237 );
3238 }
3239
3240 let session = rt.get_session_checked(&sid).await.unwrap();
3242 assert_eq!(
3243 rt.log_store.get_log(&sid).await.unwrap().len(),
3244 history_before
3245 );
3246 assert_eq!(session.seen_message_ids, dedup_before);
3247 assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3248 assert_eq!(session.state, SessionState::Open);
3249
3250 rt.process(
3253 &env(
3254 HANDOFF_MODE,
3255 "HandoffContext",
3256 "ctx-1",
3257 &sid,
3258 OWNER,
3259 handoff_context("h1"),
3260 ),
3261 None,
3262 )
3263 .await
3264 .expect("an ordinary id is accepted");
3265 assert_eq!(
3266 rt.log_store.get_log(&sid).await.unwrap().len(),
3267 history_before + 1
3268 );
3269 }
3270
3271 #[tokio::test]
3282 async fn reserved_message_id_is_rejected_on_the_session_start_path() {
3283 let rt = make_runtime();
3284 let sid = new_sid();
3285
3286 let err = rt
3287 .process(
3288 &env(
3289 HANDOFF_MODE,
3290 "SessionStart",
3291 &reserved_id("h1"),
3292 &sid,
3293 OWNER,
3294 handoff_start_payload(),
3295 ),
3296 None,
3297 )
3298 .await
3299 .unwrap_err();
3300 assert!(
3301 matches!(err, MacpError::InvalidEnvelope),
3302 "reserved id on SessionStart must be InvalidEnvelope, got {err}"
3303 );
3304
3305 assert!(rt.get_session_checked(&sid).await.is_none());
3307 assert!(rt.log_store.get_log(&sid).await.is_none());
3308 assert!(!rt.registry.sessions.read().await.contains_key(&sid));
3309
3310 rt.process(
3313 &env(
3314 HANDOFF_MODE,
3315 "SessionStart",
3316 "start-1",
3317 &sid,
3318 OWNER,
3319 handoff_start_payload(),
3320 ),
3321 None,
3322 )
3323 .await
3324 .expect("a rejected SessionStart must not reserve the session id");
3325 let session = rt.get_session_checked(&sid).await.unwrap();
3326 assert!(session.seen_message_ids.contains("start-1"));
3327 assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3328 }
3329
3330 #[tokio::test]
3360 async fn client_implicit_accept_rejected_through_the_runtime() {
3361 let rt = make_runtime();
3362 let sid = handoff_session_with_offer(&rt).await;
3363 let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3364
3365 let err = rt
3367 .process(
3368 &env(
3369 HANDOFF_MODE,
3370 "HandoffAccept",
3371 "accept-1",
3372 &sid,
3373 TARGET,
3374 handoff_accept("h1", true),
3375 ),
3376 None,
3377 )
3378 .await
3379 .unwrap_err();
3380 assert!(matches!(err, MacpError::InvalidPayload), "got {err}");
3381
3382 let err = rt
3384 .process(
3385 &env(
3386 HANDOFF_MODE,
3387 "HandoffAccept",
3388 &reserved_id("h1"),
3389 &sid,
3390 TARGET,
3391 handoff_accept("h1", true),
3392 ),
3393 None,
3394 )
3395 .await
3396 .unwrap_err();
3397 assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3398
3399 let session = rt.get_session_checked(&sid).await.unwrap();
3401 assert_eq!(
3402 rt.log_store.get_log(&sid).await.unwrap().len(),
3403 history_before
3404 );
3405 assert!(session.seen_message_ids.is_disjoint(
3406 &["accept-1".to_string(), reserved_id("h1")]
3407 .into_iter()
3408 .collect()
3409 ));
3410 let mode_state: serde_json::Value = serde_json::from_slice(&session.mode_state).unwrap();
3411 assert_eq!(mode_state["offers"]["h1"]["disposition"], "Offered");
3412
3413 rt.process(
3417 &env(
3418 HANDOFF_MODE,
3419 "HandoffAccept",
3420 "accept-2",
3421 &sid,
3422 TARGET,
3423 handoff_accept("h1", false),
3424 ),
3425 None,
3426 )
3427 .await
3428 .expect("an explicit accept is still accepted");
3429 }
3430
3431 #[tokio::test]
3437 async fn client_boundary_error_ordering_is_unchanged_at_rev2() {
3438 let rt = make_runtime();
3439 let sid = handoff_session_with_offer(&rt).await;
3440
3441 let err = rt
3443 .process(
3444 &env(
3445 HANDOFF_MODE,
3446 "HandoffAccept",
3447 &reserved_id("h1"),
3448 &sid,
3449 "agent://stranger",
3450 handoff_accept("h1", true),
3451 ),
3452 None,
3453 )
3454 .await
3455 .unwrap_err();
3456 assert!(
3457 matches!(err, MacpError::Forbidden),
3458 "authorization must be reported before the client boundary, got {err}"
3459 );
3460
3461 let err = rt
3463 .process(
3464 &env(
3465 HANDOFF_MODE,
3466 "HandoffAccept",
3467 &reserved_id("h1"),
3468 &sid,
3469 TARGET,
3470 handoff_accept("h1", true),
3471 ),
3472 None,
3473 )
3474 .await
3475 .unwrap_err();
3476 assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3477 }
3478
3479 async fn handoff_session_with_timed_offer(rt: &Runtime, timeout_ms: i64) -> String {
3485 rt.register_policy(macp_core::policy::PolicyDefinition {
3486 policy_id: "handoff-timed".into(),
3487 mode: HANDOFF_MODE.into(),
3488 description: "implicit accept".into(),
3489 rules: serde_json::json!({
3490 "acceptance": { "implicit_accept_timeout_ms": timeout_ms },
3491 "commitment": { "authority": "initiator_only" }
3492 }),
3493 schema_version: 1,
3494 })
3495 .expect("policy registers");
3496
3497 let sid = new_sid();
3498 let start = SessionStartPayload {
3499 intent: "escalate".into(),
3500 participants: vec![OWNER.into(), TARGET.into()],
3501 mode_version: "1.0.0".into(),
3502 configuration_version: "cfg-1".into(),
3503 policy_version: "handoff-timed".into(),
3504 ttl_ms: 60_000,
3505 context_id: String::new(),
3506 extensions: std::collections::HashMap::new(),
3507 roots: vec![],
3508 max_suspend_ms: 0,
3509 }
3510 .encode_to_vec();
3511 rt.process(
3512 &env(HANDOFF_MODE, "SessionStart", "start-1", &sid, OWNER, start),
3513 None,
3514 )
3515 .await
3516 .expect("session start");
3517 rt.process(
3518 &env(
3519 HANDOFF_MODE,
3520 "HandoffOffer",
3521 "offer-1",
3522 &sid,
3523 OWNER,
3524 handoff_offer("h1"),
3525 ),
3526 None,
3527 )
3528 .await
3529 .expect("offer");
3530 sid
3531 }
3532
3533 #[tokio::test]
3553 async fn synthesis_is_skipped_for_a_non_open_session() {
3554 let rt = make_runtime();
3555 let sid = handoff_session_with_timed_offer(&rt, 20).await;
3556 rt.suspend_session(&sid, "hold", OWNER)
3557 .await
3558 .expect("suspend");
3559
3560 let shared = rt.registry.get_shared(&sid).await.unwrap();
3561 let mut guard = shared.lock().await;
3562 let session = &mut *guard;
3563 assert_eq!(session.state, SessionState::Suspended);
3564 assert!(session.suspended_at_ms.is_some());
3565
3566 let long_after = session.suspended_at_ms.unwrap() + 10_000;
3569 let log_before = rt.log_store.get_log(&sid).await.unwrap_or_default().len();
3570 let dedup_before = session.seen_message_ids.len();
3571 let mode_state_before = session.mode_state.clone();
3572
3573 rt.synthesize_due_accept(&sid, session, long_after)
3574 .await
3575 .expect("the filter is a skip, not an error");
3576
3577 assert_eq!(
3578 rt.log_store.get_log(&sid).await.unwrap_or_default().len(),
3579 log_before,
3580 "a suspended session must not gain a synthetic entry"
3581 );
3582 assert_eq!(session.seen_message_ids.len(), dedup_before);
3583 assert_eq!(session.mode_state, mode_state_before);
3584
3585 let mut without_the_tripwire = session.clone();
3597 without_the_tripwire.state = SessionState::Open;
3598 without_the_tripwire.suspended_at_ms = None;
3599 let mode = rt.mode_registry.get_mode(&session.mode).unwrap();
3600 let would_have_emitted = mode
3601 .due_synthetic_envelope(&without_the_tripwire, long_after)
3602 .expect("the mode would have synthesized; only the kernel filter stopped it");
3603 let suspended_at = session.suspended_at_ms.unwrap();
3609 assert!(
3610 would_have_emitted.timestamp_unix_ms >= suspended_at
3611 && would_have_emitted.timestamp_unix_ms < long_after,
3612 "D {} must fall inside the still-open pause starting at {suspended_at}",
3613 would_have_emitted.timestamp_unix_ms
3614 );
3615
3616 drop(guard);
3619 rt.resume_session(&sid, "go", OWNER).await.expect("resume");
3620 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3621 let shared = rt.registry.get_shared(&sid).await.unwrap();
3622 let mut guard = shared.lock().await;
3623 let session = &mut *guard;
3624 let now = Utc::now().timestamp_millis();
3625 rt.synthesize_due_accept(&sid, session, now).await.unwrap();
3626 assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3627 }
3628
3629 #[tokio::test]
3639 async fn synthetic_entry_stamps_received_at_with_the_deadline() {
3640 let rt = make_runtime();
3641 let sid = handoff_session_with_timed_offer(&rt, 20).await;
3642 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3643
3644 let offer_received_at = rt
3645 .log_store
3646 .get_log(&sid)
3647 .await
3648 .unwrap()
3649 .iter()
3650 .find(|e| e.message_type == "HandoffOffer")
3651 .expect("offer entry")
3652 .received_at_ms;
3653 let expected_d = offer_received_at + 20;
3654
3655 let shared = rt.registry.get_shared(&sid).await.unwrap();
3656 let mut guard = shared.lock().await;
3657 let session = &mut *guard;
3658 let observed = Utc::now().timestamp_millis();
3660 assert!(observed > expected_d);
3661 rt.synthesize_due_accept(&sid, session, observed)
3662 .await
3663 .unwrap();
3664 drop(guard);
3665
3666 let entry = rt
3667 .log_store
3668 .get_log(&sid)
3669 .await
3670 .unwrap()
3671 .into_iter()
3672 .find(|e| e.message_id == reserved_id("h1"))
3673 .expect("the synthetic entry");
3674 assert_eq!(entry.timestamp_unix_ms, expected_d, "envelope clock is D");
3675 assert_eq!(entry.received_at_ms, expected_d, "entry clock is D");
3676 assert_ne!(
3677 entry.received_at_ms, observed,
3678 "received_at_ms must not be the observation time"
3679 );
3680 assert_eq!(entry.entry_kind, EntryKind::Incoming);
3681 }
3682
3683 #[tokio::test]
3693 async fn synthetic_accept_is_not_credited_as_participant_activity() {
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 shared = rt.registry.get_shared(&sid).await.unwrap();
3699 let mut guard = shared.lock().await;
3700 let session = &mut *guard;
3701 let before = session.participant_message_counts.get(TARGET).copied();
3702 rt.synthesize_due_accept(&sid, session, Utc::now().timestamp_millis())
3703 .await
3704 .unwrap();
3705 assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3706 assert_eq!(
3707 session.participant_message_counts.get(TARGET).copied(),
3708 before,
3709 "the target must not be credited with a message they did not send"
3710 );
3711 }
3712 #[tokio::test]
3739 async fn a_synthetic_entry_on_the_checkpoint_boundary_checkpoints_either_path() {
3740 async fn shape(rt: &Runtime, sid: &str) -> Vec<(EntryKind, String)> {
3741 rt.log_store
3742 .get_log(sid)
3743 .await
3744 .expect("log")
3745 .iter()
3746 .map(|e| (e.entry_kind.clone(), e.message_type.clone()))
3747 .collect()
3748 }
3749
3750 let mut eager = make_runtime();
3752 eager.checkpoint_interval = 3;
3753 let eager_sid = handoff_session_with_timed_offer(&eager, 20).await;
3754 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3755 assert_eq!(
3756 eager.sweep_due_synthetic_accepts().await,
3757 1,
3758 "sweep emitted"
3759 );
3760
3761 let mut lazy = make_runtime();
3765 lazy.checkpoint_interval = 3;
3766 let lazy_sid = handoff_session_with_timed_offer(&lazy, 20).await;
3767 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3768 lazy.process(
3769 &env(
3770 HANDOFF_MODE,
3771 "HandoffContext",
3772 "ctx-1",
3773 &lazy_sid,
3774 OWNER,
3775 handoff_context("h1"),
3776 ),
3777 None,
3778 )
3779 .await
3780 .expect("context accepted");
3781
3782 let expected = vec![
3783 (EntryKind::Incoming, "SessionStart".to_string()),
3784 (EntryKind::Incoming, "HandoffOffer".to_string()),
3785 (EntryKind::Incoming, "HandoffAccept".to_string()),
3786 (EntryKind::Checkpoint, "Checkpoint".to_string()),
3787 ];
3788 assert_eq!(shape(&eager, &eager_sid).await, expected, "eager sweep");
3789
3790 let lazy_shape = shape(&lazy, &lazy_sid).await;
3791 assert_eq!(
3792 lazy_shape[..4],
3793 expected[..],
3794 "the lazy path must checkpoint in the same place as the eager one"
3795 );
3796 assert_eq!(
3799 lazy_shape[4..],
3800 [(EntryKind::Incoming, "HandoffContext".to_string())],
3801 "no second checkpoint for the trigger's own append"
3802 );
3803 }
3804}