1use chrono::Utc;
2use std::sync::Arc;
3
4use crate::error::MacpError;
5use crate::extensions::ExtensionProviderRegistry;
6use crate::log_store::{EntryKind, LogEntry, LogStore};
7use crate::metrics::RuntimeMetrics;
8use crate::mode_registry::ModeRegistry;
9use crate::pb::{Envelope, ModeDescriptor};
10use crate::policy::registry::PolicyRegistry;
11use crate::policy::PolicyDefinition;
12use crate::registry::SessionRegistry;
13use crate::session::{
14 extract_ttl_ms, parse_session_start_payload, validate_canonical_session_start_payload_for_mode,
15 validate_session_id_for_acceptance, Session, SessionState,
16};
17use crate::storage::StorageBackend;
18use crate::stream_bus::SessionStreamBus;
19
20#[derive(Debug)]
21pub struct ProcessResult {
22 pub session_state: SessionState,
23 pub duplicate: bool,
24}
25
26#[derive(Clone, Debug)]
27pub enum SessionLifecycleEvent {
28 Created { session_id: String },
29 Resolved { session_id: String },
30 Expired { session_id: String },
31 Suspended { session_id: String },
32 Resumed { session_id: String },
33 Cancelled { session_id: String },
34}
35
36pub struct Runtime {
37 pub storage: Arc<dyn StorageBackend>,
38 pub registry: Arc<SessionRegistry>,
39 pub log_store: Arc<LogStore>,
40 stream_bus: Arc<SessionStreamBus>,
41 signal_bus: tokio::sync::broadcast::Sender<Envelope>,
42 session_lifecycle_bus: tokio::sync::broadcast::Sender<SessionLifecycleEvent>,
43 mode_registry: Arc<ModeRegistry>,
44 policy_registry: Arc<PolicyRegistry>,
45 #[allow(dead_code)] extensions: Arc<ExtensionProviderRegistry>,
47 metrics: Arc<RuntimeMetrics>,
48 checkpoint_interval: usize,
49}
50
51impl Runtime {
52 pub fn new(
53 storage: Arc<dyn StorageBackend>,
54 registry: Arc<SessionRegistry>,
55 log_store: Arc<LogStore>,
56 ) -> Self {
57 Self::with_mode_registry(
58 storage,
59 registry,
60 log_store,
61 Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
62 macp_policy::DefaultPolicyEvaluator,
63 ))),
64 )
65 }
66
67 pub fn with_mode_registry(
68 storage: Arc<dyn StorageBackend>,
69 registry: Arc<SessionRegistry>,
70 log_store: Arc<LogStore>,
71 mode_registry: Arc<ModeRegistry>,
72 ) -> Self {
73 Self::with_registries(
74 storage,
75 registry,
76 log_store,
77 mode_registry,
78 Arc::new(PolicyRegistry::new()),
79 )
80 }
81
82 pub fn with_registries(
83 storage: Arc<dyn StorageBackend>,
84 registry: Arc<SessionRegistry>,
85 log_store: Arc<LogStore>,
86 mode_registry: Arc<ModeRegistry>,
87 policy_registry: Arc<PolicyRegistry>,
88 ) -> Self {
89 let checkpoint_interval = std::env::var("MACP_CHECKPOINT_INTERVAL")
90 .ok()
91 .and_then(|v| v.parse().ok())
92 .unwrap_or(0); let (signal_tx, _) = tokio::sync::broadcast::channel(256);
94 let (session_lifecycle_tx, _) = tokio::sync::broadcast::channel(64);
95 Self {
96 storage,
97 registry,
98 log_store,
99 stream_bus: Arc::new(SessionStreamBus::default()),
100 signal_bus: signal_tx,
101 session_lifecycle_bus: session_lifecycle_tx,
102 mode_registry,
103 policy_registry,
104 extensions: Arc::new(ExtensionProviderRegistry::new()),
105 metrics: Arc::new(RuntimeMetrics::new()),
106 checkpoint_interval,
107 }
108 }
109
110 pub fn registered_mode_names(&self) -> Vec<String> {
113 self.mode_registry.all_mode_names()
114 }
115
116 pub fn standard_mode_descriptors(&self) -> Vec<ModeDescriptor> {
118 self.mode_registry.standard_mode_descriptors()
119 }
120
121 pub fn extension_mode_descriptors(&self) -> Vec<ModeDescriptor> {
123 self.mode_registry.extension_mode_descriptors()
124 }
125
126 pub fn register_extension(&self, descriptor: ModeDescriptor) -> Result<(), String> {
127 self.mode_registry.register_extension(descriptor)
128 }
129
130 pub fn unregister_extension(&self, mode: &str) -> Result<(), String> {
131 self.mode_registry.unregister_extension(mode)
132 }
133
134 pub fn promote_mode(&self, mode: &str, new_name: Option<&str>) -> Result<String, String> {
135 self.mode_registry.promote_mode(mode, new_name)
136 }
137
138 pub fn subscribe_mode_changes(&self) -> tokio::sync::broadcast::Receiver<()> {
139 self.mode_registry.subscribe_changes()
140 }
141
142 pub fn mode_registry(&self) -> &Arc<ModeRegistry> {
143 &self.mode_registry
144 }
145
146 pub fn register_policy(&self, definition: PolicyDefinition) -> Result<(), String> {
149 self.policy_registry.register(definition)
150 }
151
152 pub fn unregister_policy(&self, policy_id: &str) -> Result<(), String> {
153 self.policy_registry.unregister(policy_id)
154 }
155
156 pub fn get_policy(&self, policy_id: &str) -> Option<PolicyDefinition> {
157 self.policy_registry.get(policy_id)
158 }
159
160 pub fn list_policies(&self, mode_filter: Option<&str>) -> Vec<PolicyDefinition> {
161 self.policy_registry.list(mode_filter)
162 }
163
164 pub fn subscribe_policy_changes(&self) -> tokio::sync::broadcast::Receiver<()> {
165 self.policy_registry.subscribe_changes()
166 }
167
168 pub fn policy_registry(&self) -> &Arc<PolicyRegistry> {
169 &self.policy_registry
170 }
171
172 pub fn metrics(&self) -> &Arc<RuntimeMetrics> {
173 &self.metrics
174 }
175
176 pub fn subscribe_session_stream(
177 &self,
178 session_id: &str,
179 ) -> tokio::sync::broadcast::Receiver<Envelope> {
180 self.stream_bus.subscribe(session_id)
181 }
182
183 pub fn subscribe_signals(&self) -> tokio::sync::broadcast::Receiver<Envelope> {
184 self.signal_bus.subscribe()
185 }
186
187 pub fn subscribe_session_lifecycle(
188 &self,
189 ) -> tokio::sync::broadcast::Receiver<SessionLifecycleEvent> {
190 self.session_lifecycle_bus.subscribe()
191 }
192
193 pub async fn get_session_envelopes_after(
199 &self,
200 session_id: &str,
201 after_sequence: u64,
202 ) -> Result<Vec<Envelope>, u64> {
203 Ok(self
204 .log_store
205 .get_incoming_after(session_id, after_sequence)
206 .await?
207 .into_iter()
208 .map(|(_idx, entry)| Envelope {
209 macp_version: if entry.macp_version.is_empty() {
210 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 self.save_session_to_storage(session).await;
1198 self.metrics.record_session_expired(&session.mode);
1199 let _ = self
1200 .session_lifecycle_bus
1201 .send(SessionLifecycleEvent::Expired {
1202 session_id: session_id.to_string(),
1203 });
1204 Err(MacpError::TtlExpired)
1205 }
1206 }
1207 }
1208
1209 async fn maybe_compact_log(&self, session_id: &str, session: &Session) -> bool {
1212 let discarded = match self.log_store.get_log(session_id).await {
1216 Some(entries) => {
1217 let prior_base: u64 = entries
1218 .iter()
1219 .filter(|e| e.entry_kind == EntryKind::Checkpoint)
1220 .map(|e| e.compacted_incoming_ordinals)
1221 .max()
1222 .unwrap_or(0);
1223 prior_base
1224 + entries
1225 .iter()
1226 .filter(|e| e.entry_kind == EntryKind::Incoming)
1227 .count() as u64
1228 }
1229 None => 0,
1230 };
1231 match crate::storage::compaction::compact_session_log(
1232 &*self.storage,
1233 session_id,
1234 session,
1235 discarded,
1236 )
1237 .await
1238 {
1239 Ok(checkpoint) => {
1240 self.log_store
1244 .replace_session_log(session_id, vec![checkpoint])
1245 .await;
1246 true
1247 }
1248 Err(e) => {
1249 tracing::debug!(
1250 session_id,
1251 error = %e,
1252 "log compaction skipped (backend may not support it)"
1253 );
1254 false
1255 }
1256 }
1257 }
1258
1259 async fn force_insert_checkpoint(&self, session_id: &str, session: &Session) {
1262 let persisted = crate::registry::PersistedSession::from(session);
1263 let raw_payload = match serde_json::to_vec(&persisted) {
1264 Ok(bytes) => bytes,
1265 Err(e) => {
1266 tracing::warn!(session_id, error = %e, "failed to serialize forced checkpoint");
1267 return;
1268 }
1269 };
1270 let now = Utc::now().timestamp_millis();
1271 let checkpoint = LogEntry {
1272 message_id: String::new(),
1273 received_at_ms: now,
1274 sender: "_runtime".into(),
1275 message_type: "Checkpoint".into(),
1276 raw_payload,
1277 entry_kind: EntryKind::Checkpoint,
1278 session_id: session_id.into(),
1279 mode: session.mode.clone(),
1280 macp_version: String::new(),
1281 timestamp_unix_ms: now,
1282 bound_mode_version: None,
1283 semantics_rev: 0,
1284 bound_max_suspend_ms: None,
1285 compacted_incoming_ordinals: 0,
1286 };
1287 if let Err(e) = self.storage.append_log_entry(session_id, &checkpoint).await {
1288 tracing::warn!(session_id, error = %e, "failed to write forced checkpoint");
1289 return;
1290 }
1291 self.log_store.append(session_id, checkpoint).await;
1292 tracing::debug!(
1293 session_id,
1294 "forced checkpoint inserted for terminal session"
1295 );
1296 }
1297
1298 async fn maybe_insert_checkpoint(&self, session_id: &str, session: &Session) {
1308 if self.checkpoint_interval == 0 {
1309 return;
1310 }
1311 let log_len = self
1312 .log_store
1313 .get_log(session_id)
1314 .await
1315 .map(|l| l.len())
1316 .unwrap_or(0);
1317 if log_len < self.checkpoint_interval || log_len % self.checkpoint_interval != 0 {
1319 return;
1320 }
1321 self.force_insert_checkpoint(session_id, session).await;
1322 tracing::debug!(session_id, log_len, "checkpoint inserted at interval");
1323 }
1324
1325 pub async fn cleanup_expired_sessions(&self) {
1329 let now = Utc::now().timestamp_millis();
1330 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1335 let guard = self.registry.sessions.read().await;
1336 guard
1337 .iter()
1338 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1339 .collect()
1340 };
1341
1342 let mut expired_count = 0usize;
1343 for (session_id, shared) in candidates {
1344 let mut session = shared.lock().await;
1345 if session.state != SessionState::Open || now <= session.ttl_expiry {
1346 continue;
1347 }
1348 let entry =
1349 Self::make_internal_entry("TtlExpired", b"", &session_id, &session.mode, now);
1350 if let Err(e) = self.storage.append_log_entry(&session_id, &entry).await {
1351 tracing::warn!(
1352 session_id,
1353 error = %e,
1354 "failed to write TTL expiry during cleanup"
1355 );
1356 continue;
1357 }
1358 self.log_store.append(&session_id, entry).await;
1359 session.state = SessionState::Expired;
1360 self.metrics.record_session_expired(&session.mode);
1361 self.save_session_to_storage(&session).await;
1362 if !self.maybe_compact_log(&session_id, &session).await {
1363 self.force_insert_checkpoint(&session_id, &session).await;
1364 }
1365 expired_count += 1;
1366 tracing::info!(session_id = %session_id, "session expired via background cleanup");
1367 let _ = self
1368 .session_lifecycle_bus
1369 .send(SessionLifecycleEvent::Expired {
1370 session_id: session_id.clone(),
1371 });
1372 }
1373
1374 if expired_count > 0 {
1375 tracing::info!(count = expired_count, "background cleanup expired sessions");
1376 }
1377 }
1378
1379 pub async fn sweep_due_synthetic_accepts(&self) -> usize {
1439 let now = Utc::now().timestamp_millis();
1440 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1444 let guard = self.registry.sessions.read().await;
1445 guard
1446 .iter()
1447 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1448 .collect()
1449 };
1450
1451 let mut emitted = 0usize;
1452 for (session_id, shared) in candidates {
1453 let mut session = shared.lock().await;
1454 if session.state != SessionState::Open {
1465 continue;
1466 }
1467 match self
1468 .synthesize_due_accept(&session_id, &mut session, now)
1469 .await
1470 {
1471 Ok(true) => emitted += 1,
1472 Ok(false) => {}
1473 Err(e) => {
1479 tracing::warn!(
1480 session_id = %session_id,
1481 error = %e,
1482 "eager sweep could not emit a due synthetic envelope"
1483 );
1484 }
1485 }
1486 }
1487
1488 if emitted > 0 {
1489 tracing::info!(
1490 count = emitted,
1491 "eager sweep appended due synthetic envelopes"
1492 );
1493 }
1494 emitted
1495 }
1496
1497 pub async fn gc_disk_sessions(&self, retention_secs: u64) -> usize {
1505 let now = Utc::now().timestamp_millis();
1506 let cutoff = now - (retention_secs as i64 * 1000);
1507 let ids = match self.storage.list_session_ids().await {
1508 Ok(ids) => ids,
1509 Err(e) => {
1510 tracing::warn!(error = %e, "disk GC: cannot list sessions");
1511 return 0;
1512 }
1513 };
1514 let mut removed = 0usize;
1515 for id in ids {
1516 let eligible = if let Some(shared) = self.registry.get_shared(&id).await {
1519 let s = shared.lock().await;
1520 s.state.is_terminal() && s.started_at_unix_ms < cutoff
1521 } else {
1522 match self.storage.load_session(&id).await {
1523 Ok(Some(s)) => s.state.is_terminal() && s.started_at_unix_ms < cutoff,
1524 _ => false,
1527 }
1528 };
1529 if !eligible {
1530 continue;
1531 }
1532 match self.storage.delete_session(&id).await {
1533 Ok(()) => {
1534 {
1535 let mut guard = self.registry.sessions.write().await;
1536 guard.remove(&id);
1537 }
1538 self.log_store.remove_session_log(&id).await;
1539 let _ = self.stream_bus.remove_if_unused(&id);
1540 removed += 1;
1541 }
1542 Err(e) => {
1543 tracing::warn!(session_id = %id, error = %e, "disk GC: delete failed");
1544 }
1545 }
1546 }
1547 if removed > 0 {
1548 tracing::info!(count = removed, "disk GC removed terminal sessions");
1549 }
1550 removed
1551 }
1552
1553 pub async fn evict_stale_sessions(&self, retention_secs: u64) {
1559 let now = Utc::now().timestamp_millis();
1560 let cutoff = now - (retention_secs as i64 * 1000);
1561
1562 let candidates: Vec<(String, crate::registry::SharedSession)> = {
1563 let guard = self.registry.sessions.read().await;
1564 guard
1565 .iter()
1566 .map(|(id, arc)| (id.clone(), std::sync::Arc::clone(arc)))
1567 .collect()
1568 };
1569 let mut evict_ids = Vec::new();
1570 for (id, shared) in candidates {
1571 let session = shared.lock().await;
1572 if matches!(
1573 session.state,
1574 SessionState::Resolved | SessionState::Expired | SessionState::Cancelled
1575 ) && session.started_at_unix_ms < cutoff
1576 {
1577 evict_ids.push(id);
1578 }
1579 }
1580
1581 if evict_ids.is_empty() {
1582 return;
1583 }
1584 {
1585 let mut guard = self.registry.sessions.write().await;
1586 for id in &evict_ids {
1587 guard.remove(id);
1588 }
1589 }
1590 for id in &evict_ids {
1591 self.log_store.remove_session_log(id).await;
1592 let _ = self.stream_bus.remove_if_unused(id);
1595 }
1596 tracing::info!(
1597 count = evict_ids.len(),
1598 "evicted stale sessions from memory (registry + log cache + stream bus)"
1599 );
1600 }
1601}
1602
1603#[cfg(test)]
1604mod tests {
1605 use super::*;
1606 use crate::decision_pb::ProposalPayload;
1607 use crate::pb::{CommitmentPayload, SessionStartPayload};
1608 use prost::Message;
1609
1610 fn new_sid() -> String {
1611 uuid::Uuid::new_v4().as_hyphenated().to_string()
1612 }
1613
1614 fn make_runtime() -> Runtime {
1615 let storage: Arc<dyn StorageBackend> = Arc::new(crate::storage::MemoryBackend);
1616 let registry = Arc::new(SessionRegistry::new());
1617 let log_store = Arc::new(LogStore::new());
1618 Runtime::new(storage, registry, log_store)
1619 }
1620
1621 fn session_start(participants: Vec<String>) -> Vec<u8> {
1622 SessionStartPayload {
1623 intent: "intent".into(),
1624 participants,
1625 mode_version: "1.0.0".into(),
1626 configuration_version: "cfg-1".into(),
1627 policy_version: String::new(),
1628 ttl_ms: 1_000,
1629 context_id: String::new(),
1630 extensions: std::collections::HashMap::new(),
1631 roots: vec![],
1632 max_suspend_ms: 0,
1633 }
1634 .encode_to_vec()
1635 }
1636
1637 fn env(
1638 mode: &str,
1639 message_type: &str,
1640 message_id: &str,
1641 session_id: &str,
1642 sender: &str,
1643 payload: Vec<u8>,
1644 ) -> Envelope {
1645 Envelope {
1646 macp_version: "1.0".into(),
1647 mode: mode.into(),
1648 message_type: message_type.into(),
1649 message_id: message_id.into(),
1650 session_id: session_id.into(),
1651 sender: sender.into(),
1652 timestamp_unix_ms: Utc::now().timestamp_millis(),
1653 payload,
1654 }
1655 }
1656
1657 #[tokio::test]
1658 async fn standard_session_start_is_strict() {
1659 let rt = make_runtime();
1660 let sid = new_sid();
1661 let bad = SessionStartPayload {
1662 ttl_ms: 0,
1663 ..Default::default()
1664 }
1665 .encode_to_vec();
1666 let err = rt
1667 .process(
1668 &env(
1669 "macp.mode.decision.v1",
1670 "SessionStart",
1671 "m1",
1672 &sid,
1673 "agent://orchestrator",
1674 bad,
1675 ),
1676 None,
1677 )
1678 .await
1679 .unwrap_err();
1680 assert!(matches!(
1681 err,
1682 MacpError::InvalidPayload | MacpError::InvalidTtl
1683 ));
1684 }
1685
1686 #[tokio::test]
1699 async fn a_promoted_mode_still_gets_canonical_session_start_validation() {
1700 let mode_registry = Arc::new(ModeRegistry::build_default(std::sync::Arc::new(
1701 macp_policy::DefaultPolicyEvaluator,
1702 )));
1703 mode_registry
1704 .register_extension(crate::pb::ModeDescriptor {
1705 mode: "ext.promoted.v1".into(),
1706 mode_version: "1.0.0".into(),
1707 title: "Promoted".into(),
1708 description: "promotion target".into(),
1709 determinism_class: "semantic-deterministic".into(),
1710 participant_model: "declared".into(),
1711 message_types: vec!["SessionStart".into(), "Commitment".into()],
1712 terminal_message_types: vec!["Commitment".into()],
1713 ..Default::default()
1714 })
1715 .expect("register extension");
1716 assert_eq!(
1717 mode_registry.promote_mode("ext.promoted.v1", None).unwrap(),
1718 "ext.promoted.v1"
1719 );
1720 assert!(
1721 mode_registry.requires_strict_session_start("ext.promoted.v1"),
1722 "promotion must mark the entry strict"
1723 );
1724 assert!(
1725 !crate::session::requires_strict_session_start("ext.promoted.v1"),
1726 "the core's static list must NOT know this name — that disagreement is the point"
1727 );
1728
1729 let rt = Runtime::with_mode_registry(
1730 Arc::new(crate::storage::MemoryBackend),
1731 Arc::new(SessionRegistry::new()),
1732 Arc::new(LogStore::new()),
1733 mode_registry,
1734 );
1735
1736 let err = rt
1738 .process(
1739 &env(
1740 "ext.promoted.v1",
1741 "SessionStart",
1742 "m1",
1743 &new_sid(),
1744 "agent://orchestrator",
1745 session_start(vec![]),
1746 ),
1747 None,
1748 )
1749 .await
1750 .unwrap_err();
1751 assert_eq!(err.to_string(), "InvalidPayload");
1752
1753 let no_versions = SessionStartPayload {
1756 participants: vec!["agent://fraud".into()],
1757 ttl_ms: 1_000,
1758 ..Default::default()
1759 }
1760 .encode_to_vec();
1761 let err = rt
1762 .process(
1763 &env(
1764 "ext.promoted.v1",
1765 "SessionStart",
1766 "m2",
1767 &new_sid(),
1768 "agent://orchestrator",
1769 no_versions,
1770 ),
1771 None,
1772 )
1773 .await
1774 .unwrap_err();
1775 assert_eq!(err.to_string(), "InvalidPayload");
1776
1777 rt.process(
1780 &env(
1781 "ext.promoted.v1",
1782 "SessionStart",
1783 "m3",
1784 &new_sid(),
1785 "agent://orchestrator",
1786 session_start(vec!["agent://fraud".into()]),
1787 ),
1788 None,
1789 )
1790 .await
1791 .expect("a complete SessionStart must still be accepted for a promoted mode");
1792 }
1793
1794 #[tokio::test]
1795 async fn empty_mode_is_rejected() {
1796 let rt = make_runtime();
1797 let sid = new_sid();
1798 let err = rt
1799 .process(
1800 &env(
1801 "",
1802 "SessionStart",
1803 "m1",
1804 &sid,
1805 "agent://orchestrator",
1806 session_start(vec!["agent://fraud".into()]),
1807 ),
1808 None,
1809 )
1810 .await
1811 .unwrap_err();
1812 assert_eq!(err.to_string(), "InvalidEnvelope");
1813 }
1814
1815 #[tokio::test]
1816 async fn rejected_messages_do_not_enter_dedup_state() {
1817 let rt = make_runtime();
1818 let sid = new_sid();
1819 rt.process(
1820 &env(
1821 "macp.mode.decision.v1",
1822 "SessionStart",
1823 "m1",
1824 &sid,
1825 "agent://orchestrator",
1826 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
1827 ),
1828 None,
1829 )
1830 .await
1831 .unwrap();
1832
1833 let bad = rt
1834 .process(
1835 &env(
1836 "macp.mode.decision.v1",
1837 "Proposal",
1838 "m2",
1839 &sid,
1840 "agent://fraud",
1841 b"not-protobuf".to_vec(),
1842 ),
1843 None,
1844 )
1845 .await
1846 .unwrap_err();
1847 assert_eq!(bad.to_string(), "InvalidPayload");
1848
1849 let good = ProposalPayload {
1850 proposal_id: "p1".into(),
1851 option: "step-up".into(),
1852 rationale: "risk".into(),
1853 supporting_data: vec![],
1854 }
1855 .encode_to_vec();
1856 let result = rt
1857 .process(
1858 &env(
1859 "macp.mode.decision.v1",
1860 "Proposal",
1861 "m2",
1862 &sid,
1863 "agent://orchestrator",
1864 good,
1865 ),
1866 None,
1867 )
1868 .await
1869 .unwrap();
1870 assert!(!result.duplicate);
1871 }
1872
1873 #[tokio::test]
1874 async fn get_session_transitions_expired_sessions() {
1875 let rt = make_runtime();
1876 let sid = new_sid();
1877 let payload = SessionStartPayload {
1878 intent: "intent".into(),
1879 participants: vec!["agent://fraud".into()],
1880 mode_version: "1.0.0".into(),
1881 configuration_version: "cfg-1".into(),
1882 policy_version: String::new(),
1883 ttl_ms: 1,
1884 context_id: String::new(),
1885 extensions: std::collections::HashMap::new(),
1886 roots: vec![],
1887 max_suspend_ms: 0,
1888 }
1889 .encode_to_vec();
1890 rt.process(
1891 &env(
1892 "macp.mode.decision.v1",
1893 "SessionStart",
1894 "m1",
1895 &sid,
1896 "agent://orchestrator",
1897 payload,
1898 ),
1899 None,
1900 )
1901 .await
1902 .unwrap();
1903 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
1904 let session = rt.get_session_checked(&sid).await.unwrap();
1905 assert_eq!(session.state, SessionState::Expired);
1906 }
1907
1908 #[tokio::test]
1909 async fn multi_round_requires_standard_session_start() {
1910 let rt = make_runtime();
1911 let sid = new_sid();
1912 let payload = SessionStartPayload {
1914 participants: vec!["creator".into(), "other".into()],
1915 ..Default::default()
1916 }
1917 .encode_to_vec();
1918 let err = rt
1919 .process(
1920 &env(
1921 "ext.multi_round.v1",
1922 "SessionStart",
1923 "m1",
1924 &sid,
1925 "creator",
1926 payload,
1927 ),
1928 None,
1929 )
1930 .await
1931 .unwrap_err();
1932 assert!(matches!(
1933 err,
1934 MacpError::InvalidPayload | MacpError::InvalidTtl
1935 ));
1936 }
1937
1938 #[tokio::test]
1939 async fn multi_round_valid_session_start() {
1940 let rt = make_runtime();
1941 let sid = new_sid();
1942 let payload = session_start(vec!["alice".into(), "bob".into()]);
1943 rt.process(
1944 &env(
1945 "ext.multi_round.v1",
1946 "SessionStart",
1947 "m1",
1948 &sid,
1949 "coordinator",
1950 payload,
1951 ),
1952 None,
1953 )
1954 .await
1955 .unwrap();
1956 let session = rt.get_session_checked(&sid).await.unwrap();
1957 assert_eq!(session.mode, "ext.multi_round.v1");
1958 assert_eq!(session.participants, vec!["alice", "bob"]);
1959 }
1960
1961 #[tokio::test]
1962 async fn duplicate_session_start_message_id_returns_duplicate() {
1963 let rt = make_runtime();
1964 let sid = new_sid();
1965 let payload = session_start(vec!["agent://fraud".into()]);
1966 rt.process(
1967 &env(
1968 "macp.mode.decision.v1",
1969 "SessionStart",
1970 "m1",
1971 &sid,
1972 "agent://orchestrator",
1973 payload.clone(),
1974 ),
1975 None,
1976 )
1977 .await
1978 .unwrap();
1979
1980 let result = rt
1981 .process(
1982 &env(
1983 "macp.mode.decision.v1",
1984 "SessionStart",
1985 "m1",
1986 &sid,
1987 "agent://orchestrator",
1988 payload,
1989 ),
1990 None,
1991 )
1992 .await
1993 .unwrap();
1994 assert!(result.duplicate);
1995 }
1996
1997 #[tokio::test]
1998 async fn non_start_mode_mismatch_rejected() {
1999 let rt = make_runtime();
2000 let sid = new_sid();
2001 rt.process(
2002 &env(
2003 "macp.mode.decision.v1",
2004 "SessionStart",
2005 "m1",
2006 &sid,
2007 "agent://orchestrator",
2008 session_start(vec!["agent://fraud".into()]),
2009 ),
2010 None,
2011 )
2012 .await
2013 .unwrap();
2014
2015 let proposal = ProposalPayload {
2016 proposal_id: "p1".into(),
2017 option: "step-up".into(),
2018 rationale: "risk".into(),
2019 supporting_data: vec![],
2020 }
2021 .encode_to_vec();
2022 let err = rt
2023 .process(
2024 &env(
2025 "macp.mode.task.v1",
2026 "Proposal",
2027 "m2",
2028 &sid,
2029 "agent://orchestrator",
2030 proposal,
2031 ),
2032 None,
2033 )
2034 .await
2035 .unwrap_err();
2036 assert_eq!(err.to_string(), "InvalidEnvelope");
2037 }
2038
2039 #[tokio::test]
2040 async fn cancel_idempotent_on_already_expired() {
2041 let rt = make_runtime();
2042 let sid = new_sid();
2043 let payload = SessionStartPayload {
2044 intent: "intent".into(),
2045 participants: vec!["agent://fraud".into()],
2046 mode_version: "1.0.0".into(),
2047 configuration_version: "cfg-1".into(),
2048 policy_version: String::new(),
2049 ttl_ms: 1,
2050 context_id: String::new(),
2051 extensions: std::collections::HashMap::new(),
2052 roots: vec![],
2053 max_suspend_ms: 0,
2054 }
2055 .encode_to_vec();
2056 rt.process(
2057 &env(
2058 "macp.mode.decision.v1",
2059 "SessionStart",
2060 "m1",
2061 &sid,
2062 "agent://orchestrator",
2063 payload,
2064 ),
2065 None,
2066 )
2067 .await
2068 .unwrap();
2069 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2070 let result = rt
2071 .cancel_session(&sid, "cleanup", "agent://orchestrator")
2072 .await
2073 .unwrap();
2074 assert_eq!(result.session_state, SessionState::Expired);
2075 }
2076
2077 #[tokio::test]
2078 async fn accepted_envelopes_are_published_in_order() {
2079 let rt = make_runtime();
2080 let sid = new_sid();
2081 let mut events = rt.subscribe_session_stream(&sid);
2082
2083 let start = env(
2084 "macp.mode.decision.v1",
2085 "SessionStart",
2086 "m1",
2087 &sid,
2088 "agent://orchestrator",
2089 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2090 );
2091 rt.process(&start, None).await.unwrap();
2092 let first = events.recv().await.unwrap();
2093 assert_eq!(first.message_id, "m1");
2094 assert_eq!(first.message_type, "SessionStart");
2095
2096 let proposal = ProposalPayload {
2097 proposal_id: "p1".into(),
2098 option: "step-up".into(),
2099 rationale: "risk".into(),
2100 supporting_data: vec![],
2101 }
2102 .encode_to_vec();
2103 let proposal_env = env(
2104 "macp.mode.decision.v1",
2105 "Proposal",
2106 "m2",
2107 &sid,
2108 "agent://orchestrator",
2109 proposal,
2110 );
2111 rt.process(&proposal_env, None).await.unwrap();
2112 let second = events.recv().await.unwrap();
2113 assert_eq!(second.message_id, "m2");
2114 assert_eq!(second.message_type, "Proposal");
2115 }
2116
2117 #[tokio::test]
2118 async fn commitment_versions_are_carried_into_resolution() {
2119 let rt = make_runtime();
2120 let sid = new_sid();
2121 rt.process(
2122 &env(
2123 "macp.mode.proposal.v1",
2124 "SessionStart",
2125 "m1",
2126 &sid,
2127 "agent://buyer",
2128 session_start(vec!["agent://buyer".into(), "agent://seller".into()]),
2129 ),
2130 None,
2131 )
2132 .await
2133 .unwrap();
2134
2135 let proposal = crate::proposal_pb::ProposalPayload {
2136 proposal_id: "p1".into(),
2137 title: "offer".into(),
2138 summary: "summary".into(),
2139 details: vec![],
2140 tags: vec![],
2141 }
2142 .encode_to_vec();
2143 rt.process(
2144 &env(
2145 "macp.mode.proposal.v1",
2146 "Proposal",
2147 "m2",
2148 &sid,
2149 "agent://seller",
2150 proposal,
2151 ),
2152 None,
2153 )
2154 .await
2155 .unwrap();
2156 let accept = crate::proposal_pb::AcceptPayload {
2157 proposal_id: "p1".into(),
2158 reason: String::new(),
2159 }
2160 .encode_to_vec();
2161 rt.process(
2162 &env(
2163 "macp.mode.proposal.v1",
2164 "Accept",
2165 "m3",
2166 &sid,
2167 "agent://seller",
2168 accept.clone(),
2169 ),
2170 None,
2171 )
2172 .await
2173 .unwrap();
2174 rt.process(
2175 &env(
2176 "macp.mode.proposal.v1",
2177 "Accept",
2178 "m4",
2179 &sid,
2180 "agent://buyer",
2181 accept,
2182 ),
2183 None,
2184 )
2185 .await
2186 .unwrap();
2187 let commitment = CommitmentPayload {
2188 commitment_id: "c1".into(),
2189 action: "proposal.accepted".into(),
2190 authority_scope: "commercial".into(),
2191 reason: "bound".into(),
2192 mode_version: "1.0.0".into(),
2193 policy_version: "policy.default".into(),
2194 configuration_version: "cfg-1".into(),
2195 outcome_positive: true,
2196 supersedes: None,
2197 }
2198 .encode_to_vec();
2199 let result = rt
2200 .process(
2201 &env(
2202 "macp.mode.proposal.v1",
2203 "Commitment",
2204 "m5",
2205 &sid,
2206 "agent://buyer",
2207 commitment,
2208 ),
2209 None,
2210 )
2211 .await
2212 .unwrap();
2213 assert_eq!(result.session_state, SessionState::Resolved);
2214 }
2215
2216 #[tokio::test]
2217 async fn max_open_sessions_enforced_under_write_lock() {
2218 let rt = make_runtime();
2219 let sid1 = new_sid();
2220 let sid2 = new_sid();
2221 let sid3 = new_sid();
2222 rt.process(
2223 &env(
2224 "macp.mode.decision.v1",
2225 "SessionStart",
2226 "m1",
2227 &sid1,
2228 "agent://orchestrator",
2229 session_start(vec!["agent://fraud".into()]),
2230 ),
2231 Some(1),
2232 )
2233 .await
2234 .unwrap();
2235
2236 let err = rt
2237 .process(
2238 &env(
2239 "macp.mode.decision.v1",
2240 "SessionStart",
2241 "m2",
2242 &sid2,
2243 "agent://orchestrator",
2244 session_start(vec!["agent://fraud".into()]),
2245 ),
2246 Some(1),
2247 )
2248 .await
2249 .unwrap_err();
2250 assert!(matches!(err, MacpError::RateLimited));
2251
2252 rt.process(
2253 &env(
2254 "macp.mode.decision.v1",
2255 "SessionStart",
2256 "m3",
2257 &sid3,
2258 "agent://other",
2259 session_start(vec!["agent://fraud".into()]),
2260 ),
2261 Some(1),
2262 )
2263 .await
2264 .unwrap();
2265 }
2266
2267 #[tokio::test]
2268 async fn weak_session_id_rejected() {
2269 let rt = make_runtime();
2270 let err = rt
2271 .process(
2272 &env(
2273 "macp.mode.decision.v1",
2274 "SessionStart",
2275 "m1",
2276 "s1",
2277 "agent://orchestrator",
2278 session_start(vec!["agent://fraud".into()]),
2279 ),
2280 None,
2281 )
2282 .await
2283 .unwrap_err();
2284 assert_eq!(err.to_string(), "InvalidSessionId");
2285 }
2286
2287 #[tokio::test]
2288 async fn log_append_failure_rejects_session_start() {
2289 use std::io;
2290 struct FailingBackend;
2291 #[async_trait::async_trait]
2292 impl StorageBackend for FailingBackend {
2293 async fn save_session(&self, _: &Session) -> io::Result<()> {
2294 Ok(())
2295 }
2296 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2297 Ok(None)
2298 }
2299 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2300 Ok(vec![])
2301 }
2302 async fn delete_session(&self, _: &str) -> io::Result<()> {
2303 Ok(())
2304 }
2305 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2306 Ok(vec![])
2307 }
2308 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2309 Err(io::Error::other("disk full"))
2310 }
2311 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2312 Ok(vec![])
2313 }
2314 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2315 Ok(())
2316 }
2317 }
2318
2319 let storage: Arc<dyn StorageBackend> = Arc::new(FailingBackend);
2320 let registry = Arc::new(SessionRegistry::new());
2321 let log_store = Arc::new(LogStore::new());
2322 let rt = Runtime::new(storage, registry, log_store);
2323 let sid = new_sid();
2324
2325 let err = rt
2326 .process(
2327 &env(
2328 "macp.mode.decision.v1",
2329 "SessionStart",
2330 "m1",
2331 &sid,
2332 "agent://orchestrator",
2333 session_start(vec!["agent://fraud".into()]),
2334 ),
2335 None,
2336 )
2337 .await
2338 .unwrap_err();
2339 assert_eq!(err.to_string(), "StorageFailed");
2340 }
2341
2342 #[tokio::test]
2343 async fn log_append_failure_rejects_in_session_message() {
2344 use std::io;
2345 use std::sync::atomic::{AtomicUsize, Ordering};
2346
2347 struct FailOnSecondAppend {
2348 count: AtomicUsize,
2349 }
2350 #[async_trait::async_trait]
2351 impl StorageBackend for FailOnSecondAppend {
2352 async fn save_session(&self, _: &Session) -> io::Result<()> {
2353 Ok(())
2354 }
2355 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2356 Ok(None)
2357 }
2358 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2359 Ok(vec![])
2360 }
2361 async fn delete_session(&self, _: &str) -> io::Result<()> {
2362 Ok(())
2363 }
2364 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2365 Ok(vec![])
2366 }
2367 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2368 let n = self.count.fetch_add(1, Ordering::SeqCst);
2369 if n >= 1 {
2370 Err(io::Error::other("disk full"))
2371 } else {
2372 Ok(())
2373 }
2374 }
2375 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2376 Ok(vec![])
2377 }
2378 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2379 Ok(())
2380 }
2381 }
2382
2383 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2384 count: AtomicUsize::new(0),
2385 });
2386 let registry = Arc::new(SessionRegistry::new());
2387 let log_store = Arc::new(LogStore::new());
2388 let rt = Runtime::new(storage, registry, log_store);
2389 let sid = new_sid();
2390
2391 rt.process(
2393 &env(
2394 "macp.mode.decision.v1",
2395 "SessionStart",
2396 "m1",
2397 &sid,
2398 "agent://orchestrator",
2399 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2400 ),
2401 None,
2402 )
2403 .await
2404 .unwrap();
2405
2406 let proposal = ProposalPayload {
2408 proposal_id: "p1".into(),
2409 option: "step-up".into(),
2410 rationale: "risk".into(),
2411 supporting_data: vec![],
2412 }
2413 .encode_to_vec();
2414 let err = rt
2415 .process(
2416 &env(
2417 "macp.mode.decision.v1",
2418 "Proposal",
2419 "m2",
2420 &sid,
2421 "agent://orchestrator",
2422 proposal,
2423 ),
2424 None,
2425 )
2426 .await
2427 .unwrap_err();
2428 assert_eq!(err.to_string(), "StorageFailed");
2429
2430 let session = rt.get_session_checked(&sid).await.unwrap();
2432 assert!(!session.seen_message_ids.contains("m2"));
2433 }
2434
2435 #[tokio::test]
2436 async fn cancel_session_fails_if_log_append_fails() {
2437 use std::io;
2438 use std::sync::atomic::{AtomicUsize, Ordering};
2439
2440 struct FailOnSecondAppend {
2441 count: AtomicUsize,
2442 }
2443 #[async_trait::async_trait]
2444 impl StorageBackend for FailOnSecondAppend {
2445 async fn save_session(&self, _: &Session) -> io::Result<()> {
2446 Ok(())
2447 }
2448 async fn load_session(&self, _: &str) -> io::Result<Option<Session>> {
2449 Ok(None)
2450 }
2451 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
2452 Ok(vec![])
2453 }
2454 async fn delete_session(&self, _: &str) -> io::Result<()> {
2455 Ok(())
2456 }
2457 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
2458 Ok(vec![])
2459 }
2460 async fn append_log_entry(&self, _: &str, _: &LogEntry) -> io::Result<()> {
2461 let n = self.count.fetch_add(1, Ordering::SeqCst);
2462 if n >= 1 {
2463 Err(io::Error::other("disk full"))
2464 } else {
2465 Ok(())
2466 }
2467 }
2468 async fn load_log(&self, _: &str) -> io::Result<Vec<LogEntry>> {
2469 Ok(vec![])
2470 }
2471 async fn create_session_storage(&self, _: &str) -> io::Result<()> {
2472 Ok(())
2473 }
2474 }
2475
2476 let storage: Arc<dyn StorageBackend> = Arc::new(FailOnSecondAppend {
2477 count: AtomicUsize::new(0),
2478 });
2479 let registry = Arc::new(SessionRegistry::new());
2480 let log_store = Arc::new(LogStore::new());
2481 let rt = Runtime::new(storage, registry, log_store);
2482 let sid = new_sid();
2483
2484 rt.process(
2485 &env(
2486 "macp.mode.decision.v1",
2487 "SessionStart",
2488 "m1",
2489 &sid,
2490 "agent://orchestrator",
2491 session_start(vec!["agent://fraud".into()]),
2492 ),
2493 None,
2494 )
2495 .await
2496 .unwrap();
2497
2498 let err = rt
2499 .cancel_session(&sid, "test cancel", "agent://orchestrator")
2500 .await
2501 .unwrap_err();
2502 assert_eq!(err.to_string(), "StorageFailed");
2503 }
2504
2505 #[tokio::test]
2506 async fn ttl_expiration_rejects_message() {
2507 let rt = make_runtime();
2508 let sid = new_sid();
2509 let payload = SessionStartPayload {
2510 intent: "intent".into(),
2511 participants: vec!["agent://orchestrator".into(), "agent://fraud".into()],
2512 mode_version: "1.0.0".into(),
2513 configuration_version: "cfg-1".into(),
2514 policy_version: String::new(),
2515 ttl_ms: 1,
2516 context_id: String::new(),
2517 extensions: std::collections::HashMap::new(),
2518 roots: vec![],
2519 max_suspend_ms: 0,
2520 }
2521 .encode_to_vec();
2522 rt.process(
2523 &env(
2524 "macp.mode.decision.v1",
2525 "SessionStart",
2526 "m1",
2527 &sid,
2528 "agent://orchestrator",
2529 payload,
2530 ),
2531 None,
2532 )
2533 .await
2534 .unwrap();
2535 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2536 let proposal = ProposalPayload {
2537 proposal_id: "p1".into(),
2538 option: "step-up".into(),
2539 rationale: "risk".into(),
2540 supporting_data: vec![],
2541 }
2542 .encode_to_vec();
2543 let err = rt
2544 .process(
2545 &env(
2546 "macp.mode.decision.v1",
2547 "Proposal",
2548 "m2",
2549 &sid,
2550 "agent://orchestrator",
2551 proposal,
2552 ),
2553 None,
2554 )
2555 .await
2556 .unwrap_err();
2557 assert_eq!(err.to_string(), "TtlExpired");
2558 }
2559
2560 #[tokio::test]
2561 async fn cleanup_expired_sessions_marks_expired() {
2562 let rt = make_runtime();
2563 let sid = new_sid();
2564 let payload = SessionStartPayload {
2565 intent: "intent".into(),
2566 participants: vec!["agent://fraud".into()],
2567 mode_version: "1.0.0".into(),
2568 configuration_version: "cfg-1".into(),
2569 policy_version: String::new(),
2570 ttl_ms: 1,
2571 context_id: String::new(),
2572 extensions: std::collections::HashMap::new(),
2573 roots: vec![],
2574 max_suspend_ms: 0,
2575 }
2576 .encode_to_vec();
2577 rt.process(
2578 &env(
2579 "macp.mode.decision.v1",
2580 "SessionStart",
2581 "m1",
2582 &sid,
2583 "agent://orchestrator",
2584 payload,
2585 ),
2586 None,
2587 )
2588 .await
2589 .unwrap();
2590 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2591 rt.cleanup_expired_sessions().await;
2592 let session = rt.get_session_checked(&sid).await.unwrap();
2593 assert_eq!(session.state, SessionState::Expired);
2594 }
2595
2596 #[tokio::test]
2597 async fn evict_stale_sessions_removes_resolved() {
2598 let rt = make_runtime();
2599 let sid = new_sid();
2600 rt.process(
2602 &env(
2603 "macp.mode.decision.v1",
2604 "SessionStart",
2605 "m1",
2606 &sid,
2607 "agent://orchestrator",
2608 session_start(vec!["agent://orchestrator".into(), "agent://fraud".into()]),
2609 ),
2610 None,
2611 )
2612 .await
2613 .unwrap();
2614 let proposal = ProposalPayload {
2616 proposal_id: "p1".into(),
2617 option: "step-up".into(),
2618 rationale: "risk".into(),
2619 supporting_data: vec![],
2620 }
2621 .encode_to_vec();
2622 rt.process(
2623 &env(
2624 "macp.mode.decision.v1",
2625 "Proposal",
2626 "m2",
2627 &sid,
2628 "agent://orchestrator",
2629 proposal,
2630 ),
2631 None,
2632 )
2633 .await
2634 .unwrap();
2635 let commitment = CommitmentPayload {
2637 commitment_id: "c1".into(),
2638 action: "decision.selected".into(),
2639 authority_scope: "payments".into(),
2640 reason: "bound".into(),
2641 mode_version: "1.0.0".into(),
2642 policy_version: "policy.default".into(),
2643 configuration_version: "cfg-1".into(),
2644 outcome_positive: true,
2645 supersedes: None,
2646 }
2647 .encode_to_vec();
2648 let result = rt
2649 .process(
2650 &env(
2651 "macp.mode.decision.v1",
2652 "Commitment",
2653 "m3",
2654 &sid,
2655 "agent://orchestrator",
2656 commitment,
2657 ),
2658 None,
2659 )
2660 .await
2661 .unwrap();
2662 assert_eq!(result.session_state, SessionState::Resolved);
2663 tokio::time::sleep(std::time::Duration::from_millis(5)).await;
2665 rt.evict_stale_sessions(0).await;
2667 assert!(rt.registry.get_session(&sid).await.is_none());
2669 }
2670
2671 #[tokio::test]
2672 async fn session_start_with_wrong_mode_version_rejected() {
2673 let rt = make_runtime();
2674 let sid = new_sid();
2675 let payload = SessionStartPayload {
2676 intent: "test".into(),
2677 participants: vec!["agent://orchestrator".into(), "agent://worker".into()],
2678 mode_version: "99.0.0".into(), configuration_version: "cfg-1".into(),
2680 policy_version: String::new(),
2681 ttl_ms: 60_000,
2682 context_id: String::new(),
2683 extensions: std::collections::HashMap::new(),
2684 roots: vec![],
2685 max_suspend_ms: 0,
2686 }
2687 .encode_to_vec();
2688
2689 let err = rt
2690 .process(
2691 &env(
2692 "macp.mode.decision.v1",
2693 "SessionStart",
2694 "m1",
2695 &sid,
2696 "agent://orchestrator",
2697 payload,
2698 ),
2699 None,
2700 )
2701 .await
2702 .unwrap_err();
2703 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2704 }
2705
2706 #[tokio::test]
2707 async fn signal_empty_signal_type_rejected() {
2708 let rt = make_runtime();
2709 let signal_payload = crate::pb::SignalPayload {
2711 signal_type: String::new(),
2712 data: b"some data".to_vec(),
2713 confidence: 0.0,
2714 correlation_session_id: String::new(),
2715 }
2716 .encode_to_vec();
2717 let signal = Envelope {
2718 macp_version: "1.0".into(),
2719 mode: String::new(),
2720 message_type: "Signal".into(),
2721 message_id: "sig-1".into(),
2722 session_id: String::new(),
2723 sender: "agent://a".into(),
2724 timestamp_unix_ms: 0,
2725 payload: signal_payload,
2726 };
2727 let err = rt.process_signal(&signal).await.unwrap_err();
2728 assert_eq!(err.error_code(), "INVALID_ENVELOPE");
2729 }
2730
2731 #[tokio::test]
2732 async fn signal_valid_payload_accepted() {
2733 let rt = make_runtime();
2734 let signal_payload = crate::pb::SignalPayload {
2735 signal_type: "heartbeat".into(),
2736 data: vec![],
2737 confidence: 0.8,
2738 correlation_session_id: String::new(),
2739 }
2740 .encode_to_vec();
2741 let signal = Envelope {
2742 macp_version: "1.0".into(),
2743 mode: String::new(),
2744 message_type: "Signal".into(),
2745 message_id: "sig-2".into(),
2746 session_id: String::new(),
2747 sender: "agent://a".into(),
2748 timestamp_unix_ms: 0,
2749 payload: signal_payload,
2750 };
2751 rt.process_signal(&signal).await.unwrap();
2752 }
2753
2754 #[tokio::test]
2755 async fn signal_empty_payload_accepted() {
2756 let rt = make_runtime();
2757 let signal = Envelope {
2758 macp_version: "1.0".into(),
2759 mode: String::new(),
2760 message_type: "Signal".into(),
2761 message_id: "sig-3".into(),
2762 session_id: String::new(),
2763 sender: "agent://a".into(),
2764 timestamp_unix_ms: 0,
2765 payload: vec![],
2766 };
2767 rt.process_signal(&signal).await.unwrap();
2768 }
2769
2770 #[tokio::test]
2776 async fn ext_mode_empty_version_binds_descriptor_version() {
2777 let rt = make_runtime();
2778 rt.register_extension(ModeDescriptor {
2779 mode: "ext.dyn.v1".into(),
2780 mode_version: "2.5.0".into(),
2781 message_types: vec!["SessionStart".into(), "Note".into(), "Commitment".into()],
2782 terminal_message_types: vec!["Commitment".into()],
2783 ..Default::default()
2784 })
2785 .unwrap();
2786
2787 let sid = new_sid();
2788 let payload = SessionStartPayload {
2789 participants: vec!["alice".into()],
2790 configuration_version: "cfg-1".into(),
2791 ttl_ms: 60_000,
2792 ..Default::default()
2793 }
2794 .encode_to_vec();
2795 rt.process(
2796 &env("ext.dyn.v1", "SessionStart", "m1", &sid, "alice", payload),
2797 None,
2798 )
2799 .await
2800 .unwrap();
2801
2802 let session = rt.get_session_checked(&sid).await.unwrap();
2804 assert_eq!(session.mode_version, "2.5.0");
2805
2806 let bad = CommitmentPayload {
2808 commitment_id: "c1".into(),
2809 action: "work.completed".into(),
2810 authority_scope: "test".into(),
2811 reason: "done".into(),
2812 mode_version: String::new(),
2813 policy_version: "policy.default".into(),
2814 configuration_version: "cfg-1".into(),
2815 outcome_positive: true,
2816 supersedes: None,
2817 }
2818 .encode_to_vec();
2819 let err = rt
2820 .process(
2821 &env("ext.dyn.v1", "Commitment", "m2", &sid, "alice", bad),
2822 None,
2823 )
2824 .await
2825 .unwrap_err();
2826 assert_eq!(err.to_string(), "InvalidPayload");
2827
2828 let good = CommitmentPayload {
2830 commitment_id: "c1".into(),
2831 action: "work.completed".into(),
2832 authority_scope: "test".into(),
2833 reason: "done".into(),
2834 mode_version: "2.5.0".into(),
2835 policy_version: "policy.default".into(),
2836 configuration_version: "cfg-1".into(),
2837 outcome_positive: true,
2838 supersedes: None,
2839 }
2840 .encode_to_vec();
2841 let result = rt
2842 .process(
2843 &env("ext.dyn.v1", "Commitment", "m3", &sid, "alice", good),
2844 None,
2845 )
2846 .await
2847 .unwrap();
2848 assert_eq!(result.session_state, SessionState::Resolved);
2849 }
2850
2851 #[tokio::test]
2854 async fn ext_mode_binding_recorded_on_session_start_log_entry() {
2855 let rt = make_runtime();
2856 rt.register_extension(ModeDescriptor {
2857 mode: "ext.dyn2.v1".into(),
2858 mode_version: "3.0.0".into(),
2859 message_types: vec!["SessionStart".into(), "Commitment".into()],
2860 terminal_message_types: vec!["Commitment".into()],
2861 ..Default::default()
2862 })
2863 .unwrap();
2864
2865 let sid = new_sid();
2866 let payload = SessionStartPayload {
2867 participants: vec!["alice".into()],
2868 configuration_version: "cfg-1".into(),
2869 ttl_ms: 60_000,
2870 ..Default::default()
2871 }
2872 .encode_to_vec();
2873 rt.process(
2874 &env("ext.dyn2.v1", "SessionStart", "m1", &sid, "alice", payload),
2875 None,
2876 )
2877 .await
2878 .unwrap();
2879
2880 let log = rt.log_store.get_log(&sid).await.unwrap();
2881 assert_eq!(log[0].message_type, "SessionStart");
2882 assert_eq!(log[0].bound_mode_version.as_deref(), Some("3.0.0"));
2883
2884 let sid2 = new_sid();
2886 let payload2 = SessionStartPayload {
2887 participants: vec!["alice".into()],
2888 mode_version: "3.0.0".into(),
2889 configuration_version: "cfg-1".into(),
2890 ttl_ms: 60_000,
2891 ..Default::default()
2892 }
2893 .encode_to_vec();
2894 rt.process(
2895 &env(
2896 "ext.dyn2.v1",
2897 "SessionStart",
2898 "m1",
2899 &sid2,
2900 "alice",
2901 payload2,
2902 ),
2903 None,
2904 )
2905 .await
2906 .unwrap();
2907 let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2908 assert_eq!(log2[0].bound_mode_version, None);
2909 }
2910
2911 #[tokio::test]
2916 async fn session_start_binds_and_records_max_suspend_cap() {
2917 let rt = make_runtime();
2918
2919 let sid = new_sid();
2921 let payload = SessionStartPayload {
2922 participants: vec!["alice".into(), "bob".into()],
2923 mode_version: "1.0.0".into(),
2924 configuration_version: "cfg-1".into(),
2925 ttl_ms: 60_000,
2926 max_suspend_ms: 12_345,
2927 ..Default::default()
2928 }
2929 .encode_to_vec();
2930 rt.process(
2931 &env(
2932 "macp.mode.decision.v1",
2933 "SessionStart",
2934 "m1",
2935 &sid,
2936 "alice",
2937 payload,
2938 ),
2939 None,
2940 )
2941 .await
2942 .unwrap();
2943 let log = rt.log_store.get_log(&sid).await.unwrap();
2944 assert_eq!(log[0].bound_max_suspend_ms, Some(12_345));
2945
2946 let sid2 = new_sid();
2948 let payload2 = SessionStartPayload {
2949 participants: vec!["alice".into(), "bob".into()],
2950 mode_version: "1.0.0".into(),
2951 configuration_version: "cfg-1".into(),
2952 ttl_ms: 60_000,
2953 max_suspend_ms: 0,
2954 ..Default::default()
2955 }
2956 .encode_to_vec();
2957 rt.process(
2958 &env(
2959 "macp.mode.decision.v1",
2960 "SessionStart",
2961 "m2",
2962 &sid2,
2963 "alice",
2964 payload2,
2965 ),
2966 None,
2967 )
2968 .await
2969 .unwrap();
2970 let log2 = rt.log_store.get_log(&sid2).await.unwrap();
2971 assert_eq!(
2972 log2[0].bound_max_suspend_ms,
2973 Some(macp_core::session::MAX_SUSPEND_MS)
2974 );
2975 }
2976
2977 #[test]
2978 fn audit_verbosity_reads_policy_rules() {
2979 let mut session = Session::builder("s1", "macp.mode.decision.v1", "a").build();
2980 assert!(!Runtime::audit_verbose(&session));
2981
2982 session.policy_definition = Some(macp_core::policy::PolicyDefinition {
2983 policy_id: "policy.test.audit".into(),
2984 mode: "*".into(),
2985 description: "audited".into(),
2986 rules: serde_json::json!({ "audit": { "level": "info" } }),
2987 schema_version: 1,
2988 });
2989 assert!(Runtime::audit_verbose(&session));
2990
2991 session.policy_definition.as_mut().unwrap().rules =
2992 serde_json::json!({ "audit": { "level": "debug" } });
2993 assert!(!Runtime::audit_verbose(&session));
2994 }
2995
2996 #[tokio::test]
3002 async fn session_start_snapshot_failure_is_nonfatal_after_commit_point() {
3003 use std::io;
3004
3005 struct FailSnapshotBackend;
3006 #[async_trait::async_trait]
3007 impl StorageBackend for FailSnapshotBackend {
3008 async fn create_session_storage(&self, _s: &str) -> io::Result<()> {
3009 Ok(())
3010 }
3011 async fn save_session(&self, _s: &Session) -> io::Result<()> {
3012 Err(io::Error::other("snapshot disk full"))
3013 }
3014 async fn load_session(&self, _s: &str) -> io::Result<Option<Session>> {
3015 Ok(None)
3016 }
3017 async fn load_all_sessions(&self) -> io::Result<Vec<Session>> {
3018 Ok(vec![])
3019 }
3020 async fn delete_session(&self, _s: &str) -> io::Result<()> {
3021 Ok(())
3022 }
3023 async fn list_session_ids(&self) -> io::Result<Vec<String>> {
3024 Ok(vec![])
3025 }
3026 async fn append_log_entry(
3027 &self,
3028 _s: &str,
3029 _e: &crate::log_store::LogEntry,
3030 ) -> io::Result<()> {
3031 Ok(())
3032 }
3033 async fn load_log(&self, _s: &str) -> io::Result<Vec<crate::log_store::LogEntry>> {
3034 Ok(vec![])
3035 }
3036 }
3037
3038 let rt = Runtime::new(
3039 Arc::new(FailSnapshotBackend),
3040 Arc::new(SessionRegistry::new()),
3041 Arc::new(LogStore::new()),
3042 );
3043 let sid = new_sid();
3044 let result = rt
3045 .process(
3046 &env(
3047 "macp.mode.decision.v1",
3048 "SessionStart",
3049 "m1",
3050 &sid,
3051 "agent://orchestrator",
3052 session_start(vec!["agent://orchestrator".into()]),
3053 ),
3054 None,
3055 )
3056 .await
3057 .expect("start must succeed: the log append (commit point) succeeded");
3058 assert!(!result.duplicate);
3059 assert!(rt.get_session_checked(&sid).await.is_some());
3061 }
3062
3063 const HANDOFF_MODE: &str = "macp.mode.handoff.v1";
3079 const OWNER: &str = "agent://owner";
3080 const TARGET: &str = "agent://target";
3081
3082 fn reserved_id(handoff_id: &str) -> String {
3083 format!(
3084 "{}{handoff_id}",
3085 crate::mode::handoff::IMPLICIT_ACCEPT_MESSAGE_ID_PREFIX
3086 )
3087 }
3088
3089 fn handoff_start_payload() -> Vec<u8> {
3090 session_start(vec![OWNER.into(), TARGET.into()])
3091 }
3092
3093 fn handoff_commitment_payload() -> Vec<u8> {
3096 CommitmentPayload {
3097 commitment_id: "c1".into(),
3098 action: "handoff.accepted".into(),
3099 authority_scope: "support".into(),
3100 reason: "bound".into(),
3101 mode_version: "1.0.0".into(),
3102 policy_version: "policy.default".into(),
3103 configuration_version: "cfg-1".into(),
3104 outcome_positive: true,
3105 supersedes: None,
3106 }
3107 .encode_to_vec()
3108 }
3109
3110 fn handoff_offer(handoff_id: &str) -> Vec<u8> {
3111 crate::handoff_pb::HandoffOfferPayload {
3112 handoff_id: handoff_id.into(),
3113 target_participant: TARGET.into(),
3114 scope: "support".into(),
3115 reason: "escalate".into(),
3116 }
3117 .encode_to_vec()
3118 }
3119
3120 fn handoff_context(handoff_id: &str) -> Vec<u8> {
3121 crate::handoff_pb::HandoffContextPayload {
3122 handoff_id: handoff_id.into(),
3123 content_type: "text/plain".into(),
3124 context: b"background".to_vec(),
3125 }
3126 .encode_to_vec()
3127 }
3128
3129 fn handoff_accept(handoff_id: &str, implicit: bool) -> Vec<u8> {
3130 crate::handoff_pb::HandoffAcceptPayload {
3131 handoff_id: handoff_id.into(),
3132 accepted_by: TARGET.into(),
3133 reason: "ready".into(),
3134 implicit,
3135 }
3136 .encode_to_vec()
3137 }
3138
3139 async fn handoff_session_with_offer(rt: &Runtime) -> String {
3142 let sid = new_sid();
3143 rt.process(
3144 &env(
3145 HANDOFF_MODE,
3146 "SessionStart",
3147 "start-1",
3148 &sid,
3149 OWNER,
3150 handoff_start_payload(),
3151 ),
3152 None,
3153 )
3154 .await
3155 .expect("handoff session start");
3156 rt.process(
3157 &env(
3158 HANDOFF_MODE,
3159 "HandoffOffer",
3160 "offer-1",
3161 &sid,
3162 OWNER,
3163 handoff_offer("h1"),
3164 ),
3165 None,
3166 )
3167 .await
3168 .expect("handoff offer");
3169 assert_eq!(
3170 rt.get_session_checked(&sid).await.unwrap().semantics_rev,
3171 macp_core::session::CURRENT_SEMANTICS_REV
3172 );
3173 sid
3174 }
3175
3176 #[tokio::test]
3189 async fn reserved_message_id_namespace_is_rejected_at_rev2() {
3190 let rt = make_runtime();
3191 let sid = handoff_session_with_offer(&rt).await;
3192
3193 let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3194 let dedup_before = rt
3195 .get_session_checked(&sid)
3196 .await
3197 .unwrap()
3198 .seen_message_ids
3199 .clone();
3200
3201 for (message_type, sender, payload) in [
3202 ("HandoffContext", OWNER, handoff_context("h1")),
3203 ("Commitment", OWNER, handoff_commitment_payload()),
3204 ("HandoffAccept", TARGET, handoff_accept("h1", false)),
3205 ] {
3206 let err = rt
3207 .process(
3208 &env(
3209 HANDOFF_MODE,
3210 message_type,
3211 &reserved_id("h1"),
3212 &sid,
3213 sender,
3214 payload,
3215 ),
3216 None,
3217 )
3218 .await
3219 .unwrap_err();
3220 assert!(
3221 matches!(err, MacpError::InvalidEnvelope),
3222 "{message_type} with a reserved id must be InvalidEnvelope, got {err}"
3223 );
3224 }
3225
3226 let session = rt.get_session_checked(&sid).await.unwrap();
3228 assert_eq!(
3229 rt.log_store.get_log(&sid).await.unwrap().len(),
3230 history_before
3231 );
3232 assert_eq!(session.seen_message_ids, dedup_before);
3233 assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3234 assert_eq!(session.state, SessionState::Open);
3235
3236 rt.process(
3239 &env(
3240 HANDOFF_MODE,
3241 "HandoffContext",
3242 "ctx-1",
3243 &sid,
3244 OWNER,
3245 handoff_context("h1"),
3246 ),
3247 None,
3248 )
3249 .await
3250 .expect("an ordinary id is accepted");
3251 assert_eq!(
3252 rt.log_store.get_log(&sid).await.unwrap().len(),
3253 history_before + 1
3254 );
3255 }
3256
3257 #[tokio::test]
3268 async fn reserved_message_id_is_rejected_on_the_session_start_path() {
3269 let rt = make_runtime();
3270 let sid = new_sid();
3271
3272 let err = rt
3273 .process(
3274 &env(
3275 HANDOFF_MODE,
3276 "SessionStart",
3277 &reserved_id("h1"),
3278 &sid,
3279 OWNER,
3280 handoff_start_payload(),
3281 ),
3282 None,
3283 )
3284 .await
3285 .unwrap_err();
3286 assert!(
3287 matches!(err, MacpError::InvalidEnvelope),
3288 "reserved id on SessionStart must be InvalidEnvelope, got {err}"
3289 );
3290
3291 assert!(rt.get_session_checked(&sid).await.is_none());
3293 assert!(rt.log_store.get_log(&sid).await.is_none());
3294 assert!(!rt.registry.sessions.read().await.contains_key(&sid));
3295
3296 rt.process(
3299 &env(
3300 HANDOFF_MODE,
3301 "SessionStart",
3302 "start-1",
3303 &sid,
3304 OWNER,
3305 handoff_start_payload(),
3306 ),
3307 None,
3308 )
3309 .await
3310 .expect("a rejected SessionStart must not reserve the session id");
3311 let session = rt.get_session_checked(&sid).await.unwrap();
3312 assert!(session.seen_message_ids.contains("start-1"));
3313 assert!(!session.seen_message_ids.contains(&reserved_id("h1")));
3314 }
3315
3316 #[tokio::test]
3343 async fn client_implicit_accept_rejected_through_the_runtime() {
3344 let rt = make_runtime();
3345 let sid = handoff_session_with_offer(&rt).await;
3346 let history_before = rt.log_store.get_log(&sid).await.unwrap().len();
3347
3348 let err = rt
3350 .process(
3351 &env(
3352 HANDOFF_MODE,
3353 "HandoffAccept",
3354 "accept-1",
3355 &sid,
3356 TARGET,
3357 handoff_accept("h1", true),
3358 ),
3359 None,
3360 )
3361 .await
3362 .unwrap_err();
3363 assert!(matches!(err, MacpError::InvalidPayload), "got {err}");
3364
3365 let err = rt
3367 .process(
3368 &env(
3369 HANDOFF_MODE,
3370 "HandoffAccept",
3371 &reserved_id("h1"),
3372 &sid,
3373 TARGET,
3374 handoff_accept("h1", true),
3375 ),
3376 None,
3377 )
3378 .await
3379 .unwrap_err();
3380 assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3381
3382 let session = rt.get_session_checked(&sid).await.unwrap();
3384 assert_eq!(
3385 rt.log_store.get_log(&sid).await.unwrap().len(),
3386 history_before
3387 );
3388 assert!(session.seen_message_ids.is_disjoint(
3389 &["accept-1".to_string(), reserved_id("h1")]
3390 .into_iter()
3391 .collect()
3392 ));
3393 let mode_state: serde_json::Value = serde_json::from_slice(&session.mode_state).unwrap();
3394 assert_eq!(mode_state["offers"]["h1"]["disposition"], "Offered");
3395
3396 rt.process(
3400 &env(
3401 HANDOFF_MODE,
3402 "HandoffAccept",
3403 "accept-2",
3404 &sid,
3405 TARGET,
3406 handoff_accept("h1", false),
3407 ),
3408 None,
3409 )
3410 .await
3411 .expect("an explicit accept is still accepted");
3412 }
3413
3414 #[tokio::test]
3420 async fn client_boundary_error_ordering_is_unchanged_at_rev2() {
3421 let rt = make_runtime();
3422 let sid = handoff_session_with_offer(&rt).await;
3423
3424 let err = rt
3426 .process(
3427 &env(
3428 HANDOFF_MODE,
3429 "HandoffAccept",
3430 &reserved_id("h1"),
3431 &sid,
3432 "agent://stranger",
3433 handoff_accept("h1", true),
3434 ),
3435 None,
3436 )
3437 .await
3438 .unwrap_err();
3439 assert!(
3440 matches!(err, MacpError::Forbidden),
3441 "authorization must be reported before the client boundary, got {err}"
3442 );
3443
3444 let err = rt
3446 .process(
3447 &env(
3448 HANDOFF_MODE,
3449 "HandoffAccept",
3450 &reserved_id("h1"),
3451 &sid,
3452 TARGET,
3453 handoff_accept("h1", true),
3454 ),
3455 None,
3456 )
3457 .await
3458 .unwrap_err();
3459 assert!(matches!(err, MacpError::InvalidEnvelope), "got {err}");
3460 }
3461
3462 async fn handoff_session_with_timed_offer(rt: &Runtime, timeout_ms: i64) -> String {
3468 rt.register_policy(macp_core::policy::PolicyDefinition {
3469 policy_id: "handoff-timed".into(),
3470 mode: HANDOFF_MODE.into(),
3471 description: "implicit accept".into(),
3472 rules: serde_json::json!({
3473 "acceptance": { "implicit_accept_timeout_ms": timeout_ms },
3474 "commitment": { "authority": "initiator_only" }
3475 }),
3476 schema_version: 1,
3477 })
3478 .expect("policy registers");
3479
3480 let sid = new_sid();
3481 let start = SessionStartPayload {
3482 intent: "escalate".into(),
3483 participants: vec![OWNER.into(), TARGET.into()],
3484 mode_version: "1.0.0".into(),
3485 configuration_version: "cfg-1".into(),
3486 policy_version: "handoff-timed".into(),
3487 ttl_ms: 60_000,
3488 context_id: String::new(),
3489 extensions: std::collections::HashMap::new(),
3490 roots: vec![],
3491 max_suspend_ms: 0,
3492 }
3493 .encode_to_vec();
3494 rt.process(
3495 &env(HANDOFF_MODE, "SessionStart", "start-1", &sid, OWNER, start),
3496 None,
3497 )
3498 .await
3499 .expect("session start");
3500 rt.process(
3501 &env(
3502 HANDOFF_MODE,
3503 "HandoffOffer",
3504 "offer-1",
3505 &sid,
3506 OWNER,
3507 handoff_offer("h1"),
3508 ),
3509 None,
3510 )
3511 .await
3512 .expect("offer");
3513 sid
3514 }
3515
3516 #[tokio::test]
3536 async fn synthesis_is_skipped_for_a_non_open_session() {
3537 let rt = make_runtime();
3538 let sid = handoff_session_with_timed_offer(&rt, 20).await;
3539 rt.suspend_session(&sid, "hold", OWNER)
3540 .await
3541 .expect("suspend");
3542
3543 let shared = rt.registry.get_shared(&sid).await.unwrap();
3544 let mut guard = shared.lock().await;
3545 let session = &mut *guard;
3546 assert_eq!(session.state, SessionState::Suspended);
3547 assert!(session.suspended_at_ms.is_some());
3548
3549 let long_after = session.suspended_at_ms.unwrap() + 10_000;
3552 let log_before = rt.log_store.get_log(&sid).await.unwrap_or_default().len();
3553 let dedup_before = session.seen_message_ids.len();
3554 let mode_state_before = session.mode_state.clone();
3555
3556 rt.synthesize_due_accept(&sid, session, long_after)
3557 .await
3558 .expect("the filter is a skip, not an error");
3559
3560 assert_eq!(
3561 rt.log_store.get_log(&sid).await.unwrap_or_default().len(),
3562 log_before,
3563 "a suspended session must not gain a synthetic entry"
3564 );
3565 assert_eq!(session.seen_message_ids.len(), dedup_before);
3566 assert_eq!(session.mode_state, mode_state_before);
3567
3568 let mut without_the_tripwire = session.clone();
3580 without_the_tripwire.state = SessionState::Open;
3581 without_the_tripwire.suspended_at_ms = None;
3582 let mode = rt.mode_registry.get_mode(&session.mode).unwrap();
3583 let would_have_emitted = mode
3584 .due_synthetic_envelope(&without_the_tripwire, long_after)
3585 .expect("the mode would have synthesized; only the kernel filter stopped it");
3586 let suspended_at = session.suspended_at_ms.unwrap();
3592 assert!(
3593 would_have_emitted.timestamp_unix_ms >= suspended_at
3594 && would_have_emitted.timestamp_unix_ms < long_after,
3595 "D {} must fall inside the still-open pause starting at {suspended_at}",
3596 would_have_emitted.timestamp_unix_ms
3597 );
3598
3599 drop(guard);
3602 rt.resume_session(&sid, "go", OWNER).await.expect("resume");
3603 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3604 let shared = rt.registry.get_shared(&sid).await.unwrap();
3605 let mut guard = shared.lock().await;
3606 let session = &mut *guard;
3607 let now = Utc::now().timestamp_millis();
3608 rt.synthesize_due_accept(&sid, session, now).await.unwrap();
3609 assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3610 }
3611
3612 #[tokio::test]
3622 async fn synthetic_entry_stamps_received_at_with_the_deadline() {
3623 let rt = make_runtime();
3624 let sid = handoff_session_with_timed_offer(&rt, 20).await;
3625 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3626
3627 let offer_received_at = rt
3628 .log_store
3629 .get_log(&sid)
3630 .await
3631 .unwrap()
3632 .iter()
3633 .find(|e| e.message_type == "HandoffOffer")
3634 .expect("offer entry")
3635 .received_at_ms;
3636 let expected_d = offer_received_at + 20;
3637
3638 let shared = rt.registry.get_shared(&sid).await.unwrap();
3639 let mut guard = shared.lock().await;
3640 let session = &mut *guard;
3641 let observed = Utc::now().timestamp_millis();
3643 assert!(observed > expected_d);
3644 rt.synthesize_due_accept(&sid, session, observed)
3645 .await
3646 .unwrap();
3647 drop(guard);
3648
3649 let entry = rt
3650 .log_store
3651 .get_log(&sid)
3652 .await
3653 .unwrap()
3654 .into_iter()
3655 .find(|e| e.message_id == reserved_id("h1"))
3656 .expect("the synthetic entry");
3657 assert_eq!(entry.timestamp_unix_ms, expected_d, "envelope clock is D");
3658 assert_eq!(entry.received_at_ms, expected_d, "entry clock is D");
3659 assert_ne!(
3660 entry.received_at_ms, observed,
3661 "received_at_ms must not be the observation time"
3662 );
3663 assert_eq!(entry.entry_kind, EntryKind::Incoming);
3664 }
3665
3666 #[tokio::test]
3676 async fn synthetic_accept_is_not_credited_as_participant_activity() {
3677 let rt = make_runtime();
3678 let sid = handoff_session_with_timed_offer(&rt, 20).await;
3679 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3680
3681 let shared = rt.registry.get_shared(&sid).await.unwrap();
3682 let mut guard = shared.lock().await;
3683 let session = &mut *guard;
3684 let before = session.participant_message_counts.get(TARGET).copied();
3685 rt.synthesize_due_accept(&sid, session, Utc::now().timestamp_millis())
3686 .await
3687 .unwrap();
3688 assert!(session.seen_message_ids.contains(&reserved_id("h1")));
3689 assert_eq!(
3690 session.participant_message_counts.get(TARGET).copied(),
3691 before,
3692 "the target must not be credited with a message they did not send"
3693 );
3694 }
3695 #[tokio::test]
3722 async fn a_synthetic_entry_on_the_checkpoint_boundary_checkpoints_either_path() {
3723 async fn shape(rt: &Runtime, sid: &str) -> Vec<(EntryKind, String)> {
3724 rt.log_store
3725 .get_log(sid)
3726 .await
3727 .expect("log")
3728 .iter()
3729 .map(|e| (e.entry_kind.clone(), e.message_type.clone()))
3730 .collect()
3731 }
3732
3733 let mut eager = make_runtime();
3735 eager.checkpoint_interval = 3;
3736 let eager_sid = handoff_session_with_timed_offer(&eager, 20).await;
3737 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3738 assert_eq!(
3739 eager.sweep_due_synthetic_accepts().await,
3740 1,
3741 "sweep emitted"
3742 );
3743
3744 let mut lazy = make_runtime();
3748 lazy.checkpoint_interval = 3;
3749 let lazy_sid = handoff_session_with_timed_offer(&lazy, 20).await;
3750 tokio::time::sleep(std::time::Duration::from_millis(60)).await;
3751 lazy.process(
3752 &env(
3753 HANDOFF_MODE,
3754 "HandoffContext",
3755 "ctx-1",
3756 &lazy_sid,
3757 OWNER,
3758 handoff_context("h1"),
3759 ),
3760 None,
3761 )
3762 .await
3763 .expect("context accepted");
3764
3765 let expected = vec![
3766 (EntryKind::Incoming, "SessionStart".to_string()),
3767 (EntryKind::Incoming, "HandoffOffer".to_string()),
3768 (EntryKind::Incoming, "HandoffAccept".to_string()),
3769 (EntryKind::Checkpoint, "Checkpoint".to_string()),
3770 ];
3771 assert_eq!(shape(&eager, &eager_sid).await, expected, "eager sweep");
3772
3773 let lazy_shape = shape(&lazy, &lazy_sid).await;
3774 assert_eq!(
3775 lazy_shape[..4],
3776 expected[..],
3777 "the lazy path must checkpoint in the same place as the eager one"
3778 );
3779 assert_eq!(
3782 lazy_shape[4..],
3783 [(EntryKind::Incoming, "HandoffContext".to_string())],
3784 "no second checkpoint for the trigger's own append"
3785 );
3786 }
3787}