Skip to content
Closed
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: 1 addition & 1 deletion FEATURE_PARITY.md
Original file line number Diff line number Diff line change
Expand Up @@ -205,7 +205,7 @@ This document tracks feature parity between IronClaw (Rust implementation) and O
| Sender_id in trusted metadata | ✅ | ❌ | Exposed in system metadata |
| Per-group `systemPrompt` injection | ✅ | ❌ | Per-group/per-direct system prompts injected via `GroupSystemPrompt` (Telegram, Discord, WhatsApp, BlueBubbles) |
| Visible reply enforcement | ✅ | ❌ | `messages.visibleReplies` requires output via `message(action=send)`; group-scope override available |
| Active-run steering queue | ✅ | ❌ | `messages.queue` `steer` mode (default) drains queued messages at next model boundary; `queue` legacy one-at-a-time |
| Active-run steering queue | ✅ | 🚧 | Reborn queues busy-thread user messages as steering input for the active run and WebUI shows them as queued until the loop consumes them; `queue` legacy one-at-a-time remains follow-up |
| Tool-progress streaming into previews | ✅ | ❌ | Tool progress shown in live preview edits (Discord/Slack/Telegram/Mattermost/Matrix) |
| `dmPolicy="open"` semantics | ✅ | 🚧 | Public open-DM only with effective wildcard; pairing-store senders no longer count for DM audits (OpenClaw fixed across all channels) |

Expand Down
17 changes: 16 additions & 1 deletion crates/ironclaw_agent_loop/src/executor/canonical.rs
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,19 @@ impl DefaultExecutorPipeline {
}
InputStep::Exit(exit) => return Ok(exit),
}
if !pending_input_ack.is_empty() {
state = CheckpointStage
.process(
ctx,
CheckpointInput {
state,
kind: CheckpointKind::BeforeModel,
},
)
.await?
.state;
pending_input_ack.ack(host).await?;
}

match self
.prompt
Expand All @@ -96,7 +109,9 @@ impl DefaultExecutorPipeline {
)
.await?
{
PromptStep::Exit(exit) => return Ok(exit),
PromptStep::Exit(exit) => {
return Ok(exit);
}

PromptStep::Prepared(prompt) => {
let prompt = *prompt;
Expand Down
14 changes: 14 additions & 0 deletions crates/ironclaw_agent_loop/src/executor/capabilities.rs
Original file line number Diff line number Diff line change
Expand Up @@ -983,6 +983,20 @@ impl CapabilityStage {
diagnostic_ref: None,
};
}
CapabilityOutcome::Completed(result) => {
clear_matching_pending_approval_resume(&mut state, &call);
clear_matching_pending_auth_resume(&mut state, &call);
clear_matching_pending_external_tool_resume(&mut state, &call);
append_completed_capability_result(
ctx.host,
&mut state,
&call,
result,
capability_batch,
)
.await?;
return Ok(BatchStep::Continue(Box::new(state)));
}
promoted => {
return Box::pin(self.handle_capability_outcome(
ctx,
Expand Down
7 changes: 6 additions & 1 deletion crates/ironclaw_agent_loop/src/executor/input.rs
Original file line number Diff line number Diff line change
Expand Up @@ -200,6 +200,9 @@ pub(super) fn consume_drainable_inputs(
});
}
let last_ack = &batch.input_acks[consumed_len - 1];
if drained && state.prompt_context_cursor.is_none() {
state.prompt_context_cursor = Some(state.input_cursor.clone());
}
state.input_cursor = last_ack.cursor.clone();
let ack_tokens = batch
.input_acks
Expand All @@ -221,7 +224,9 @@ fn user_facing_input_matches_drain_mode(input: &LoopInput, mode: UserFacingInput
UserFacingInputDrainMode::FollowUp => {
matches!(
input,
LoopInput::FollowUp { .. } | LoopInput::UserMessage { .. }
LoopInput::FollowUp { .. }
| LoopInput::UserMessage { .. }
| LoopInput::Steering { .. }
)
}
}
Expand Down
5 changes: 5 additions & 0 deletions crates/ironclaw_agent_loop/src/executor/loop_exit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,11 @@ pub(super) async fn try_final_answer_nudge(
// `None`s) is what actually forces a tool-free provider call. See the comment
// on `LoopModelRequest.capability_view` construction further down.
let context_plan = ctx.planner.context().plan_context_request(state).await;
// Consume the prompt-context handoff once, mirroring the main prompt-build
// path (`BuiltPromptBundle`). This nudge builds its own prompt bundle from
// the plan, so the lagging cursor must be cleared here too — otherwise it
// would persist in the checkpointed final state.
state.prompt_context_cursor = None;
let mut request = context_plan.request;
request.surface_version = None;
request.capability_view = None;
Expand Down
2 changes: 2 additions & 0 deletions crates/ironclaw_agent_loop/src/executor/prompt.rs
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ impl BuiltPromptBundle {
) -> Result<Self, AgentLoopExecutorError> {
let bundle =
build_prompt_bundle_for_surface(ctx, state, surface_version, capability_view).await?;
state.prompt_context_cursor = None;
refresh_compaction_prompt_from_index(state, &bundle.compaction_message_index);
Ok(bundle)
}
Expand All @@ -115,6 +116,7 @@ impl BuiltPromptBundle {
self,
state: &mut LoopExecutionState,
) -> Vec<LoopModelMessage> {
state.prompt_context_cursor = None;
refresh_compaction_prompt_from_index(state, &self.compaction_message_index);
self.messages
}
Expand Down
86 changes: 86 additions & 0 deletions crates/ironclaw_agent_loop/src/executor/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,7 @@ async fn reply_only_drains_follow_up_before_stop_strategy_completes() {
assert_eq!(
host.checkpoint_kinds(),
vec![
LoopCheckpointKind::BeforeModel,
LoopCheckpointKind::BeforeModel,
LoopCheckpointKind::BeforeModel,
LoopCheckpointKind::Final,
Expand All @@ -167,6 +168,45 @@ async fn reply_only_drains_follow_up_before_stop_strategy_completes() {
assert_eq!(final_staged_state(&host).stop_state.turns_completed, 2);
}

#[tokio::test]
async fn reply_only_drains_steering_before_stop_strategy_completes() {
let host = MockHost::new(vec![reply_response(), reply_response()]);
let run_context = host.run_context().clone();
let host = host.with_input_batches(vec![
LoopInputBatch {
inputs: Vec::new(),
input_acks: Vec::new(),
next_cursor: input_cursor(&run_context, "input-cursor:no-input"),
},
LoopInputBatch {
inputs: vec![LoopInput::Steering {
message_ref: message_ref("msg:steering-follow-up"),
}],
input_acks: vec![input_ack(
&run_context,
"input-cursor:after-steering-follow-up",
"input-ack:after-steering-follow-up",
)],
next_cursor: input_cursor(&run_context, "input-cursor:after-steering-follow-up"),
},
]);
let executor = CanonicalAgentLoopExecutor;
let state = LoopExecutionState::initial_for_run(host.run_context());

let exit = executor
.execute_family(&crate::families::default(), &host, state)
.await
.expect("execute");

assert!(matches!(exit, LoopExit::Completed(_)));
assert_eq!(host.model_requests().len(), 2);
assert_eq!(
host.acked_input_tokens(),
vec![LoopInputAckToken::new("input-ack:after-steering-follow-up").expect("valid")]
);
assert_eq!(final_staged_state(&host).stop_state.turns_completed, 2);
}

#[tokio::test]
async fn reply_only_uses_configured_stop_strategy_decision() {
let host = MockHost::new(vec![reply_response(), reply_response()]);
Expand Down Expand Up @@ -1065,6 +1105,10 @@ async fn input_stage_steering_drain_carries_pending_ack() {
state.input_cursor,
input_cursor(&run_context, "input-cursor:after-user")
);
assert_eq!(
state.prompt_context_cursor,
Some(LoopInputCursor::origin_for_run(&run_context))
);
assert!(host.acked_input_tokens().is_empty());
pending_input_ack.ack(&host).await.expect("ack inputs");
assert_eq!(
Expand Down Expand Up @@ -1117,6 +1161,10 @@ async fn input_stage_steering_input_is_drained_like_user_message() {
state.input_cursor,
input_cursor(&run_context, "input-cursor:after-steering")
);
assert_eq!(
state.prompt_context_cursor,
Some(LoopInputCursor::origin_for_run(&run_context))
);
}
InputStep::Exit(exit) => panic!("expected continue, got {exit:?}"),
}
Expand Down Expand Up @@ -1684,6 +1732,44 @@ async fn no_progress_nudge_synthesizes_reply_when_gate_enabled() {
}
}

#[tokio::test]
async fn nudge_clears_prompt_context_cursor_in_final_state() {
// Regression: the final-answer nudge builds its own prompt bundle from
// `plan_context_request`, so it must consume the prompt-context handoff
// once — like the main prompt-build path (`BuiltPromptBundle`) — instead of
// leaving a lagging `prompt_context_cursor` in the checkpointed final state.
let host = MockHost::new(vec![reply_response_with_text("Here is the final answer.")])
.with_driver_nudges_enabled();
let run_context = host.run_context().clone();
let family = crate::families::default();
let ctx = StageContext {
planner: family.planner(),
host: &host,
};
let mut state = LoopExecutionState::initial_for_run(host.run_context());
state.prompt_context_cursor = Some(LoopInputCursor::origin_for_run(&run_context));

let exit = ExitStage
.process(
ctx,
ExitInput {
state,
kind: StopKind::NoProgressDetected,
},
)
.await
.expect("exit stage");

assert!(
matches!(exit, LoopExit::Completed(_)),
"nudge with a queued reply should complete the turn"
);
assert!(
final_staged_state(&host).prompt_context_cursor.is_none(),
"nudge must clear prompt_context_cursor before the final checkpoint"
);
}

#[tokio::test]
async fn no_progress_skips_nudge_when_gate_disabled() {
// Gate OFF: even with a model reply available, no tool-free nudge call is
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -292,6 +292,7 @@ async fn steering_drain_acks_only_after_cursor_checkpoint_is_durable() {
vec![
"checkpoint:before_model".to_string(),
"ack_inputs".to_string(),
"checkpoint:before_model".to_string(),
"checkpoint:final".to_string(),
]
);
Expand Down
9 changes: 7 additions & 2 deletions crates/ironclaw_agent_loop/src/executor/tests/support.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,10 @@ use ironclaw_turns::{
use crate::{
default_planner::DefaultPlanner,
family::{ComponentDigest, ComponentIdentity, LoopFamily, LoopFamilyId},
state::{CheckpointKind, GateStrategyState, LoopExecutionState, StopStrategyState},
state::{
CheckpointKind, GateStrategyState, LoopExecutionState, RecoveryAttemptClass,
StopStrategyState,
},
strategies::{
CapabilityErrorClass, CapabilityErrorSummary, CapabilityFilter, CapabilityStrategy,
ContextStrategy, DefaultBudgetStrategy, DefaultCompactionStrategy, GateHandlingStrategy,
Expand Down Expand Up @@ -464,7 +467,9 @@ impl RecoveryStrategy for RetryPolicyDeniedRecoveryStrategy {
) -> RecoveryOutcome {
if err.class == CapabilityErrorClass::PolicyDenied {
return RecoveryOutcome::Retry {
recovery: state.recovery_state.clone(),
recovery: state
.recovery_state
.with_incremented_attempts_for(RecoveryAttemptClass::CapabilityInternal),
scope: RetryScope::Call,
alter: None,
};
Expand Down
7 changes: 7 additions & 0 deletions crates/ironclaw_agent_loop/src/state.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,6 +55,12 @@ pub struct LoopExecutionState {
pub result_refs: Vec<LoopResultRef>,
pub last_gate: Option<LoopGateRef>,
pub input_cursor: LoopInputCursor,
/// Cursor to use for the next prompt context after draining queued
/// user-facing input. This intentionally lags behind `input_cursor` so the
/// message can be acked without disappearing from the prompt bundle's
/// `after` window.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub prompt_context_cursor: Option<LoopInputCursor>,
pub surface_version: Option<CapabilitySurfaceVersion>,

// executor-observed (populated by executor; read-only to strategies)
Expand Down Expand Up @@ -258,6 +264,7 @@ impl LoopExecutionState {
result_refs: Vec::new(),
last_gate: None,
input_cursor: LoopInputCursor::origin_for_run(context),
prompt_context_cursor: None,
surface_version: None,
recent_call_signatures: BoundedRing::new(),
seen_capability_output_digests: BoundedRing::new(),
Expand Down
34 changes: 30 additions & 4 deletions crates/ironclaw_agent_loop/src/strategies/context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,12 @@ impl ContextStrategy for DefaultContextStrategy {
ContextPlan {
request: LoopPromptBundleRequest {
mode: PromptMode::TextOnly,
context_cursor: None,
context_cursor: Some(
state
.prompt_context_cursor
.clone()
.unwrap_or_else(|| state.input_cursor.clone()),
),
surface_version: None,
checkpoint_state_ref: None,
max_messages: Some(self.max_messages.max(1)),
Expand Down Expand Up @@ -130,8 +135,9 @@ mod tests {
AgentLoopDriverDescriptor, RunProfileId, RunProfileVersion, TurnId, TurnRunId, TurnScope,
run_profile::{
CancellationPolicy, CapabilitySurfaceProfileId, CheckpointPolicy, CheckpointSchemaId,
ConcurrencyClass, ContextProfileId, LoopDriverId, LoopRunContext, ModelProfileId,
PromptMode, RedactedRunProfileProvenance, ResolvedRunProfile, ResourceBudgetPolicy,
ConcurrencyClass, ContextProfileId, LoopDriverId, LoopInputCursor,
LoopInputCursorToken, LoopRunContext, ModelProfileId, PromptMode,
RedactedRunProfileProvenance, ResolvedRunProfile, ResourceBudgetPolicy,
ResourceBudgetTier, RunClassId, RunProfileFingerprint, RuntimeProfileConstraints,
SchedulingClass, SteeringPolicy,
},
Expand Down Expand Up @@ -237,11 +243,31 @@ mod tests {
assert!(request.request.inline_messages.is_empty());
assert!(!request.emitted_admission_control);
assert!(!request.emitted_repeated_call_warning);
assert!(request.request.context_cursor.is_none());
assert_eq!(
request.request.context_cursor,
Some(state.input_cursor.clone())
);
assert!(request.request.surface_version.is_none());
assert!(request.request.checkpoint_state_ref.is_none());
}

#[tokio::test]
async fn plan_context_request_prefers_pending_prompt_context_cursor() {
let strategy = DefaultContextStrategy::default();
let run_context = test_run_context();
let mut state = LoopExecutionState::initial_for_run(&run_context);
let original_cursor = state.input_cursor.clone();
state.input_cursor = LoopInputCursor::from_host_token(
&run_context,
LoopInputCursorToken::new("input-cursor:after-queued").expect("valid cursor"),
);
state.prompt_context_cursor = Some(original_cursor.clone());

let request = strategy.plan_context_request(&state).await;

assert_eq!(request.request.context_cursor, Some(original_cursor));
}

#[tokio::test]
async fn plan_context_request_clamps_zero_to_one() {
let strategy = DefaultContextStrategy { max_messages: 0 };
Expand Down
Loading
Loading