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