Skip to main content

macp_modes/mode/
task.rs

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/// The accepted `TaskRequest` as it stands in serialized `mode_state`.
17///
18/// **`#[non_exhaustive]`** (0.8.0, `DECISIONS.md` D7), for the same reason as
19/// the handoff, quorum and proposal records: this is runtime-produced
20/// coordination state, deserialized from accepted envelopes, that grows a
21/// field whenever the mode learns something new. With all-`pub` fields and no
22/// seal each added field is a `constructible_struct_adds_field` major, and
23/// `release-plz.toml`'s `semver_check = true` turns that into a blocked
24/// release PR across all seven lockstep crates. 0.8.0 is already being taken
25/// for the handoff field, so sealing the rest of the class here costs nothing
26/// extra and makes every future field additive.
27///
28/// Fields stay `pub` and readable; only construction by struct literal from
29/// another crate is refused, and nothing outside `macp-modes` constructs one —
30/// the mode mints them from accepted envelopes, which is why no constructor is
31/// offered in their place.
32#[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/// A `TaskReject` as it stands in serialized `mode_state`.
45///
46/// `#[non_exhaustive]` for the reason on [`TaskRecord`].
47#[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/// A `TaskUpdate` as it stands in serialized `mode_state`.
56///
57/// `#[non_exhaustive]` for the reason on [`TaskRecord`].
58#[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/// A `TaskComplete` as it stands in serialized `mode_state`.
70///
71/// `#[non_exhaustive]` for the reason on [`TaskRecord`].
72#[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/// A `TaskFail` as it stands in serialized `mode_state`.
82///
83/// `#[non_exhaustive]` for the reason on [`TaskRecord`].
84#[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/// The task mode's whole serialized `mode_state`.
101///
102/// `#[non_exhaustive]` for the reason on [`TaskRecord`]. `Default` is still
103/// derived and still reachable from other crates (`TaskState::default()`);
104/// `#[non_exhaustive]` refuses only the struct-literal form, including
105/// `TaskState { ..Default::default() }`.
106#[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    /// Construct the mode with an injected governance policy evaluator.
122    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        // RFC-MACP-0009 §2 (orchestrated model): the requester coordinates
168        // work among eligible assignees, so the pool must contain at least
169        // one participant other than the initiator. The initiator's own
170        // authority to emit TaskRequest is role-based (§4) and does NOT
171        // require participant-list membership — an external orchestrator may
172        // drive a task without listing itself (mirrors the role-vs-membership
173        // principle RFC-0007 §2 and RFC-0011 §2 make explicit).
174        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                // RFC-MACP-0012: check allow_reassignment_on_reject if prior rejections exist
228                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                // RFC-MACP-0009 §5.3b: TaskAccept is irrevocable unless policy
260                // permits reassignment. §5.3c: when allow_reassignment_on_reject
261                // is true, the active assignee may send TaskReject to return the
262                // session to the pre-assignment state.
263                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                        // Policy permits: clear assignee, return to pre-assignment state
281                    } else {
282                        // Someone other than the active assignee trying to reject
283                        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 active assignee is rejecting with policy permission, clear assignment
298                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                // Governance policy gate (shared): fail closed, only
377                // an explicit Allow proceeds.
378                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    // --- Session Start ---
513
514    #[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        // The initiator alone leaves no eligible assignee (RFC-0009 §2).
536        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    /// RFC-MACP-0009 §2/§4: TaskRequest authority is role-based (initiator),
555    /// not membership-based — an external orchestrator driving a pool of
556    /// assignees is a legitimate orchestrated session.
557    #[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()]; // planner not included
562        let result = mode
563            .on_session_start(&session, &env("planner", "SessionStart", vec![]))
564            .unwrap();
565        assert!(matches!(result, ModeResponse::PersistState(_)));
566    }
567
568    // --- TaskRequest ---
569
570    #[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    // --- TaskAccept / TaskReject ---
636
637    #[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()); // still unassigned
717            }
718            _ => panic!("Expected PersistState"),
719        }
720    }
721
722    // --- TaskUpdate ---
723
724    #[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    // --- TaskComplete / TaskFail ---
796
797    #[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    // --- Commitment ---
878
879    #[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    // --- Full lifecycle ---
1021
1022    #[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    // --- Negative outcome commitment ---
1068
1069    #[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        // Send TaskRequest
1078        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        // Worker accepts
1086        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        // Worker reports failure
1094        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        // Commit with negative outcome
1102        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    // --- Duplicate rejection ---
1121
1122    #[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    // --- Commitment version mismatch ---
1193
1194    #[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    // --- Unknown message type ---
1242
1243    #[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    // --- Policy ---
1258
1259    #[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        // Commitment should succeed: output is required and task completion has output
1296        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        // Commitment should succeed: require_output is false
1342        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    // --- Competing TaskAccept ---
1352
1353    #[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        // Three participants: planner (initiator), w1, w2
1358        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        // Open-assignee task (no specific requested_assignee)
1364        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        // First accept from w1 succeeds
1372        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        // Second accept from w2 is rejected (active_assignee already set)
1382        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    // --- Only active assignee can send TaskUpdate ---
1392
1393    #[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        // Open-assignee task
1403        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        // agentA accepts the task
1411        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        // agentB (non-assignee) attempts to send TaskUpdate — expect Forbidden
1419        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    // --- TaskReject with reassignment policy (RFC-MACP-0009 §5.3c) ---
1429
1430    #[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        // Active assignee rejects with policy permission — returns to pre-assignment
1463        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        // No policy = no reassignment
1487        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        // Active assignee rejects without policy permission — denied
1506        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}