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