Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -127,3 +127,6 @@ mutants.out.old/

# Skill bundles staged into a workspace at activation. Build output, never source.
.skills/

# Added by cleanup skill: campaign state, local-only.
Comment thread
henrypark133 marked this conversation as resolved.
.cleanup/
Original file line number Diff line number Diff line change
Expand Up @@ -875,7 +875,14 @@ fn reborn_contracts_crates_carry_a_checked_size_ceiling() {
// belong beside the turn contract consumed across loop families.
// 20_156 -> 20_334 (2026-08-19, #7686 restack): capability dispatch-result
// declarations retained beside main's provider-neutral output contracts.
("ironclaw_host_api", 20_334),
// 20_334 -> 20_369 (2026-08-19, #7752 merge): +14 for `turn::ActivationProvenance`,
// the enum tagging why a run was created (Human / ParentAgent / System).
// It is turn vocabulary, and `ironclaw_turns` may not own turn
// vocabulary — a new turn type goes to `host_api` by that crate's own
// charter — so there is no lower crate to move it to. Pure DTO: three
// unit variants and a serde derive, no behavior. Value re-captured from
// this test's own report on the merged tree, never counted by eye.
("ironclaw_host_api", 20_369),
// 14_479 -> 13_949 (2026-08-07, #7157): downward re-capture after the
// delivery-heuristic vocabulary (stored trigger delivery targets and
// their run-profile plumbing) left this crate with the two-lane
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,7 @@ fn coordinator_submit_request(submission: ConversationTurnSubmission) -> SubmitT
product_context.execution_policy = submission.execution_policy;
}
SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
scope: submission.scope,
actor: submission.actor,
Expand Down Expand Up @@ -139,6 +140,13 @@ pub(crate) fn turn_submission_error(error: TurnError) -> TurnSubmissionError {
TurnSubmissionErrorCategory::InvalidRequest,
TurnSubmissionRetry::Permanent,
),
// A thread that spent its autonomous-wake budget is parked pending
// human attention, not permanently refused: the trigger poller may
// retry, and a human activation restores the budget.
AdmissionRejectionReason::SystemWakeStreak => (
TurnSubmissionErrorCategory::AdmissionRejected,
TurnSubmissionRetry::RetryableAfterKeyRotation,
),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
AdmissionRejectionReason::Policy | AdmissionRejectionReason::Unauthorized => (
TurnSubmissionErrorCategory::Unauthorized,
TurnSubmissionRetry::Permanent,
Expand Down Expand Up @@ -232,6 +240,13 @@ mod tests {
TurnSubmissionErrorCategory::InvalidRequest,
TurnSubmissionRetry::Permanent,
),
(
TurnError::AdmissionRejected(AdmissionRejection::new(
AdmissionRejectionReason::SystemWakeStreak,
)),
TurnSubmissionErrorCategory::AdmissionRejected,
TurnSubmissionRetry::RetryableAfterKeyRotation,
),
(
TurnError::AdmissionRejected(AdmissionRejection::new(
AdmissionRejectionReason::Policy,
Expand Down Expand Up @@ -351,7 +366,7 @@ mod tests {
distinct.len(),
12,
"TurnError has 12 variants; the class table must name every one \
(AdmissionRejected appears four times, once per rejection reason)"
(AdmissionRejected appears six times, once per rejection reason)"
);
}

Expand Down
3 changes: 3 additions & 0 deletions crates/app/ironclaw_composition/src/factory/auth_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -192,6 +192,7 @@ async fn standalone_oauth_turn_gate_callback_resumes_default_turn_coordinator()
let actor = TurnActor::new(UserId::new("alice").unwrap());
let submit = turn_coordinator
.submit_turn(SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope: scope.clone(),
Expand Down Expand Up @@ -824,6 +825,7 @@ async fn submit_and_block_provider_auth_run(
) -> TurnRunId {
let submit = turn_coordinator
.submit_turn(SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope: scope.clone(),
Expand Down Expand Up @@ -968,6 +970,7 @@ async fn submit_and_block_auth_run(
) -> ironclaw_host_api::turn::TurnRunId {
let submit = turn_coordinator
.submit_turn(SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope: scope.clone(),
Expand Down
2 changes: 2 additions & 0 deletions crates/app/ironclaw_composition/src/factory/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2040,6 +2040,7 @@ async fn production_libsql_turn_state_uses_configured_runtime_identity() {
Some(owner.clone()),
);
let submit = ironclaw_turns::SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope,
Expand Down Expand Up @@ -2129,6 +2130,7 @@ async fn production_libsql_turn_state_uses_default_runtime_identity_when_unconfi
Some(owner.clone()),
);
let submit = ironclaw_turns::SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope,
Expand Down
1 change: 1 addition & 0 deletions crates/app/ironclaw_composition/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2328,6 +2328,7 @@ impl RebornRuntime {
let response = match self
.turn_coordinator
.submit_turn(SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: accepted.replay_metadata.resolved_model.clone(),
scope: scope.clone(),
actor: TurnActor::new(self.actor_user_id.clone()),
Expand Down
2 changes: 2 additions & 0 deletions crates/app/ironclaw_composition/src/runtime/tests/core.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4257,6 +4257,7 @@ async fn cancel_run_propagates_to_children_when_event_sink_is_unavailable() {
let parent = runtime
.turn_coordinator
.submit_turn(SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope: parent_scope.clone(),
Expand Down Expand Up @@ -7195,6 +7196,7 @@ async fn deferred_busy_message_not_auto_submitted_after_run_cancellation() {
let submitted_a = runtime
.turn_coordinator
.submit_turn(SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope: scope.clone(),
Expand Down
1 change: 1 addition & 0 deletions crates/app/ironclaw_composition/tests/service_factory.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1584,6 +1584,7 @@ async fn production_postgres_process_journal_pool_writes_rows_the_data_plane_rea
ironclaw_turns::TurnCoordinator::submit_turn(
services.turn_coordinator_for_test().as_ref(),
ironclaw_turns::SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope,
Expand Down
1 change: 1 addition & 0 deletions crates/app/ironclaw_composition/tests/webui_v2_serve.rs
Original file line number Diff line number Diff line change
Expand Up @@ -664,6 +664,7 @@ mod openai_compat_mount_tests {
let response = self
.coordinator
.submit_turn(SubmitTurnRequest {
subagent_activation_provenance: None,
scope,
actor: caller.actor(),
accepted_message_ref: accepted_message_ref.clone(),
Expand Down
35 changes: 35 additions & 0 deletions crates/contracts/ironclaw_host_api/src/turn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -859,6 +859,21 @@ impl From<RunOriginAdapter> for String {
}
}

/// Why a run was created on its thread. Set once at run creation and immutable
/// thereafter — the derived activation-streak caps read bounded windows of this
/// field instead of maintaining a stored counter.
///
/// `Human` is the ordinary case and resets both streaks. `ParentAgent` tags a
/// parent re-activating one of its own children. `System` tags a background
/// subagent completion waking its parent.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ActivationProvenance {
Human,
ParentAgent,
System,
}

/// The lifecycle position of one turn run. Terminal, gate-parked, and in-flight
/// statuses are distinguished by the predicates below rather than by scattered
/// match tables at the call sites.
Expand Down Expand Up @@ -1185,6 +1200,26 @@ pub enum SubmitTurnResponse {
mod tests {
use super::*;

#[test]
fn activation_provenance_wire_strings_are_snake_case() {
assert_eq!(
serde_json::to_value(ActivationProvenance::Human).expect("serialize"),
serde_json::json!("human")
);
assert_eq!(
serde_json::to_value(ActivationProvenance::ParentAgent).expect("serialize"),
serde_json::json!("parent_agent")
);
assert_eq!(
serde_json::to_value(ActivationProvenance::System).expect("serialize"),
serde_json::json!("system")
);

let round_tripped: ActivationProvenance =
serde_json::from_value(serde_json::json!("parent_agent")).expect("deserialize");
assert_eq!(round_tripped, ActivationProvenance::ParentAgent);
}

#[test]
fn blocked_external_tool_status_is_non_terminal_and_keeps_lock() {
assert!(!TurnStatus::BlockedExternalTool.is_terminal());
Expand Down
1 change: 1 addition & 0 deletions crates/domains/ironclaw_conversations/src/inbound.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1381,6 +1381,7 @@ mod tests {
product_context.execution_policy = submission.execution_policy;
}
SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
scope: submission.scope,
actor: submission.actor,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -4110,6 +4110,7 @@ fn submit_turn_request(submission: ConversationTurnSubmission) -> SubmitTurnRequ
product_context.execution_policy = submission.execution_policy;
}
SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope: submission.scope,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -2257,6 +2257,7 @@ pub(crate) fn http_without_body_then_operation_failed_wat() -> String {

pub(crate) fn submit_turn_request(thread: &str, idempotency_key: &str) -> SubmitTurnRequest {
SubmitTurnRequest {
subagent_activation_provenance: None,
requested_model: None,
output_contract: None,
scope: TurnScope::new(
Expand Down
10 changes: 10 additions & 0 deletions crates/kernel/ironclaw_processes/src/journal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -951,6 +951,16 @@ pub trait ProcessSnapshotSource: Send + Sync {
&self,
scope: &ResourceScope,
) -> Result<Vec<JournaledProcessSnapshot>, Self::Error>;

/// Newest-first, bounded read of one scope's agent-turn processes.
///
/// Unlike [`Self::process_snapshots`] this is explicitly bounded: callers
/// that need a fixed recent window must not enumerate a whole scope.
async fn recent_agent_turn_snapshots(
&self,
scope: &ResourceScope,
limit: u32,
) -> Result<Vec<JournaledProcessSnapshot>, Self::Error>;
}

#[async_trait]
Expand Down
18 changes: 18 additions & 0 deletions crates/kernel/ironclaw_processes/src/journal_store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -733,6 +733,24 @@ where
snapshots.sort_by_key(|snapshot| snapshot.process_id.as_uuid());
Ok(snapshots)
}

async fn recent_agent_turn_snapshots(
&self,
scope: &ResourceScope,
limit: u32,
) -> Result<Vec<JournaledProcessSnapshot>, Self::Error> {
self.ensure_materialized().await?;
// `is_system()`, not `== ResourceScope::system()`: the constructor
// mints a fresh `invocation_id` on every call, so an equality check
// against it can never match and the guard would be dead.
if scope.is_system() {
return Err(ProcessJournalStoreError::InvalidRequest(
"system-wide process snapshot reads are unbounded; use paged process journal reads"
.to_string(),
));
}
rows::recent_agent_turn_processes_for_scope(self.filesystem.as_ref(), scope, limit).await
}
}

#[async_trait]
Expand Down
74 changes: 74 additions & 0 deletions crates/kernel/ironclaw_processes/src/journal_store/rows.rs
Original file line number Diff line number Diff line change
Expand Up @@ -710,6 +710,80 @@ where
decode_process_rows(&rows)
}

/// Newest-first, bounded read of one scope's `AgentTurn` processes.
///
/// `process_scope_v3` is keyed on `(scope_key, created_at, process_id)` and
/// deliberately not on `process_kind`, while a thread's scope also holds
/// capability-invocation processes. A flat `LIMIT` could therefore be
/// satisfied entirely by non-`AgentTurn` rows, so this walks the descending
/// keyset a page at a time and filters by kind in memory until it has `limit`
/// agent-turn rows.
///
/// Bounded by `MAX_RECENT_PROCESS_PAGES` rather than walking a scope's whole
Comment thread
henrypark133 marked this conversation as resolved.
/// history: a thread whose newest ~8 pages hold no agent-turn run returns a
/// short window. The only caller is the derived activation-streak cap, which
/// treats a short window as "streak not established" — the same disposition it
/// gives a genuinely young thread.
pub(super) async fn recent_agent_turn_processes_for_scope<F>(
filesystem: &ScopedFilesystem<F>,
scope: &ResourceScope,
limit: u32,
) -> Result<Vec<JournaledProcessSnapshot>, ProcessJournalStoreError>
where
F: RootFilesystem,
{
const MAX_RECENT_PROCESS_PAGES: u32 = 8;

if limit == 0 {
return Ok(Vec::new());
}
let prefix = scoped_path(&format!("{MATERIALIZED_PREFIX}/process"))?;
let filter = eq_value("scope_key", IndexValue::Text(scope_owner_key(scope)?))?;
// Over-fetch per page so a scope interleaving capability-invocation
// processes with its runs still makes progress on every round trip.
let page_limit = limit.saturating_mul(4).clamp(1, Page::MAX_LIMIT);

let mut collected: Vec<JournaledProcessSnapshot> = Vec::new();
let mut cursor: Option<ironclaw_filesystem::OrderedQueryCursor> = None;
for _ in 0..MAX_RECENT_PROCESS_PAGES {
let mut page = ironclaw_filesystem::OrderedPage::new(
index_name("process_scope_v3")?,
index_key("created_at")?,
index_key("process_id")?,
SortDirection::Descending,
page_limit,
);
if let Some(cursor) = cursor.clone() {
page = page.after(cursor);
}
let rows = filesystem
.query_ordered(&ResourceScope::system(), &prefix, &filter, &page)
.await?;
let exhausted = (rows.len() as u32) < page_limit;
let snapshots = decode_process_rows(&rows)?;
let Some(last) = snapshots.last() else {
// Every row in this page decoded away (tombstones); without a
// decodable row there is no cursor to advance past, so stop rather
// than re-reading the same page.
break;
};
cursor = Some(ironclaw_filesystem::OrderedQueryCursor {
value: IndexValue::I64(last.created_at.timestamp_micros()),
tie_breaker: IndexValue::Text(last.process_id.as_uuid().to_string()),
});
Comment thread
coderabbitai[bot] marked this conversation as resolved.
collected.extend(
snapshots
.into_iter()
.filter(|snapshot| snapshot.process_kind == ProcessKind::AgentTurn),
);
if collected.len() >= limit as usize || exhausted {
break;
}
}
collected.truncate(limit as usize);
Ok(collected)
}

pub(super) async fn child_processes<F>(
filesystem: &ScopedFilesystem<F>,
parent_process_id: ProcessId,
Expand Down
Loading
Loading