From 4f37ae041cec36e04c57bf4999a16af7f1a146b2 Mon Sep 17 00:00:00 2001 From: Coffee Date: Sat, 27 Jun 2026 01:09:27 +0800 Subject: [PATCH 1/2] fix(reborn): harden activity identity invariants --- crates/ironclaw_agent_loop/src/executor.rs | 2 + .../src/executor/canonical.rs | 36 +++-- .../src/executor/capabilities.rs | 15 ++ .../ironclaw_agent_loop/src/executor/tests.rs | 142 +++++++++++++++++- .../src/factory/auth_tests.rs | 2 + .../src/runtime.rs | 1 + .../src/runtime/tests/auth_interaction.rs | 1 + crates/ironclaw_turns/src/memory/mod.rs | 2 +- crates/ironclaw_turns/src/runner.rs | 2 + .../tests/per_inbound_type_concurrency_cap.rs | 1 + .../tests/per_user_concurrency_cap.rs | 1 + .../tests/turn_coordinator_contract.rs | 79 +++++++++- 12 files changed, 267 insertions(+), 17 deletions(-) diff --git a/crates/ironclaw_agent_loop/src/executor.rs b/crates/ironclaw_agent_loop/src/executor.rs index 6fdf9140a69..1d3959a992b 100644 --- a/crates/ironclaw_agent_loop/src/executor.rs +++ b/crates/ironclaw_agent_loop/src/executor.rs @@ -22,6 +22,8 @@ mod turn_stop; use assistant_reply::{AssistantReplyInput, AssistantReplyStage}; use budget::{BudgetInput, BudgetStage, BudgetStep}; +#[cfg(test)] +use canonical::resume_capability_input; use capabilities::{CapabilityInput, CapabilityStage}; use capability_helpers::{ CapabilitySurfaceIndex, append_capability_error_ref, append_capability_result_ref, diff --git a/crates/ironclaw_agent_loop/src/executor/canonical.rs b/crates/ironclaw_agent_loop/src/executor/canonical.rs index 77bb494cc9b..2efc4d5c0b0 100644 --- a/crates/ironclaw_agent_loop/src/executor/canonical.rs +++ b/crates/ironclaw_agent_loop/src/executor/canonical.rs @@ -1,6 +1,9 @@ use ironclaw_turns::{ LoopExit, - run_profile::{AgentLoopDriverHost, LoopDriverNoteKind, LoopProgressEvent, ParentLoopOutput}, + run_profile::{ + AgentLoopDriverHost, CapabilityCallCandidate, LoopDriverNoteKind, LoopProgressEvent, + ParentLoopOutput, VisibleCapabilitySurface, + }, }; use tracing::debug; @@ -304,17 +307,9 @@ impl DefaultExecutorPipeline { let resume = *resume; pending_input_ack = resume.pending_input_ack; pending_input_ack.ack(host).await?; - let completed = self - .capabilities - .process( - ctx, - CapabilityInput { - state: resume.state, - surface: resume.surface, - calls: vec![resume.call], - }, - ) - .await?; + let capability_input = + resume_capability_input(resume.state, resume.surface, vec![resume.call])?; + let completed = self.capabilities.process(ctx, capability_input).await?; let completed = self.post_capability.process(ctx, completed).await?; @@ -459,3 +454,20 @@ impl DefaultExecutorPipeline { } } } + +pub(super) fn resume_capability_input( + state: LoopExecutionState, + surface: VisibleCapabilitySurface, + calls: Vec, +) -> Result { + if calls.len() != 1 { + return Err(AgentLoopExecutorError::PlannerContract { + detail: "resume dispatch must contain exactly one capability call", + }); + } + Ok(CapabilityInput { + state, + surface, + calls, + }) +} diff --git a/crates/ironclaw_agent_loop/src/executor/capabilities.rs b/crates/ironclaw_agent_loop/src/executor/capabilities.rs index 02ecf797589..23c8ccfcc3b 100644 --- a/crates/ironclaw_agent_loop/src/executor/capabilities.rs +++ b/crates/ironclaw_agent_loop/src/executor/capabilities.rs @@ -1104,6 +1104,7 @@ impl CapabilityStage { ControlFlow)>, AgentLoopExecutorError, > { + ensure_unique_activity_ids_for_denied_resume(&visible_calls)?; let (denied_calls, remaining_calls): (Vec<_>, Vec<_>) = visible_calls .into_iter() .partition(|call| call.activity_id == denied_activity_id); @@ -1162,6 +1163,20 @@ impl CapabilityStage { } } +fn ensure_unique_activity_ids_for_denied_resume( + calls: &[CapabilityCallCandidate], +) -> Result<(), AgentLoopExecutorError> { + let mut activity_ids = HashSet::new(); + for call in calls { + if !activity_ids.insert(call.activity_id) { + return Err(AgentLoopExecutorError::PlannerContract { + detail: "denied resume batch contains duplicate activity id", + }); + } + } + Ok(()) +} + fn clear_matching_pending_approval_resume( state: &mut LoopExecutionState, call: &CapabilityCallCandidate, diff --git a/crates/ironclaw_agent_loop/src/executor/tests.rs b/crates/ironclaw_agent_loop/src/executor/tests.rs index 11be7d18bdd..bb04cce2ebc 100644 --- a/crates/ironclaw_agent_loop/src/executor/tests.rs +++ b/crates/ironclaw_agent_loop/src/executor/tests.rs @@ -35,7 +35,8 @@ use super::{ CapabilityStage, DrainInput, ExecutorStage, ExitInput, ExitStage, GateInput, GateStage, HostStage, InputStage, InputStep, PendingInputAck, PromptInput, PromptStage, PromptStep, StageContext, StopInput, StopStage, StopStep, TurnCompletedStep, UserFacingInputDrainMode, - consume_drainable_inputs, sanitize_result_ref_suffix, synthetic_provider_error_result_ref, + consume_drainable_inputs, resume_capability_input, sanitize_result_ref_suffix, + synthetic_provider_error_result_ref, }; #[allow(dead_code)] @@ -6668,6 +6669,145 @@ async fn capability_stage_denied_auth_resume_only_fails_matching_activity_when_c assert!(final_state.result_refs.contains(&y_result_ref)); } +#[tokio::test] +async fn capability_stage_denied_auth_resume_rejects_duplicate_activity_id_before_side_effects() { + let host = MockHost::new(Vec::new()); + let family = crate::families::default(); + let ctx = StageContext { + planner: family.planner(), + host: &host, + }; + let denied_activity_id = CapabilityActivityId::new(); + let mut state = LoopExecutionState::initial_for_run(host.run_context()); + state.pending_auth_resume = Some(PendingAuthResume { + gate_ref: LoopGateRef::new("gate:auth-deny-duplicate-activity").expect("valid"), + capability_id: capability_id(), + surface_version: surface_version(), + input_ref: CapabilityInputRef::new("input:duplicate-activity-resume").expect("valid"), + effective_capability_ids: vec![capability_id()], + provider_replay: None, + resume_token: None, + activity_id: denied_activity_id, + prior_approval: None, + replay: None, + disposition: Some(ironclaw_turns::GateResumeDisposition::Denied), + }); + + let calls = vec![ + CapabilityCallCandidate { + activity_id: denied_activity_id, + surface_version: surface_version(), + capability_id: capability_id(), + input_ref: CapabilityInputRef::new("input:duplicate-activity-a").expect("valid"), + effective_capability_ids: vec![capability_id()], + provider_replay: Some(ProviderToolCallReplay { + provider_id: "test-provider".to_string(), + provider_model_id: "test-model".to_string(), + provider_turn_id: "turn_1".to_string(), + provider_call_id: "call_duplicate_a".to_string(), + provider_tool_name: ProviderToolName::new("demo__echo") + .expect("provider tool name"), + arguments: serde_json::json!({"message": "a"}), + response_reasoning: None, + reasoning: None, + signature: None, + }), + }, + CapabilityCallCandidate { + activity_id: denied_activity_id, + surface_version: surface_version(), + capability_id: capability_id(), + input_ref: CapabilityInputRef::new("input:duplicate-activity-b").expect("valid"), + effective_capability_ids: vec![capability_id()], + provider_replay: Some(ProviderToolCallReplay { + provider_id: "test-provider".to_string(), + provider_model_id: "test-model".to_string(), + provider_turn_id: "turn_1".to_string(), + provider_call_id: "call_duplicate_b".to_string(), + provider_tool_name: ProviderToolName::new("demo__echo") + .expect("provider tool name"), + arguments: serde_json::json!({"message": "b"}), + response_reasoning: None, + reasoning: None, + signature: None, + }), + }, + ]; + + let result = CapabilityStage + .process( + ctx, + CapabilityInput { + state, + surface: ironclaw_turns::run_profile::LoopCapabilityPort::visible_capabilities( + &host, + VisibleCapabilityRequest, + ) + .await + .expect("visible surface"), + calls, + }, + ) + .await; + + assert!(matches!( + result, + Err(AgentLoopExecutorError::PlannerContract { detail }) + if detail == "denied resume batch contains duplicate activity id" + )); + assert!( + host.batch_invocations().is_empty(), + "duplicate denied activity ids must not reach capability dispatch" + ); + assert!( + host.appended_result_refs().is_empty(), + "duplicate denied activity ids must not append model-visible failures" + ); + assert!( + !host.progress_events().iter().any(|event| matches!( + event, + LoopProgressEvent::CapabilityActivityFailed { + activity_id, + reason_kind: CapabilityFailureKind::GateDeclined, + .. + } if *activity_id == denied_activity_id + )), + "ambiguous denied activity ids must not emit a terminal activity milestone" + ); +} + +#[test] +fn resume_capability_input_rejects_batched_resume_calls() { + let host = MockHost::new(Vec::new()); + let state = LoopExecutionState::initial_for_run(host.run_context()); + let call = CapabilityCallCandidate { + activity_id: CapabilityActivityId::new(), + surface_version: surface_version(), + capability_id: capability_id(), + input_ref: CapabilityInputRef::new("input:resume-single").expect("valid"), + effective_capability_ids: vec![capability_id()], + provider_replay: None, + }; + let mut second_call = call.clone(); + second_call.activity_id = CapabilityActivityId::new(); + second_call.input_ref = CapabilityInputRef::new("input:resume-extra").expect("valid"); + + let result = resume_capability_input( + state, + ironclaw_turns::run_profile::VisibleCapabilitySurface { + version: surface_version(), + descriptors: Vec::new(), + }, + vec![call, second_call], + ); + + assert!(matches!( + result, + Err(AgentLoopExecutorError::PlannerContract { detail }) + if detail == "resume dispatch must contain exactly one capability call" + )); +} + /// Regression test: partition + sizing invariant with 1 denied + 2 remaining calls. /// /// When the denied auth-resume batch contains one denied call (X) and TWO diff --git a/crates/ironclaw_reborn_composition/src/factory/auth_tests.rs b/crates/ironclaw_reborn_composition/src/factory/auth_tests.rs index 06c75fac8c3..a3001d8254c 100644 --- a/crates/ironclaw_reborn_composition/src/factory/auth_tests.rs +++ b/crates/ironclaw_reborn_composition/src/factory/auth_tests.rs @@ -223,6 +223,7 @@ async fn local_dev_oauth_turn_gate_callback_resumes_default_turn_coordinator() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: LoopCheckpointStateRef::new("checkpoint:auth-callback").unwrap(), + blocked_activity_id: None, reason: BlockedReason::Auth { gate_ref: gate_ref.clone(), credential_requirements: Vec::new(), @@ -678,6 +679,7 @@ async fn submit_and_block_auth_run( lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: LoopCheckpointStateRef::new("checkpoint:auth-callback-2").unwrap(), + blocked_activity_id: None, reason: BlockedReason::Auth { gate_ref: ironclaw_turns::GateRef::new(gate_ref).unwrap(), credential_requirements: Vec::new(), diff --git a/crates/ironclaw_reborn_composition/src/runtime.rs b/crates/ironclaw_reborn_composition/src/runtime.rs index 2b00974a1a6..015a21d893f 100644 --- a/crates/ironclaw_reborn_composition/src/runtime.rs +++ b/crates/ironclaw_reborn_composition/src/runtime.rs @@ -9184,6 +9184,7 @@ command = "test" lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: LoopCheckpointStateRef::new("checkpoint:audit").unwrap(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, diff --git a/crates/ironclaw_reborn_composition/src/runtime/tests/auth_interaction.rs b/crates/ironclaw_reborn_composition/src/runtime/tests/auth_interaction.rs index 0e86f8d482b..b07160bc694 100644 --- a/crates/ironclaw_reborn_composition/src/runtime/tests/auth_interaction.rs +++ b/crates/ironclaw_reborn_composition/src/runtime/tests/auth_interaction.rs @@ -236,6 +236,7 @@ async fn submit_and_block_auth_run( checkpoint_id: TurnCheckpointId::new(), state_ref: LoopCheckpointStateRef::new("checkpoint:runtime-auth-read-model") .expect("checkpoint ref"), + blocked_activity_id: None, reason: BlockedReason::Auth { gate_ref: gate_ref.clone(), credential_requirements: Vec::new(), diff --git a/crates/ironclaw_turns/src/memory/mod.rs b/crates/ironclaw_turns/src/memory/mod.rs index 6280d3453c9..377e3856df6 100644 --- a/crates/ironclaw_turns/src/memory/mod.rs +++ b/crates/ironclaw_turns/src/memory/mod.rs @@ -1454,7 +1454,7 @@ impl TurnRunTransitionPort for InMemoryTurnStateStore { inner.apply_status_transition(transition, &record); record.checkpoint_id = Some(request.checkpoint_id); record.gate_ref = Some(request.reason.gate_ref().clone()); - record.blocked_activity_id = None; + record.blocked_activity_id = request.blocked_activity_id; record.credential_requirements = request.reason.credential_requirements().to_vec(); record.runner_id = None; record.lease_token = None; diff --git a/crates/ironclaw_turns/src/runner.rs b/crates/ironclaw_turns/src/runner.rs index 4b5806a085e..99b82466845 100644 --- a/crates/ironclaw_turns/src/runner.rs +++ b/crates/ironclaw_turns/src/runner.rs @@ -57,6 +57,8 @@ pub struct BlockRunRequest { pub lease_token: TurnLeaseToken, pub checkpoint_id: TurnCheckpointId, pub state_ref: LoopCheckpointStateRef, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub blocked_activity_id: Option, pub reason: BlockedReason, } diff --git a/crates/ironclaw_turns/tests/per_inbound_type_concurrency_cap.rs b/crates/ironclaw_turns/tests/per_inbound_type_concurrency_cap.rs index bdf21b631d2..8eb428f0c9f 100644 --- a/crates/ironclaw_turns/tests/per_inbound_type_concurrency_cap.rs +++ b/crates/ironclaw_turns/tests/per_inbound_type_concurrency_cap.rs @@ -271,6 +271,7 @@ async fn trigger_counter_decrements_on_block_and_resets_on_resume() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate.clone(), }, diff --git a/crates/ironclaw_turns/tests/per_user_concurrency_cap.rs b/crates/ironclaw_turns/tests/per_user_concurrency_cap.rs index 1fa9eb75c3d..a3ed79394e0 100644 --- a/crates/ironclaw_turns/tests/per_user_concurrency_cap.rs +++ b/crates/ironclaw_turns/tests/per_user_concurrency_cap.rs @@ -202,6 +202,7 @@ async fn running_counter_decrements_on_block_and_resets_on_resume() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate.clone(), }, diff --git a/crates/ironclaw_turns/tests/turn_coordinator_contract.rs b/crates/ironclaw_turns/tests/turn_coordinator_contract.rs index 18baa9a6afc..221d097e3cc 100644 --- a/crates/ironclaw_turns/tests/turn_coordinator_contract.rs +++ b/crates/ironclaw_turns/tests/turn_coordinator_contract.rs @@ -12,9 +12,9 @@ use chrono::{DateTime, Duration as ChronoDuration, TimeZone, Utc}; use ironclaw_host_api::{AgentId, ProjectId, TenantId, ThreadId, UserId}; use ironclaw_turns::{ AcceptedMessageRef, AdmissionRejection, AdmissionRejectionReason, AllowAllTurnAdmissionPolicy, - BlockedReason, CancelRunRequest, CancelRunResponse, DefaultTurnCoordinator, - DefaultTurnLifecycleEventBus, GateRef, GetRunStateRequest, IdempotencyKey, - InMemoryRunProfileResolver, InMemoryTurnEventSink, InMemoryTurnStateStore, + BlockedReason, CancelRunRequest, CancelRunResponse, CapabilityActivityId, + DefaultTurnCoordinator, DefaultTurnLifecycleEventBus, GateRef, GetRunStateRequest, + IdempotencyKey, InMemoryRunProfileResolver, InMemoryTurnEventSink, InMemoryTurnStateStore, InMemoryTurnStateStoreLimits, LifecyclePublicationErrorPort, LifecyclePublishingTurnStateStore, LoopBlockedKind, LoopCheckpointStateRef, LoopExitMapping, LoopGateRef, LoopResultRef, ProductTurnContext, ReplyTargetBindingRef, ResolvedRunProfile, ResumeTurnRequest, @@ -1770,6 +1770,61 @@ async fn lifecycle_publishing_store_publishes_failed_run_event() { })); } +#[tokio::test] +async fn lifecycle_publishing_store_publishes_blocked_event_activity_id_from_block_run() { + let raw_store = Arc::new(InMemoryTurnStateStore::default()); + let sink = Arc::new(InMemoryTurnEventSink::default()); + let transition_port = lifecycle_publishing_store(raw_store, None, Some(sink.clone())); + let coordinator = DefaultTurnCoordinator::new(transition_port.clone()); + let run_id = accepted_run_id( + &coordinator + .submit_turn(submit_request( + "thread-lifecycle-blocked-activity", + "idem-lifecycle-blocked-activity", + )) + .await + .unwrap(), + ); + let runner_id = TurnRunnerId::new(); + let lease_token = TurnLeaseToken::new(); + transition_port + .claim_next_run(ClaimRunRequest { + runner_id, + lease_token, + scope_filter: Some(scope("thread-lifecycle-blocked-activity")), + }) + .await + .unwrap() + .unwrap(); + + let gate_ref = GateRef::new("gate:lifecycle-blocked-activity").unwrap(); + let blocked_activity_id = CapabilityActivityId::new(); + let blocked_state = transition_port + .block_run(BlockRunRequest { + run_id, + runner_id, + lease_token, + checkpoint_id: TurnCheckpointId::new(), + state_ref: block_state_ref(), + blocked_activity_id: Some(blocked_activity_id), + reason: BlockedReason::Approval { + gate_ref: gate_ref.clone(), + }, + }) + .await + .unwrap(); + + assert_eq!(blocked_state.blocked_activity_id, Some(blocked_activity_id)); + let blocked_event = sink + .events() + .into_iter() + .find(|event| event.run_id == run_id && event.kind == TurnEventKind::Blocked) + .expect("blocked lifecycle event"); + let blocked_gate = blocked_event.blocked_gate.expect("blocked gate metadata"); + assert_eq!(blocked_gate.gate_ref, gate_ref); + assert_eq!(blocked_gate.activity_id, Some(blocked_activity_id)); +} + #[tokio::test] async fn lifecycle_publishing_store_publishes_record_runner_failure_as_failed_event() { let raw_store = Arc::new(InMemoryTurnStateStore::default()); @@ -1847,6 +1902,7 @@ async fn turn_lifecycle_projection_replays_submit_block_resume_complete_without_ lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -2478,6 +2534,7 @@ async fn resume_turn_wakes_runner_for_same_run_after_requeue() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -2605,6 +2662,7 @@ async fn resume_turn_ignores_wake_notification_panic_after_requeue() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -3806,6 +3864,7 @@ async fn blocked_resume_then_recovery_failure_releases_admission_reservation() { run_id, runner_id, lease_token, + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -3957,6 +4016,7 @@ async fn runner_claim_and_block_update_persistent_run_lock_and_checkpoint_record lease_token, checkpoint_id, state_ref: state_ref.clone(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -4040,6 +4100,7 @@ async fn block_run_with_reason_yields_expected_blocked_gate( lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref, + blocked_activity_id: None, reason: reason_builder(gate_ref.clone()), }) .await @@ -4112,6 +4173,7 @@ async fn resume_updates_persisted_run_binding_refs_and_replay_envelope() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -4459,6 +4521,7 @@ async fn idempotency_persistence_snapshot_retains_each_operation_kind_capacity() lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -4558,6 +4621,7 @@ async fn idempotency_replay_helpers_require_matching_operation_kind() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -5351,6 +5415,7 @@ async fn blocked_run_persists_checkpoint_and_keeps_same_thread_lock_until_resume lease_token, checkpoint_id, state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -5417,6 +5482,7 @@ async fn resume_turn_rejects_unexpected_blocked_status_without_requeueing_run() lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -5485,6 +5551,7 @@ async fn resume_turn_from_foreign_actor_is_denied_without_requeueing_run() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -5558,6 +5625,7 @@ async fn cancel_run_from_foreign_actor_is_denied_without_mutating_run() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: gate_ref.clone(), }, @@ -5628,6 +5696,7 @@ async fn resume_turn_with_wrong_gate_resolution_ref_is_invalid_request() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: GateRef::new("approval-gate").unwrap(), }, @@ -5826,6 +5895,7 @@ async fn cancelled_running_run_cannot_be_reopened_as_blocked() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Approval { gate_ref: GateRef::new("approval-gate").unwrap(), }, @@ -6057,6 +6127,7 @@ async fn any_blocked_gate_resume_does_not_resume_dependent_run_gate() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::AwaitDependentRun { gate_ref: gate_ref.clone(), }, @@ -7287,6 +7358,7 @@ async fn resume_turn_resume_disposition_is_persisted_and_visible_on_claim() { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Auth { gate_ref: gate_ref.clone(), credential_requirements: Vec::new(), @@ -7353,6 +7425,7 @@ async fn resume_turn_resume_disposition_is_persisted_and_visible_on_claim() { lease_token: claimed.lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: block_state_ref(), + blocked_activity_id: None, reason: BlockedReason::Auth { gate_ref: gate_ref2.clone(), credential_requirements: Vec::new(), From 5f19625d4929937295397ba23c299251ac6553f9 Mon Sep 17 00:00:00 2001 From: Coffee Date: Sat, 27 Jun 2026 11:17:37 +0000 Subject: [PATCH 2/2] test(reborn): cover batched resume dispatch guard --- .../src/executor/canonical.rs | 22 +++++++- .../ironclaw_agent_loop/src/executor/tests.rs | 56 +++++++++++++++++-- .../tests/turn_event_publisher_contract.rs | 1 + .../src/auth_continuation.rs | 1 + 4 files changed, 72 insertions(+), 8 deletions(-) diff --git a/crates/ironclaw_agent_loop/src/executor/canonical.rs b/crates/ironclaw_agent_loop/src/executor/canonical.rs index 2efc4d5c0b0..61bbc414953 100644 --- a/crates/ironclaw_agent_loop/src/executor/canonical.rs +++ b/crates/ironclaw_agent_loop/src/executor/canonical.rs @@ -307,9 +307,14 @@ impl DefaultExecutorPipeline { let resume = *resume; pending_input_ack = resume.pending_input_ack; pending_input_ack.ack(host).await?; - let capability_input = - resume_capability_input(resume.state, resume.surface, vec![resume.call])?; - let completed = self.capabilities.process(ctx, capability_input).await?; + let completed = self + .process_resume_capability_calls( + ctx, + resume.state, + resume.surface, + vec![resume.call], + ) + .await?; let completed = self.post_capability.process(ctx, completed).await?; @@ -453,6 +458,17 @@ impl DefaultExecutorPipeline { } } } + + pub(super) async fn process_resume_capability_calls( + &self, + ctx: StageContext<'_>, + state: LoopExecutionState, + surface: VisibleCapabilitySurface, + calls: Vec, + ) -> Result { + let capability_input = resume_capability_input(state, surface, calls)?; + self.capabilities.process(ctx, capability_input).await + } } pub(super) fn resume_capability_input( diff --git a/crates/ironclaw_agent_loop/src/executor/tests.rs b/crates/ironclaw_agent_loop/src/executor/tests.rs index bb04cce2ebc..edc2e4fc5e9 100644 --- a/crates/ironclaw_agent_loop/src/executor/tests.rs +++ b/crates/ironclaw_agent_loop/src/executor/tests.rs @@ -32,11 +32,11 @@ use crate::test_support::compaction::{ use super::{ AgentLoopExecutor, AgentLoopExecutorError, AssistantReplyInput, AssistantReplyStage, BatchStep, BudgetInput, BudgetStage, BudgetStep, CanonicalAgentLoopExecutor, CapabilityInput, - CapabilityStage, DrainInput, ExecutorStage, ExitInput, ExitStage, GateInput, GateStage, - HostStage, InputStage, InputStep, PendingInputAck, PromptInput, PromptStage, PromptStep, - StageContext, StopInput, StopStage, StopStep, TurnCompletedStep, UserFacingInputDrainMode, - consume_drainable_inputs, resume_capability_input, sanitize_result_ref_suffix, - synthetic_provider_error_result_ref, + CapabilityStage, DefaultExecutorPipeline, DrainInput, ExecutorStage, ExitInput, ExitStage, + GateInput, GateStage, HostStage, InputStage, InputStep, PendingInputAck, PromptInput, + PromptStage, PromptStep, StageContext, StopInput, StopStage, StopStep, TurnCompletedStep, + UserFacingInputDrainMode, consume_drainable_inputs, resume_capability_input, + sanitize_result_ref_suffix, synthetic_provider_error_result_ref, }; #[allow(dead_code)] @@ -6808,6 +6808,52 @@ fn resume_capability_input_rejects_batched_resume_calls() { )); } +#[tokio::test] +async fn executor_resume_dispatch_rejects_batched_calls_before_host_invocation() { + let host = MockHost::new(Vec::new()); + let family = crate::families::default(); + let ctx = StageContext { + planner: family.planner(), + host: &host, + }; + let state = LoopExecutionState::initial_for_run(host.run_context()); + let call = CapabilityCallCandidate { + activity_id: CapabilityActivityId::new(), + surface_version: surface_version(), + capability_id: capability_id(), + input_ref: CapabilityInputRef::new("input:resume-dispatch-single").expect("valid"), + effective_capability_ids: vec![capability_id()], + provider_replay: None, + }; + let mut second_call = call.clone(); + second_call.activity_id = CapabilityActivityId::new(); + second_call.input_ref = CapabilityInputRef::new("input:resume-dispatch-extra").expect("valid"); + let surface = ironclaw_turns::run_profile::LoopCapabilityPort::visible_capabilities( + &host, + VisibleCapabilityRequest, + ) + .await + .expect("visible surface"); + + let result = DefaultExecutorPipeline::default() + .process_resume_capability_calls(ctx, state, surface, vec![call, second_call]) + .await; + + assert!(matches!( + result, + Err(AgentLoopExecutorError::PlannerContract { detail }) + if detail == "resume dispatch must contain exactly one capability call" + )); + assert!( + host.batch_invocations().is_empty(), + "batched resume dispatch must not reach capability batch invocation" + ); + assert!( + host.single_invocations().is_empty(), + "batched resume dispatch must not reach single capability invocation" + ); +} + /// Regression test: partition + sizing invariant with 1 denied + 2 remaining calls. /// /// When the denied auth-resume batch contains one denied call (X) and TWO diff --git a/crates/ironclaw_loop_support/tests/turn_event_publisher_contract.rs b/crates/ironclaw_loop_support/tests/turn_event_publisher_contract.rs index 62634960ed0..21d17ddc3f5 100644 --- a/crates/ironclaw_loop_support/tests/turn_event_publisher_contract.rs +++ b/crates/ironclaw_loop_support/tests/turn_event_publisher_contract.rs @@ -89,6 +89,7 @@ async fn event_publishing_transition_port_publishes_blocked_and_terminal_events( lease_token: blocked_lease, checkpoint_id: TurnCheckpointId::new(), state_ref: LoopCheckpointStateRef::new("checkpoint:block-state").unwrap(), + blocked_activity_id: None, reason: BlockedReason::AwaitDependentRun { gate_ref: GateRef::new("gate-dependent-run").unwrap(), }, diff --git a/crates/ironclaw_product_workflow/src/auth_continuation.rs b/crates/ironclaw_product_workflow/src/auth_continuation.rs index 0368b0dc49b..cc48ef276e3 100644 --- a/crates/ironclaw_product_workflow/src/auth_continuation.rs +++ b/crates/ironclaw_product_workflow/src/auth_continuation.rs @@ -703,6 +703,7 @@ mod tests { lease_token, checkpoint_id: TurnCheckpointId::new(), state_ref: LoopCheckpointStateRef::new("checkpoint:auth-real").unwrap(), + blocked_activity_id: None, reason: BlockedReason::Auth { gate_ref: GateRef::new("gate:auth-real").unwrap(), credential_requirements: Vec::new(),