1use crate::mode::util::{
2 check_commitment_authority, is_declared_participant, validate_commitment_payload_for_session,
3};
4use crate::mode::{Mode, ModeResponse};
5use macp_core::error::MacpError;
6use macp_core::session::Session;
7use macp_pb::pb::Envelope;
8use macp_pb::task_pb::{
9 TaskAcceptPayload, TaskCompletePayload, TaskFailPayload, TaskRejectPayload, TaskRequestPayload,
10 TaskUpdatePayload,
11};
12use prost::Message;
13use serde::{Deserialize, Serialize};
14
15#[derive(Debug, Clone, Serialize, Deserialize)]
16pub struct TaskRecord {
17 pub task_id: String,
18 pub title: String,
19 pub instructions: String,
20 pub requested_assignee: String,
21 pub input: Vec<u8>,
22 pub deadline_unix_ms: i64,
23 pub requester: String,
24}
25
26#[derive(Debug, Clone, Serialize, Deserialize)]
27pub struct TaskRejectRecord {
28 pub task_id: String,
29 pub assignee: String,
30 pub reason: String,
31}
32
33#[derive(Debug, Clone, Serialize, Deserialize)]
34pub struct TaskUpdateRecord {
35 pub task_id: String,
36 pub status: String,
37 pub progress: f64,
38 pub message: String,
39 pub partial_output: Vec<u8>,
40 pub sender: String,
41}
42
43#[derive(Debug, Clone, Serialize, Deserialize)]
44pub struct TaskCompleteRecord {
45 pub task_id: String,
46 pub assignee: String,
47 pub output: Vec<u8>,
48 pub summary: String,
49}
50
51#[derive(Debug, Clone, Serialize, Deserialize)]
52pub struct TaskFailRecord {
53 pub task_id: String,
54 pub assignee: String,
55 pub error_code: String,
56 pub reason: String,
57 pub retryable: bool,
58}
59
60#[derive(Debug, Clone, Serialize, Deserialize)]
61pub enum TaskTerminalReport {
62 Complete(TaskCompleteRecord),
63 Fail(TaskFailRecord),
64}
65
66#[derive(Debug, Clone, Serialize, Deserialize, Default)]
67pub struct TaskState {
68 pub task: Option<TaskRecord>,
69 pub active_assignee: Option<String>,
70 pub rejections: Vec<TaskRejectRecord>,
71 pub updates: Vec<TaskUpdateRecord>,
72 pub terminal_report: Option<TaskTerminalReport>,
73}
74
75pub struct TaskMode {
76 evaluator: std::sync::Arc<dyn macp_core::policy::PolicyEvaluator>,
77}
78
79impl TaskMode {
80 pub fn new(evaluator: std::sync::Arc<dyn macp_core::policy::PolicyEvaluator>) -> Self {
82 Self { evaluator }
83 }
84
85 fn encode_state(state: &TaskState) -> Vec<u8> {
86 serde_json::to_vec(state).expect("TaskState is always serializable")
87 }
88
89 fn decode_state(data: &[u8]) -> Result<TaskState, MacpError> {
90 serde_json::from_slice(data).map_err(|_| MacpError::InvalidModeState)
91 }
92
93 fn ensure_task_matches(expected_task_id: &str, actual_task_id: &str) -> Result<(), MacpError> {
94 if expected_task_id.is_empty() || expected_task_id != actual_task_id {
95 return Err(MacpError::InvalidPayload);
96 }
97 Ok(())
98 }
99
100 fn can_assignee_respond(session: &Session, task: &TaskRecord, sender: &str) -> bool {
101 if !task.requested_assignee.is_empty() {
102 sender == task.requested_assignee
103 } else {
104 is_declared_participant(&session.participants, sender)
105 && sender != session.initiator_sender
106 }
107 }
108}
109
110impl Mode for TaskMode {
111 fn authorize_sender(&self, session: &Session, env: &Envelope) -> Result<(), MacpError> {
112 match env.message_type.as_str() {
113 "TaskRequest" if env.sender == session.initiator_sender => Ok(()),
114 "TaskRequest" => Err(MacpError::Forbidden),
115 "Commitment" => check_commitment_authority(session, &env.sender),
116 _ if is_declared_participant(&session.participants, &env.sender) => Ok(()),
117 _ => Err(MacpError::Forbidden),
118 }
119 }
120
121 fn on_session_start(
122 &self,
123 session: &Session,
124 _env: &Envelope,
125 ) -> Result<ModeResponse, MacpError> {
126 if session.participants.len() < 2 {
127 return Err(MacpError::InvalidPayload);
128 }
129 if !session
130 .participants
131 .iter()
132 .any(|p| p == &session.initiator_sender)
133 {
134 return Err(MacpError::InvalidPayload);
135 }
136 Ok(ModeResponse::PersistState(Self::encode_state(
137 &TaskState::default(),
138 )))
139 }
140
141 fn on_message(&self, session: &Session, env: &Envelope) -> Result<ModeResponse, MacpError> {
142 let mut state = if session.mode_state.is_empty() {
143 TaskState::default()
144 } else {
145 Self::decode_state(&session.mode_state)?
146 };
147
148 match env.message_type.as_str() {
149 "TaskRequest" => {
150 if env.sender != session.initiator_sender {
151 return Err(MacpError::Forbidden);
152 }
153 let payload = TaskRequestPayload::decode(&*env.payload)
154 .map_err(|_| MacpError::InvalidPayload)?;
155 if state.task.is_some() || payload.task_id.is_empty() {
156 return Err(MacpError::InvalidPayload);
157 }
158 if !payload.requested_assignee.is_empty()
159 && !is_declared_participant(&session.participants, &payload.requested_assignee)
160 {
161 return Err(MacpError::InvalidPayload);
162 }
163 state.task = Some(TaskRecord {
164 task_id: payload.task_id,
165 title: payload.title,
166 instructions: payload.instructions,
167 requested_assignee: payload.requested_assignee,
168 input: payload.input,
169 deadline_unix_ms: payload.deadline_unix_ms,
170 requester: env.sender.clone(),
171 });
172 Ok(ModeResponse::PersistState(Self::encode_state(&state)))
173 }
174 "TaskAccept" => {
175 let payload = TaskAcceptPayload::decode(&*env.payload)
176 .map_err(|_| MacpError::InvalidPayload)?;
177 let task = state.task.as_ref().ok_or(MacpError::InvalidPayload)?;
178 Self::ensure_task_matches(&payload.task_id, &task.task_id)?;
179 if state.active_assignee.is_some() {
180 return Err(MacpError::InvalidPayload);
181 }
182 if !state.rejections.is_empty() {
184 let allow = session.policy_definition.as_ref().is_some_and(|p| {
185 serde_json::from_value::<macp_core::policy::rules::TaskPolicyRules>(
186 p.rules.clone(),
187 )
188 .unwrap_or_default()
189 .assignment
190 .allow_reassignment_on_reject
191 });
192 if !allow {
193 return Err(MacpError::PolicyDenied {
194 reasons: vec![
195 "reassignment after rejection not allowed by policy".into()
196 ],
197 });
198 }
199 }
200 if !payload.assignee.is_empty() && payload.assignee != env.sender {
201 return Err(MacpError::InvalidPayload);
202 }
203 if !Self::can_assignee_respond(session, task, &env.sender) {
204 return Err(MacpError::Forbidden);
205 }
206 state.active_assignee = Some(env.sender.clone());
207 Ok(ModeResponse::PersistState(Self::encode_state(&state)))
208 }
209 "TaskReject" => {
210 let payload = TaskRejectPayload::decode(&*env.payload)
211 .map_err(|_| MacpError::InvalidPayload)?;
212 let task = state.task.as_ref().ok_or(MacpError::InvalidPayload)?;
213 Self::ensure_task_matches(&payload.task_id, &task.task_id)?;
214 if let Some(ref active) = state.active_assignee {
219 if active == &env.sender {
220 let allow = session.policy_definition.as_ref().is_some_and(|p| {
221 serde_json::from_value::<macp_core::policy::rules::TaskPolicyRules>(
222 p.rules.clone(),
223 )
224 .unwrap_or_default()
225 .assignment
226 .allow_reassignment_on_reject
227 });
228 if !allow {
229 return Err(MacpError::PolicyDenied {
230 reasons: vec![
231 "active assignee cannot reject without allow_reassignment_on_reject policy".into(),
232 ],
233 });
234 }
235 } else {
237 return Err(MacpError::InvalidPayload);
239 }
240 }
241 if !payload.assignee.is_empty() && payload.assignee != env.sender {
242 return Err(MacpError::InvalidPayload);
243 }
244 if state.active_assignee.is_none()
245 && !Self::can_assignee_respond(session, task, &env.sender)
246 {
247 return Err(MacpError::Forbidden);
248 }
249 if state.rejections.iter().any(|r| r.assignee == env.sender) {
250 return Err(MacpError::InvalidPayload);
251 }
252 if state.active_assignee.as_deref() == Some(env.sender.as_str()) {
254 state.active_assignee = None;
255 }
256 state.rejections.push(TaskRejectRecord {
257 task_id: payload.task_id,
258 assignee: env.sender.clone(),
259 reason: payload.reason,
260 });
261 Ok(ModeResponse::PersistState(Self::encode_state(&state)))
262 }
263 "TaskUpdate" => {
264 let payload = TaskUpdatePayload::decode(&*env.payload)
265 .map_err(|_| MacpError::InvalidPayload)?;
266 let task = state.task.as_ref().ok_or(MacpError::InvalidPayload)?;
267 Self::ensure_task_matches(&payload.task_id, &task.task_id)?;
268 if state.terminal_report.is_some()
269 || state.active_assignee.as_deref() != Some(env.sender.as_str())
270 {
271 return Err(MacpError::Forbidden);
272 }
273 state.updates.push(TaskUpdateRecord {
274 task_id: payload.task_id,
275 status: payload.status,
276 progress: payload.progress,
277 message: payload.message,
278 partial_output: payload.partial_output,
279 sender: env.sender.clone(),
280 });
281 Ok(ModeResponse::PersistState(Self::encode_state(&state)))
282 }
283 "TaskComplete" => {
284 let payload = TaskCompletePayload::decode(&*env.payload)
285 .map_err(|_| MacpError::InvalidPayload)?;
286 let task = state.task.as_ref().ok_or(MacpError::InvalidPayload)?;
287 Self::ensure_task_matches(&payload.task_id, &task.task_id)?;
288 if state.terminal_report.is_some()
289 || state.active_assignee.as_deref() != Some(env.sender.as_str())
290 {
291 return Err(MacpError::Forbidden);
292 }
293 if !payload.assignee.is_empty() && payload.assignee != env.sender {
294 return Err(MacpError::InvalidPayload);
295 }
296 state.terminal_report = Some(TaskTerminalReport::Complete(TaskCompleteRecord {
297 task_id: payload.task_id,
298 assignee: env.sender.clone(),
299 output: payload.output,
300 summary: payload.summary,
301 }));
302 Ok(ModeResponse::PersistState(Self::encode_state(&state)))
303 }
304 "TaskFail" => {
305 let payload = TaskFailPayload::decode(&*env.payload)
306 .map_err(|_| MacpError::InvalidPayload)?;
307 let task = state.task.as_ref().ok_or(MacpError::InvalidPayload)?;
308 Self::ensure_task_matches(&payload.task_id, &task.task_id)?;
309 if state.terminal_report.is_some()
310 || state.active_assignee.as_deref() != Some(env.sender.as_str())
311 {
312 return Err(MacpError::Forbidden);
313 }
314 if !payload.assignee.is_empty() && payload.assignee != env.sender {
315 return Err(MacpError::InvalidPayload);
316 }
317 state.terminal_report = Some(TaskTerminalReport::Fail(TaskFailRecord {
318 task_id: payload.task_id,
319 assignee: env.sender.clone(),
320 error_code: payload.error_code,
321 reason: payload.reason,
322 retryable: payload.retryable,
323 }));
324 Ok(ModeResponse::PersistState(Self::encode_state(&state)))
325 }
326 "Commitment" => {
327 validate_commitment_payload_for_session(session, &env.payload)?;
328 if state.terminal_report.is_none() {
329 return Err(MacpError::InvalidPayload);
330 }
331 if let Some(ref policy) = session.policy_definition {
333 let has_output = matches!(
334 &state.terminal_report,
335 Some(TaskTerminalReport::Complete(record)) if !record.output.is_empty()
336 );
337 let decision = self.evaluator.evaluate_task_commitment(policy, has_output);
338 if let macp_core::policy::PolicyDecision::Deny { reasons } = decision {
339 tracing::warn!(
340 session_id = %session.session_id,
341 policy_id = %policy.policy_id,
342 reasons = ?reasons,
343 "policy denied commitment"
344 );
345 return Err(MacpError::PolicyDenied { reasons });
346 }
347 }
348 Ok(ModeResponse::PersistAndResolve {
349 state: Self::encode_state(&state),
350 resolution: env.payload.clone(),
351 })
352 }
353 _ => Err(MacpError::InvalidPayload),
354 }
355 }
356}
357
358#[cfg(test)]
359mod tests {
360 use super::*;
361 use macp_core::session::{Session, SessionState};
362 use macp_pb::pb::CommitmentPayload;
363 use std::collections::HashSet;
364
365 fn base_session() -> Session {
366 Session {
367 session_id: "s1".into(),
368 state: SessionState::Open,
369 ttl_expiry: i64::MAX,
370 ttl_ms: 60_000,
371 started_at_unix_ms: 0,
372 resolution: None,
373 mode: "macp.mode.task.v1".into(),
374 mode_state: vec![],
375 participants: vec!["planner".into(), "worker".into()],
376 seen_message_ids: HashSet::new(),
377 intent: String::new(),
378 mode_version: "1.0.0".into(),
379 configuration_version: "config".into(),
380 policy_version: "policy".into(),
381 context_id: String::new(),
382 extensions: std::collections::HashMap::new(),
383 roots: vec![],
384 initiator_sender: "planner".into(),
385 participant_message_counts: std::collections::HashMap::new(),
386 participant_last_seen: std::collections::HashMap::new(),
387 policy_definition: None,
388 suspended_at_ms: None,
389 accumulated_suspended_ms: 0,
390 }
391 }
392
393 fn env(sender: &str, message_type: &str, payload: Vec<u8>) -> Envelope {
394 Envelope {
395 macp_version: "1.0".into(),
396 mode: "macp.mode.task.v1".into(),
397 message_type: message_type.into(),
398 message_id: format!("{}-{}", sender, message_type),
399 session_id: "s1".into(),
400 sender: sender.into(),
401 timestamp_unix_ms: 0,
402 payload,
403 }
404 }
405
406 fn commitment_payload() -> Vec<u8> {
407 CommitmentPayload {
408 commitment_id: "c1".into(),
409 action: "task.completed".into(),
410 authority_scope: "ops".into(),
411 reason: "done".into(),
412 mode_version: "1.0.0".into(),
413 policy_version: "policy".into(),
414 configuration_version: "config".into(),
415 outcome_positive: true,
416 supersedes: None,
417 }
418 .encode_to_vec()
419 }
420
421 fn apply(session: &mut Session, result: ModeResponse) {
422 match result {
423 ModeResponse::PersistState(data) => session.mode_state = data,
424 ModeResponse::PersistAndResolve { state, .. } => session.mode_state = state,
425 _ => {}
426 }
427 }
428
429 fn make_task_request(task_id: &str, assignee: &str) -> Vec<u8> {
430 TaskRequestPayload {
431 task_id: task_id.into(),
432 title: "Test Task".into(),
433 instructions: "Do the thing".into(),
434 requested_assignee: assignee.into(),
435 input: vec![],
436 deadline_unix_ms: 0,
437 }
438 .encode_to_vec()
439 }
440
441 fn make_task_accept(task_id: &str, assignee: &str) -> Vec<u8> {
442 TaskAcceptPayload {
443 task_id: task_id.into(),
444 assignee: assignee.into(),
445 reason: "ready".into(),
446 }
447 .encode_to_vec()
448 }
449
450 fn make_task_reject(task_id: &str, assignee: &str) -> Vec<u8> {
451 TaskRejectPayload {
452 task_id: task_id.into(),
453 assignee: assignee.into(),
454 reason: "busy".into(),
455 }
456 .encode_to_vec()
457 }
458
459 fn make_task_update(task_id: &str) -> Vec<u8> {
460 TaskUpdatePayload {
461 task_id: task_id.into(),
462 status: "in_progress".into(),
463 progress: 0.5,
464 message: "halfway done".into(),
465 partial_output: vec![],
466 }
467 .encode_to_vec()
468 }
469
470 fn make_task_complete(task_id: &str, assignee: &str) -> Vec<u8> {
471 TaskCompletePayload {
472 task_id: task_id.into(),
473 assignee: assignee.into(),
474 output: b"result".to_vec(),
475 summary: "done".into(),
476 }
477 .encode_to_vec()
478 }
479
480 fn make_task_fail(task_id: &str, assignee: &str) -> Vec<u8> {
481 TaskFailPayload {
482 task_id: task_id.into(),
483 assignee: assignee.into(),
484 error_code: "E001".into(),
485 reason: "failed".into(),
486 retryable: true,
487 }
488 .encode_to_vec()
489 }
490
491 #[test]
494 fn session_start_initializes_state() {
495 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
496 let session = base_session();
497 let result = mode
498 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
499 .unwrap();
500 match result {
501 ModeResponse::PersistState(data) => {
502 let state: TaskState = serde_json::from_slice(&data).unwrap();
503 assert!(state.task.is_none());
504 assert!(state.active_assignee.is_none());
505 }
506 _ => panic!("Expected PersistState"),
507 }
508 }
509
510 #[test]
511 fn session_start_requires_at_least_two_participants() {
512 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
513 let mut session = base_session();
514 session.participants = vec!["planner".into()]; let err = mode
516 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
517 .unwrap_err();
518 assert_eq!(err.to_string(), "InvalidPayload");
519 }
520
521 #[test]
522 fn session_start_rejects_empty_participants() {
523 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
524 let mut session = base_session();
525 session.participants.clear();
526 let err = mode
527 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
528 .unwrap_err();
529 assert_eq!(err.to_string(), "InvalidPayload");
530 }
531
532 #[test]
533 fn session_start_rejects_when_initiator_not_participant() {
534 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
535 let mut session = base_session();
536 session.participants = vec!["worker".into(), "other".into()]; let err = mode
538 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
539 .unwrap_err();
540 assert_eq!(err.to_string(), "InvalidPayload");
541 }
542
543 #[test]
546 fn task_request_from_initiator() {
547 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
548 let mut session = base_session();
549 let result = mode
550 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
551 .unwrap();
552 apply(&mut session, result);
553 let result = mode
554 .on_message(
555 &session,
556 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
557 )
558 .unwrap();
559 match result {
560 ModeResponse::PersistState(data) => {
561 let state: TaskState = serde_json::from_slice(&data).unwrap();
562 assert!(state.task.is_some());
563 assert_eq!(state.task.unwrap().task_id, "t1");
564 }
565 _ => panic!("Expected PersistState"),
566 }
567 }
568
569 #[test]
570 fn duplicate_task_request_rejected() {
571 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
572 let mut session = base_session();
573 let result = mode
574 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
575 .unwrap();
576 apply(&mut session, result);
577 let result = mode
578 .on_message(
579 &session,
580 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
581 )
582 .unwrap();
583 apply(&mut session, result);
584 let err = mode
585 .on_message(
586 &session,
587 &env("planner", "TaskRequest", make_task_request("t2", "worker")),
588 )
589 .unwrap_err();
590 assert_eq!(err.to_string(), "InvalidPayload");
591 }
592
593 #[test]
594 fn non_initiator_task_request_rejected() {
595 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
596 let mut session = base_session();
597 let result = mode
598 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
599 .unwrap();
600 apply(&mut session, result);
601 let err = mode
602 .on_message(
603 &session,
604 &env("worker", "TaskRequest", make_task_request("t1", "worker")),
605 )
606 .unwrap_err();
607 assert_eq!(err.to_string(), "Forbidden");
608 }
609
610 #[test]
613 fn correct_assignee_can_accept() {
614 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
615 let mut session = base_session();
616 let result = mode
617 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
618 .unwrap();
619 apply(&mut session, result);
620 let result = mode
621 .on_message(
622 &session,
623 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
624 )
625 .unwrap();
626 apply(&mut session, result);
627 let result = mode
628 .on_message(
629 &session,
630 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
631 )
632 .unwrap();
633 match result {
634 ModeResponse::PersistState(data) => {
635 let state: TaskState = serde_json::from_slice(&data).unwrap();
636 assert_eq!(state.active_assignee, Some("worker".into()));
637 }
638 _ => panic!("Expected PersistState"),
639 }
640 }
641
642 #[test]
643 fn wrong_assignee_cannot_accept() {
644 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
645 let mut session = base_session();
646 let result = mode
647 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
648 .unwrap();
649 apply(&mut session, result);
650 let result = mode
651 .on_message(
652 &session,
653 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
654 )
655 .unwrap();
656 apply(&mut session, result);
657 let err = mode
658 .on_message(
659 &session,
660 &env("planner", "TaskAccept", make_task_accept("t1", "planner")),
661 )
662 .unwrap_err();
663 assert_eq!(err.to_string(), "Forbidden");
664 }
665
666 #[test]
667 fn assignee_can_reject() {
668 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
669 let mut session = base_session();
670 let result = mode
671 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
672 .unwrap();
673 apply(&mut session, result);
674 let result = mode
675 .on_message(
676 &session,
677 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
678 )
679 .unwrap();
680 apply(&mut session, result);
681 let result = mode
682 .on_message(
683 &session,
684 &env("worker", "TaskReject", make_task_reject("t1", "worker")),
685 )
686 .unwrap();
687 match result {
688 ModeResponse::PersistState(data) => {
689 let state: TaskState = serde_json::from_slice(&data).unwrap();
690 assert_eq!(state.rejections.len(), 1);
691 assert!(state.active_assignee.is_none()); }
693 _ => panic!("Expected PersistState"),
694 }
695 }
696
697 #[test]
700 fn active_assignee_can_update() {
701 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
702 let mut session = base_session();
703 let result = mode
704 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
705 .unwrap();
706 apply(&mut session, result);
707 let result = mode
708 .on_message(
709 &session,
710 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
711 )
712 .unwrap();
713 apply(&mut session, result);
714 let result = mode
715 .on_message(
716 &session,
717 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
718 )
719 .unwrap();
720 apply(&mut session, result);
721 let result = mode
722 .on_message(
723 &session,
724 &env("worker", "TaskUpdate", make_task_update("t1")),
725 )
726 .unwrap();
727 match result {
728 ModeResponse::PersistState(data) => {
729 let state: TaskState = serde_json::from_slice(&data).unwrap();
730 assert_eq!(state.updates.len(), 1);
731 assert_eq!(state.updates[0].task_id, "t1");
732 assert_eq!(state.updates[0].sender, "worker");
733 assert_eq!(state.updates[0].status, "in_progress");
734 }
735 _ => panic!("Expected PersistState"),
736 }
737 }
738
739 #[test]
740 fn non_assignee_cannot_update() {
741 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
742 let mut session = base_session();
743 let result = mode
744 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
745 .unwrap();
746 apply(&mut session, result);
747 let result = mode
748 .on_message(
749 &session,
750 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
751 )
752 .unwrap();
753 apply(&mut session, result);
754 let result = mode
755 .on_message(
756 &session,
757 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
758 )
759 .unwrap();
760 apply(&mut session, result);
761 let err = mode
762 .on_message(
763 &session,
764 &env("planner", "TaskUpdate", make_task_update("t1")),
765 )
766 .unwrap_err();
767 assert_eq!(err.to_string(), "Forbidden");
768 }
769
770 #[test]
773 fn assignee_can_complete() {
774 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
775 let mut session = base_session();
776 let result = mode
777 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
778 .unwrap();
779 apply(&mut session, result);
780 let result = mode
781 .on_message(
782 &session,
783 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
784 )
785 .unwrap();
786 apply(&mut session, result);
787 let result = mode
788 .on_message(
789 &session,
790 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
791 )
792 .unwrap();
793 apply(&mut session, result);
794 let result = mode
795 .on_message(
796 &session,
797 &env("worker", "TaskComplete", make_task_complete("t1", "worker")),
798 )
799 .unwrap();
800 match result {
801 ModeResponse::PersistState(data) => {
802 let state: TaskState = serde_json::from_slice(&data).unwrap();
803 assert!(matches!(
804 state.terminal_report,
805 Some(TaskTerminalReport::Complete(_))
806 ));
807 }
808 _ => panic!("Expected PersistState"),
809 }
810 }
811
812 #[test]
813 fn assignee_can_fail() {
814 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
815 let mut session = base_session();
816 let result = mode
817 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
818 .unwrap();
819 apply(&mut session, result);
820 let result = mode
821 .on_message(
822 &session,
823 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
824 )
825 .unwrap();
826 apply(&mut session, result);
827 let result = mode
828 .on_message(
829 &session,
830 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
831 )
832 .unwrap();
833 apply(&mut session, result);
834 let result = mode
835 .on_message(
836 &session,
837 &env("worker", "TaskFail", make_task_fail("t1", "worker")),
838 )
839 .unwrap();
840 match result {
841 ModeResponse::PersistState(data) => {
842 let state: TaskState = serde_json::from_slice(&data).unwrap();
843 assert!(matches!(
844 state.terminal_report,
845 Some(TaskTerminalReport::Fail(_))
846 ));
847 }
848 _ => panic!("Expected PersistState"),
849 }
850 }
851
852 #[test]
855 fn commitment_after_complete() {
856 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
857 let mut session = base_session();
858 let result = mode
859 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
860 .unwrap();
861 apply(&mut session, result);
862 let result = mode
863 .on_message(
864 &session,
865 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
866 )
867 .unwrap();
868 apply(&mut session, result);
869 let result = mode
870 .on_message(
871 &session,
872 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
873 )
874 .unwrap();
875 apply(&mut session, result);
876 let result = mode
877 .on_message(
878 &session,
879 &env("worker", "TaskComplete", make_task_complete("t1", "worker")),
880 )
881 .unwrap();
882 apply(&mut session, result);
883 let result = mode
884 .on_message(
885 &session,
886 &env("planner", "Commitment", commitment_payload()),
887 )
888 .unwrap();
889 assert!(matches!(result, ModeResponse::PersistAndResolve { .. }));
890 }
891
892 #[test]
893 fn commitment_after_fail() {
894 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
895 let mut session = base_session();
896 let result = mode
897 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
898 .unwrap();
899 apply(&mut session, result);
900 let result = mode
901 .on_message(
902 &session,
903 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
904 )
905 .unwrap();
906 apply(&mut session, result);
907 let result = mode
908 .on_message(
909 &session,
910 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
911 )
912 .unwrap();
913 apply(&mut session, result);
914 let result = mode
915 .on_message(
916 &session,
917 &env("worker", "TaskFail", make_task_fail("t1", "worker")),
918 )
919 .unwrap();
920 apply(&mut session, result);
921 let result = mode
922 .on_message(
923 &session,
924 &env("planner", "Commitment", commitment_payload()),
925 )
926 .unwrap();
927 assert!(matches!(result, ModeResponse::PersistAndResolve { .. }));
928 }
929
930 #[test]
931 fn commitment_before_terminal_report_rejected() {
932 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
933 let mut session = base_session();
934 let result = mode
935 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
936 .unwrap();
937 apply(&mut session, result);
938 let result = mode
939 .on_message(
940 &session,
941 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
942 )
943 .unwrap();
944 apply(&mut session, result);
945 let result = mode
946 .on_message(
947 &session,
948 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
949 )
950 .unwrap();
951 apply(&mut session, result);
952 let err = mode
953 .on_message(
954 &session,
955 &env("planner", "Commitment", commitment_payload()),
956 )
957 .unwrap_err();
958 assert_eq!(err.to_string(), "InvalidPayload");
959 }
960
961 #[test]
962 fn non_initiator_commitment_rejected() {
963 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
964 let mut session = base_session();
965 let result = mode
966 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
967 .unwrap();
968 apply(&mut session, result);
969 let result = mode
970 .on_message(
971 &session,
972 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
973 )
974 .unwrap();
975 apply(&mut session, result);
976 let result = mode
977 .on_message(
978 &session,
979 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
980 )
981 .unwrap();
982 apply(&mut session, result);
983 let result = mode
984 .on_message(
985 &session,
986 &env("worker", "TaskComplete", make_task_complete("t1", "worker")),
987 )
988 .unwrap();
989 apply(&mut session, result);
990 let commit_env = env("worker", "Commitment", commitment_payload());
991 let err = mode.authorize_sender(&session, &commit_env).unwrap_err();
992 assert_eq!(err.to_string(), "Forbidden");
993 }
994
995 #[test]
998 fn full_task_lifecycle() {
999 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1000 let mut session = base_session();
1001 let result = mode
1002 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1003 .unwrap();
1004 apply(&mut session, result);
1005 let result = mode
1006 .on_message(
1007 &session,
1008 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
1009 )
1010 .unwrap();
1011 apply(&mut session, result);
1012 let result = mode
1013 .on_message(
1014 &session,
1015 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
1016 )
1017 .unwrap();
1018 apply(&mut session, result);
1019 let result = mode
1020 .on_message(
1021 &session,
1022 &env("worker", "TaskUpdate", make_task_update("t1")),
1023 )
1024 .unwrap();
1025 apply(&mut session, result);
1026 let result = mode
1027 .on_message(
1028 &session,
1029 &env("worker", "TaskComplete", make_task_complete("t1", "worker")),
1030 )
1031 .unwrap();
1032 apply(&mut session, result);
1033 let result = mode
1034 .on_message(
1035 &session,
1036 &env("planner", "Commitment", commitment_payload()),
1037 )
1038 .unwrap();
1039 assert!(matches!(result, ModeResponse::PersistAndResolve { .. }));
1040 }
1041
1042 #[test]
1045 fn negative_outcome_task_failed() {
1046 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1047 let mut session = base_session();
1048 let result = mode
1049 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1050 .unwrap();
1051 apply(&mut session, result);
1052 let result = mode
1054 .on_message(
1055 &session,
1056 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
1057 )
1058 .unwrap();
1059 apply(&mut session, result);
1060 let result = mode
1062 .on_message(
1063 &session,
1064 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
1065 )
1066 .unwrap();
1067 apply(&mut session, result);
1068 let result = mode
1070 .on_message(
1071 &session,
1072 &env("worker", "TaskFail", make_task_fail("t1", "worker")),
1073 )
1074 .unwrap();
1075 apply(&mut session, result);
1076 let negative_commitment = CommitmentPayload {
1078 commitment_id: "c1".into(),
1079 action: "task.failed".into(),
1080 authority_scope: "ops".into(),
1081 reason: "task failed".into(),
1082 mode_version: "1.0.0".into(),
1083 policy_version: "policy".into(),
1084 configuration_version: "config".into(),
1085 outcome_positive: false,
1086 supersedes: None,
1087 }
1088 .encode_to_vec();
1089 let result = mode
1090 .on_message(&session, &env("planner", "Commitment", negative_commitment))
1091 .unwrap();
1092 assert!(matches!(result, ModeResponse::PersistAndResolve { .. }));
1093 }
1094
1095 #[test]
1098 fn duplicate_rejection_from_same_sender_rejected() {
1099 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1100 let mut session = base_session();
1101 session.participants = vec!["planner".into(), "w1".into(), "w2".into()];
1102 let result = mode
1103 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1104 .unwrap();
1105 apply(&mut session, result);
1106 let result = mode
1107 .on_message(
1108 &session,
1109 &env("planner", "TaskRequest", make_task_request("t1", "")),
1110 )
1111 .unwrap();
1112 apply(&mut session, result);
1113 let result = mode
1114 .on_message(
1115 &session,
1116 &env("w1", "TaskReject", make_task_reject("t1", "w1")),
1117 )
1118 .unwrap();
1119 apply(&mut session, result);
1120 let err = mode
1121 .on_message(
1122 &session,
1123 &env("w1", "TaskReject", make_task_reject("t1", "w1")),
1124 )
1125 .unwrap_err();
1126 assert_eq!(err.to_string(), "InvalidPayload");
1127 }
1128
1129 #[test]
1130 fn different_senders_can_both_reject() {
1131 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1132 let mut session = base_session();
1133 session.participants = vec!["planner".into(), "w1".into(), "w2".into()];
1134 let result = mode
1135 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1136 .unwrap();
1137 apply(&mut session, result);
1138 let result = mode
1139 .on_message(
1140 &session,
1141 &env("planner", "TaskRequest", make_task_request("t1", "")),
1142 )
1143 .unwrap();
1144 apply(&mut session, result);
1145 let result = mode
1146 .on_message(
1147 &session,
1148 &env("w1", "TaskReject", make_task_reject("t1", "w1")),
1149 )
1150 .unwrap();
1151 apply(&mut session, result);
1152 let result = mode
1153 .on_message(
1154 &session,
1155 &env("w2", "TaskReject", make_task_reject("t1", "w2")),
1156 )
1157 .unwrap();
1158 match result {
1159 ModeResponse::PersistState(data) => {
1160 let state: TaskState = serde_json::from_slice(&data).unwrap();
1161 assert_eq!(state.rejections.len(), 2);
1162 }
1163 _ => panic!("Expected PersistState"),
1164 }
1165 }
1166
1167 #[test]
1170 fn commitment_version_mismatch_rejected() {
1171 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1172 let mut session = base_session();
1173 let result = mode
1174 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1175 .unwrap();
1176 apply(&mut session, result);
1177 let result = mode
1178 .on_message(
1179 &session,
1180 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
1181 )
1182 .unwrap();
1183 apply(&mut session, result);
1184 let result = mode
1185 .on_message(
1186 &session,
1187 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
1188 )
1189 .unwrap();
1190 apply(&mut session, result);
1191 let result = mode
1192 .on_message(
1193 &session,
1194 &env("worker", "TaskComplete", make_task_complete("t1", "worker")),
1195 )
1196 .unwrap();
1197 apply(&mut session, result);
1198 let bad_commitment = CommitmentPayload {
1199 commitment_id: "c1".into(),
1200 action: "task.completed".into(),
1201 authority_scope: "ops".into(),
1202 reason: "done".into(),
1203 mode_version: "wrong".into(),
1204 policy_version: "policy".into(),
1205 configuration_version: "config".into(),
1206 outcome_positive: true,
1207 supersedes: None,
1208 }
1209 .encode_to_vec();
1210 let err = mode
1211 .on_message(&session, &env("planner", "Commitment", bad_commitment))
1212 .unwrap_err();
1213 assert_eq!(err.to_string(), "InvalidPayload");
1214 }
1215
1216 #[test]
1219 fn unknown_message_type_rejected() {
1220 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1221 let mut session = base_session();
1222 let result = mode
1223 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1224 .unwrap();
1225 apply(&mut session, result);
1226 let err = mode
1227 .on_message(&session, &env("worker", "CustomType", vec![]))
1228 .unwrap_err();
1229 assert_eq!(err.to_string(), "InvalidPayload");
1230 }
1231
1232 #[test]
1235 fn policy_allows_commitment_when_output_present() {
1236 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1237 let mut session = base_session();
1238 session.policy_definition = Some(macp_core::policy::PolicyDefinition {
1239 policy_id: "test-strict".into(),
1240 mode: "macp.mode.task.v1".into(),
1241 description: "strict".into(),
1242 rules: serde_json::json!({ "completion": { "require_output": true } }),
1243 schema_version: 1,
1244 });
1245 let result = mode
1246 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1247 .unwrap();
1248 apply(&mut session, result);
1249 let result = mode
1250 .on_message(
1251 &session,
1252 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
1253 )
1254 .unwrap();
1255 apply(&mut session, result);
1256 let result = mode
1257 .on_message(
1258 &session,
1259 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
1260 )
1261 .unwrap();
1262 apply(&mut session, result);
1263 let result = mode
1264 .on_message(
1265 &session,
1266 &env("worker", "TaskComplete", make_task_complete("t1", "worker")),
1267 )
1268 .unwrap();
1269 apply(&mut session, result);
1270 let result = mode
1272 .on_message(
1273 &session,
1274 &env("planner", "Commitment", commitment_payload()),
1275 )
1276 .unwrap();
1277 assert!(matches!(result, ModeResponse::PersistAndResolve { .. }));
1278 }
1279
1280 #[test]
1281 fn policy_with_no_output_requirement_allows_commitment() {
1282 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1283 let mut session = base_session();
1284 session.policy_definition = Some(macp_core::policy::PolicyDefinition {
1285 policy_id: "test-permissive".into(),
1286 mode: "macp.mode.task.v1".into(),
1287 description: "permissive".into(),
1288 rules: serde_json::json!({ "completion": { "require_output": false } }),
1289 schema_version: 1,
1290 });
1291 let result = mode
1292 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1293 .unwrap();
1294 apply(&mut session, result);
1295 let result = mode
1296 .on_message(
1297 &session,
1298 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
1299 )
1300 .unwrap();
1301 apply(&mut session, result);
1302 let result = mode
1303 .on_message(
1304 &session,
1305 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
1306 )
1307 .unwrap();
1308 apply(&mut session, result);
1309 let result = mode
1310 .on_message(
1311 &session,
1312 &env("worker", "TaskComplete", make_task_complete("t1", "worker")),
1313 )
1314 .unwrap();
1315 apply(&mut session, result);
1316 let result = mode
1318 .on_message(
1319 &session,
1320 &env("planner", "Commitment", commitment_payload()),
1321 )
1322 .unwrap();
1323 assert!(matches!(result, ModeResponse::PersistAndResolve { .. }));
1324 }
1325
1326 #[test]
1329 fn competing_task_accept_second_rejected() {
1330 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1331 let mut session = base_session();
1332 session.participants = vec!["planner".into(), "w1".into(), "w2".into()];
1334 let result = mode
1335 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1336 .unwrap();
1337 apply(&mut session, result);
1338 let result = mode
1340 .on_message(
1341 &session,
1342 &env("planner", "TaskRequest", make_task_request("t1", "")),
1343 )
1344 .unwrap();
1345 apply(&mut session, result);
1346 let result = mode
1348 .on_message(
1349 &session,
1350 &env("w1", "TaskAccept", make_task_accept("t1", "w1")),
1351 )
1352 .unwrap();
1353 apply(&mut session, result);
1354 let state: TaskState = serde_json::from_slice(&session.mode_state).unwrap();
1355 assert_eq!(state.active_assignee, Some("w1".into()));
1356 let err = mode
1358 .on_message(
1359 &session,
1360 &env("w2", "TaskAccept", make_task_accept("t1", "w2")),
1361 )
1362 .unwrap_err();
1363 assert_eq!(err.to_string(), "InvalidPayload");
1364 }
1365
1366 #[test]
1369 fn only_active_assignee_can_send_task_update() {
1370 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1371 let mut session = base_session();
1372 session.participants = vec!["planner".into(), "agentA".into(), "agentB".into()];
1373 let result = mode
1374 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1375 .unwrap();
1376 apply(&mut session, result);
1377 let result = mode
1379 .on_message(
1380 &session,
1381 &env("planner", "TaskRequest", make_task_request("t1", "")),
1382 )
1383 .unwrap();
1384 apply(&mut session, result);
1385 let result = mode
1387 .on_message(
1388 &session,
1389 &env("agentA", "TaskAccept", make_task_accept("t1", "agentA")),
1390 )
1391 .unwrap();
1392 apply(&mut session, result);
1393 let err = mode
1395 .on_message(
1396 &session,
1397 &env("agentB", "TaskUpdate", make_task_update("t1")),
1398 )
1399 .unwrap_err();
1400 assert_eq!(err.to_string(), "Forbidden");
1401 }
1402
1403 #[test]
1406 fn active_assignee_can_reject_with_reassignment_policy() {
1407 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1408 let mut session = base_session();
1409 session.policy_definition = Some(macp_core::policy::PolicyDefinition {
1410 policy_id: "test".into(),
1411 mode: "macp.mode.task.v1".into(),
1412 description: "allows reassignment".into(),
1413 rules: serde_json::json!({
1414 "assignment": { "allow_reassignment_on_reject": true },
1415 "completion": { "require_output": false }
1416 }),
1417 schema_version: 1,
1418 });
1419 let result = mode
1420 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1421 .unwrap();
1422 apply(&mut session, result);
1423 let result = mode
1424 .on_message(
1425 &session,
1426 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
1427 )
1428 .unwrap();
1429 apply(&mut session, result);
1430 let result = mode
1431 .on_message(
1432 &session,
1433 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
1434 )
1435 .unwrap();
1436 apply(&mut session, result);
1437 let result = mode
1439 .on_message(
1440 &session,
1441 &env("worker", "TaskReject", make_task_reject("t1", "worker")),
1442 )
1443 .unwrap();
1444 match result {
1445 ModeResponse::PersistState(data) => {
1446 let state: TaskState = serde_json::from_slice(&data).unwrap();
1447 assert!(
1448 state.active_assignee.is_none(),
1449 "should clear active assignee"
1450 );
1451 assert_eq!(state.rejections.len(), 1);
1452 }
1453 _ => panic!("Expected PersistState"),
1454 }
1455 }
1456
1457 #[test]
1458 fn active_assignee_cannot_reject_without_reassignment_policy() {
1459 let mode = TaskMode::new(std::sync::Arc::new(macp_policy::DefaultPolicyEvaluator));
1460 let mut session = base_session();
1461 let result = mode
1463 .on_session_start(&session, &env("planner", "SessionStart", vec![]))
1464 .unwrap();
1465 apply(&mut session, result);
1466 let result = mode
1467 .on_message(
1468 &session,
1469 &env("planner", "TaskRequest", make_task_request("t1", "worker")),
1470 )
1471 .unwrap();
1472 apply(&mut session, result);
1473 let result = mode
1474 .on_message(
1475 &session,
1476 &env("worker", "TaskAccept", make_task_accept("t1", "worker")),
1477 )
1478 .unwrap();
1479 apply(&mut session, result);
1480 let err = mode
1482 .on_message(
1483 &session,
1484 &env("worker", "TaskReject", make_task_reject("t1", "worker")),
1485 )
1486 .unwrap_err();
1487 assert_eq!(err.to_string(), "PolicyDenied");
1488 }
1489}