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
11 changes: 10 additions & 1 deletion crates/app/ironclaw_composition/src/runtime.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,8 @@ use ironclaw_assistant::{
ApprovalResolverPort, ApprovalTurnRunLocator, AuthInteractionService,
DefaultApprovalInteractionService, DefaultAuthInteractionService,
OutboundPreferencesProductService, PersistentApprovalGranteeResolver,
RunStateApprovalInteractionReadModel, SuggestionsProcessCommitObserver,
RunOutcomeProcessCommitObserver, RunStateApprovalInteractionReadModel,
SuggestionsProcessCommitObserver,
};
use ironclaw_event_log::{
DurableAuditLog, DurableEventLog, EventError, NonBlockingEventSink, RuntimeEvent,
Expand Down Expand Up @@ -3241,6 +3242,14 @@ pub(crate) async fn build_runtime_with_resource_governor(
.map_err(|error| RebornRuntimeError::MalformedConfig {
reason: format!("suggestion generation observer wiring failed: {error}"),
})?;
processes
.subscribe_process_observer(Arc::new(RunOutcomeProcessCommitObserver::new(
Arc::clone(&services.notification_inbox),
Arc::clone(&thread_service),
)))
.map_err(|error| RebornRuntimeError::MalformedConfig {
reason: format!("run outcome notification observer wiring failed: {error}"),
})?;
let filesystem_skill_context_runtime = filesystem_skill_context_runtime(&services);
let (skill_context_source, skill_activation_source, skill_execution_adapter) = match (
configured_skill_context_source,
Expand Down
154 changes: 148 additions & 6 deletions crates/app/ironclaw_composition/tests/trigger_poller_e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -61,7 +61,12 @@ use ironclaw_outbound::{
CommunicationPreferenceRecord, DeliveryDefaultScope, TriggeredRunDeliveryOutcomeKind,
TriggeredRunDeliveryStore,
};
use ironclaw_product_contracts::surface::{ProductSurfaceCaller, ProductSurfaceInvokeRequest};
use ironclaw_product_contracts::{
notification_inbox::{
NOTIFICATIONS_VIEW, ProductListNotificationsResponse, ProductNotificationKind,
},
surface::{ProductSurfaceCaller, ProductSurfaceInvokeRequest, ProductSurfaceQueryRequest},
};
use ironclaw_triggers::{
TRIGGER_TRUSTED_ADAPTER_INSTALLATION_ID, TRIGGER_TRUSTED_ADAPTER_KIND,
TRIGGER_TRUSTED_EXTERNAL_ACTOR_NAMESPACE, TriggerDeliveryTargetId, TriggerExecutionSpec,
Expand Down Expand Up @@ -1548,8 +1553,8 @@ where
/// model (spec §8): a scheduled fire's RESULT is never pushed to a channel.
///
/// due trigger -> trusted ingress -> real Reborn run -> persisted final reply
/// in the fire's own run thread -> background-run notifier -> **nothing on any
/// channel**.
/// in the fire's own run thread -> durable Inbox outcome, while the background
/// delivery notifier still sends **nothing on any channel**.
///
/// The two arms are the two ways the retired routing used to pick a
/// destination: QA-9B had none (it inherited the creator's default), QA-9D
Expand Down Expand Up @@ -1601,7 +1606,7 @@ async fn scheduled_trigger_results_are_never_pushed_to_a_channel_across_restart(
)
.await;

wait_for_recorded_outcome(
let default_run_id = wait_for_recorded_outcome(
&repository,
&delivery_store,
default_target_trigger,
Expand All @@ -1615,7 +1620,7 @@ async fn scheduled_trigger_results_are_never_pushed_to_a_channel_across_restart(
TriggeredRunDeliveryOutcomeKind::Suppressed,
)
.await;
wait_for_recorded_outcome(
let explicit_run_id = wait_for_recorded_outcome(
&repository,
&delivery_store,
explicit_target_trigger,
Expand Down Expand Up @@ -1682,6 +1687,71 @@ async fn scheduled_trigger_results_are_never_pushed_to_a_channel_across_restart(
"the suppressible routine must perform prior tool work before its typed result"
);

let caller = ProductSurfaceCaller::new(
TenantId::new(TENANT).expect("tenant"),
UserId::new(USER).expect("user"),
Some(AgentId::new(AGENT).expect("agent")),
None,
);
let surface = runtime.product_surface(None).expect("product surface");
let inbox = tokio::time::timeout(Duration::from_secs(5), async {
loop {
let page = surface
.query(
caller.clone(),
ProductSurfaceQueryRequest {
view_id: NOTIFICATIONS_VIEW.id.to_string(),
input: json!({ "limit": 10 }),
cursor: None,
limit: None,
},
)
.await
.expect("query notification Inbox");
let response: ProductListNotificationsResponse = serde_json::from_value(
page.items
.into_iter()
.next()
.expect("notification response payload"),
)
.expect("decode notification response");
if response.notifications.len() == 2 {
break response;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("background completion notifications arrive");
assert!(
inbox
.notifications
.iter()
.all(|notification| notification.kind == ProductNotificationKind::RunCompleted),
"only selected completed results are materialized: {:?}",
inbox.notifications
);
let notified_runs = inbox
.notifications
.iter()
.filter_map(|notification| notification.turn_run_id.clone())
.collect::<std::collections::BTreeSet<_>>();
assert_eq!(
notified_runs,
std::collections::BTreeSet::from(
[default_run_id.to_string(), explicit_run_id.to_string(),]
),
"ordinary scheduled completions notify, while NothingToReport stays suppressed"
);

// The identities themselves are what dedupe on replay.
let record_count_before_restart = inbox.notifications.len();
let identities_before_restart = inbox
.notifications
.iter()
.map(|notification| notification.id.clone())
.collect::<std::collections::BTreeSet<_>>();

tokio::time::sleep(Duration::from_millis(250)).await;
assert!(
slack_provider.provider_messages().is_empty(),
Expand All @@ -1696,7 +1766,79 @@ async fn scheduled_trigger_results_are_never_pushed_to_a_channel_across_restart(
)
.await;
register_delivery_targets(&restarted);
tokio::time::sleep(Duration::from_millis(300)).await;
let restarted_surface = restarted.product_surface(None).expect("restarted surface");
let replayed = tokio::time::timeout(Duration::from_secs(5), async {
loop {
let replayed_page = restarted_surface
.query(
caller.clone(),
ProductSurfaceQueryRequest {
view_id: NOTIFICATIONS_VIEW.id.to_string(),
input: json!({ "limit": 10 }),
cursor: None,
limit: None,
},
)
.await
.expect("query Inbox while observer replay settles");
let replayed: ProductListNotificationsResponse = serde_json::from_value(
replayed_page
.items
.into_iter()
.next()
.expect("replayed notification payload"),
)
.expect("decode replayed Inbox");
if replayed.notifications.len() >= record_count_before_restart {
break replayed;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("observer replay reaches the pre-restart notification count");
assert_eq!(
replayed.notifications.len(),
record_count_before_restart,
"restart/replay must not add a second record under an already-existing id"
);

tokio::time::sleep(Duration::from_millis(100)).await;
let settled_page = restarted_surface
.query(
caller,
ProductSurfaceQueryRequest {
view_id: NOTIFICATIONS_VIEW.id.to_string(),
input: json!({ "limit": 10 }),
cursor: None,
limit: None,
},
)
.await
.expect("query Inbox after observer replay settles");
let settled: ProductListNotificationsResponse = serde_json::from_value(
settled_page
.items
.into_iter()
.next()
.expect("settled notification payload"),
)
.expect("decode settled Inbox");
assert_eq!(
settled.notifications.len(),
record_count_before_restart,
"restart/replay must remain at the pre-restart count after a settle window"
);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
assert_eq!(
settled
.notifications
.iter()
.map(|notification| notification.id.clone())
.collect::<std::collections::BTreeSet<_>>(),
identities_before_restart,
"restart/replay resolves to the same notification identities, so a replayed \
commit is absorbed by its stable id instead of arriving twice"
);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
restarted
.shutdown()
.await
Expand Down
6 changes: 6 additions & 0 deletions crates/domains/ironclaw_notifications/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,12 @@ publication, and the filesystem-backed store contract.
attempt; those remain in `ironclaw_outbound`. Product-specific projection and
read policy remain in `ironclaw_assistant`.

Run outcome records are produced by `ironclaw_assistant` from committed process
journal transitions. Successful scheduled runs require an exact durable
finalized-reply lookup; external delivery failures remain separate from the run
outcome and come from the outbound delivery caller. The domain store does not
infer either fact.

The store persists through `ScopedFilesystem`, so composition selects the real
backend and provides tenant/user-rewritten mounts. Records contain metadata and
typed references only; prompts, replies, tool payloads, secrets, and backend
Expand Down
2 changes: 1 addition & 1 deletion crates/kernel/ironclaw_processes/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,6 @@ ironclaw_event_store = { path = "../../events/ironclaw_event_store" }
libsql = { version = "0.9", default-features = false, features = ["core", "replication", "remote", "tls"] }
deadpool-postgres = "0.14"
tempfile = "3"
tokio = { version = "1", features = ["macros", "rt", "time"] }
tokio = { version = "1", features = ["macros", "rt", "test-util", "time"] }
tokio-postgres = { version = "0.7", features = ["with-uuid-1", "with-chrono-0_4", "with-serde_json-1"] }
uuid = { version = "1", features = ["v4", "serde"] }
50 changes: 50 additions & 0 deletions crates/kernel/ironclaw_processes/src/journal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -513,6 +513,10 @@ pub struct ProcessJournalEntry {
pub struct ProcessJournalCommit {
pub state: JournaledProcessSnapshot,
pub kind: ProcessJournalKind,
/// Timestamp of the journal transition that produced this committed
/// state. Observers use it for replay-stable materialized records.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub occurred_at: Option<Timestamp>,
Comment thread
coderabbitai[bot] marked this conversation as resolved.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub sanitized_reason: Option<String>,
}
Expand Down Expand Up @@ -1489,6 +1493,52 @@ mod tests {
assert_eq!(absent.checkpoint_kind, None);
}

#[test]
fn process_journal_commit_without_occurred_at_remains_readable() {
let commit = ProcessJournalCommit {
state: JournaledProcessSnapshot {
process_id: ProcessId::new(),
process_kind: ProcessKind::Internal,
scope: ResourceScope {
tenant_id: TenantId::new("tenant-legacy-commit").expect("tenant"),
user_id: UserId::new("user-legacy-commit").expect("user"),
agent_id: None,
project_id: None,
mission_id: None,
thread_id: None,
invocation_id: ironclaw_host_api::ids::InvocationId::new(),
},
status: ProcessLifecycleStatus::Completed,
suspension: None,
checkpoint_ref: None,
checkpoint_kind: None,
input_ref: None,
failure: None,
journal_cursor: ProcessJournalCursor(1),
lease: None,
crash_reclaim_count: 0,
created_at: chrono::Utc::now(),
owner_user_id: None,
concurrency_class: None,
parent_process_id: None,
root_process_id: None,
metadata: Value::Null,
},
kind: ProcessJournalKind::Completed,
occurred_at: Some(chrono::Utc::now()),
sanitized_reason: None,
};
let mut legacy = serde_json::to_value(commit).expect("serialize current commit");
legacy
.as_object_mut()
.expect("commit object")
.remove("occurred_at");

let decoded: ProcessJournalCommit =
serde_json::from_value(legacy).expect("legacy commit remains readable");
assert_eq!(decoded.occurred_at, None);
}

#[test]
fn opaque_refs_reject_empty_oversized_and_control_values() {
assert!(ProcessCheckpointRef::new("").is_err());
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -548,6 +548,7 @@ fn committed_batch(entries: Vec<ProcessJournalEntry>) -> Option<CommittedBatch>
entry.committed_state.map(|state| ProcessJournalCommit {
state: *state,
kind: entry.kind,
occurred_at: entry.occurred_at,
sanitized_reason: entry.sanitized_reason,
})
})
Expand Down
16 changes: 16 additions & 0 deletions crates/kernel/ironclaw_processes/src/journal_store/observer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@ use crate::{
ProcessJournalObserverRegistry,
};

const MAX_OBSERVER_REPLAY_ATTEMPTS: u32 = 8;

#[derive(Clone)]
pub(super) struct RegisteredProcessObserver {
pub(super) id: String,
Expand Down Expand Up @@ -266,6 +268,7 @@ where
commits.push(ProcessJournalCommit {
state: state.clone(),
kind: entry.kind,
occurred_at: entry.occurred_at,
sanitized_reason: entry.sanitized_reason.clone(),
});
}
Expand Down Expand Up @@ -340,12 +343,25 @@ where
tokio::spawn(async move {
let _replay_guard = ObserverReplayGuard(Arc::clone(&observer.replay_running));
let mut delay = Duration::from_millis(250);
let mut attempt = 0_u32;
loop {
attempt += 1;
match store.replay_durable_observer_once(&observer, None).await {
Ok(()) => break,
Err(error) if attempt >= MAX_OBSERVER_REPLAY_ATTEMPTS => {
tracing::error!(
observer_id = %observer.id,
attempts = attempt,
%error,
"durable process observer replay exhausted its retry budget; \
cursor remains unacknowledged"
);
break;
}
Err(error) => {
tracing::debug!(
observer_id = %observer.id,
attempt,
%error,
"durable process observer delivery will retry"
);
Expand Down
Loading
Loading