Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions crates/ironclaw_agent_loop/src/executor.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
44 changes: 36 additions & 8 deletions crates/ironclaw_agent_loop/src/executor/canonical.rs
Original file line number Diff line number Diff line change
@@ -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;

Expand Down Expand Up @@ -305,14 +308,11 @@ impl DefaultExecutorPipeline {
pending_input_ack = resume.pending_input_ack;
pending_input_ack.ack(host).await?;
let completed = self
.capabilities
.process(
.process_resume_capability_calls(
ctx,
CapabilityInput {
state: resume.state,
surface: resume.surface,
calls: vec![resume.call],
},
resume.state,
resume.surface,
vec![resume.call],
)
.await?;

Expand Down Expand Up @@ -458,4 +458,32 @@ impl DefaultExecutorPipeline {
}
}
}

pub(super) async fn process_resume_capability_calls(
&self,
ctx: StageContext<'_>,
state: LoopExecutionState,
surface: VisibleCapabilitySurface,
calls: Vec<CapabilityCallCandidate>,
) -> Result<TurnCompletedStep, AgentLoopExecutorError> {
let capability_input = resume_capability_input(state, surface, calls)?;
self.capabilities.process(ctx, capability_input).await
}
}

pub(super) fn resume_capability_input(
state: LoopExecutionState,
surface: VisibleCapabilitySurface,
calls: Vec<CapabilityCallCandidate>,
) -> Result<CapabilityInput, AgentLoopExecutorError> {
if calls.len() != 1 {
return Err(AgentLoopExecutorError::PlannerContract {
detail: "resume dispatch must contain exactly one capability call",
});
}
Ok(CapabilityInput {
state,
surface,
calls,
})
}
15 changes: 15 additions & 0 deletions crates/ironclaw_agent_loop/src/executor/capabilities.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1104,6 +1104,7 @@ impl CapabilityStage {
ControlFlow<TurnCompletedStep, (LoopExecutionState, Vec<CapabilityCallCandidate>)>,
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);
Expand Down Expand Up @@ -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,
Expand Down
194 changes: 190 additions & 4 deletions crates/ironclaw_agent_loop/src/executor/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,10 +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, 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)]
Expand Down Expand Up @@ -6715,6 +6716,191 @@ 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"
));
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

#[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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
},
Expand Down
1 change: 1 addition & 0 deletions crates/ironclaw_product_workflow/src/auth_continuation.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down
2 changes: 2 additions & 0 deletions crates/ironclaw_reborn_composition/src/factory/auth_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down Expand Up @@ -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(),
Expand Down
1 change: 1 addition & 0 deletions crates/ironclaw_reborn_composition/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9185,6 +9185,7 @@ output_schema_ref = "schemas/write.output.json"
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(),
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
Expand Down
2 changes: 1 addition & 1 deletion crates/ironclaw_turns/src/memory/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
2 changes: 2 additions & 0 deletions crates/ironclaw_turns/src/runner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<CapabilityActivityId>,
pub reason: BlockedReason,
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
},
Expand Down
1 change: 1 addition & 0 deletions crates/ironclaw_turns/tests/per_user_concurrency_cap.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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(),
},
Expand Down
Loading
Loading