Repository navigation
[DRAFT] Optimize hosted Postgres turn-state latency - #5667
serrrfirat wants to merge 36 commits into
Conversation
📝 WalkthroughWalkthroughAdds a StorageTxn::reserve_sequence primitive across libSQL/Postgres/scoped filesystems; replaces PersistentResourceGovernor with a new journaling FilesystemResourceGovernor everywhere it is wired; adds a PostgresSecretStore and a filesystem row-store turn-state backend; hardens trigger repositories; and introduces a latency benchmarking harness with goal/spec docs. ChangesFilesystem reserve_sequence primitive
FilesystemResourceGovernor rollout
Postgres secret store & secret index scoping
Trigger repository hardening
Filesystem turn-state row store & lifecycle churn
RebornBuildInput Postgres constructor
Hosted single-tenant Postgres latency harness
Estimated code review effort: 5 (Critical) | ~150 minutes Possibly related issues
Possibly related PRs
Suggested reviewers: Advisory-xact-lock migration path (postgres.rs run_migrations): verify lock released on early-return error paths — CLAUDE.md transaction-safety invariant, not caught by clippy. Delta journal ack propagation on flusher-thread panic (filesystem_governor.rs): confirm pending oneshot receivers don't hang forever — sandbox/trust concern (silent stall vs. hard fail). PostgresSecretStore AAD derivation from scope/handle: confirm no cross-tenant AAD collision — secrets invariant, needs explicit test, not typechecked. Row store durable-delta suppression for Tier-2 cache evictions: confirm this can't cause silent data loss on process crash before compaction — flag as durability risk, not style. 🚥 Pre-merge checks | ✅ 3 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (3 passed)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Code Review
This pull request introduces a journaled filesystem-backed resource governor (FilesystemResourceGovernor) and a filesystem turn-state row store (FilesystemTurnStateRowStore) to optimize performance and reduce contention on single-blob state documents. It also adds a comprehensive latency harness to compare libSQL and PostgreSQL backends across several production-shaped workloads. The code review identified three key issues: a correctness bug where historical usage from previous periods is incorrectly restored upon restart or compaction, a race condition in turn_store_for_context that can instantiate duplicate row stores for the same user, and an O(N) read complexity bottleneck in read_delta_log that should be documented for future optimization.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| fn rebuild_tallies_from_reservations(state: &mut ResourceState) { | ||
| state.reserved_by_account.clear(); | ||
| state.usage_by_account.clear(); | ||
| for record in state.reservations.values() { |
There was a problem hiding this comment.
Correctness Issue: Historical usage from previous periods is incorrectly restored upon restart/compaction
In rebuild_tallies_from_reservations, the tallies for usage_by_account are cleared and rebuilt by iterating over all reservations in state.reservations and adding their usage:
state.usage_by_account.entry(account.clone()).or_default().add_assign(&spent);However, this does not take into account the period boundaries or whether a reservation was reconciled in a previous, already-rolled-over budget period.
When the governor restarts or performs compaction, rebuild_tallies_from_reservations will blindly add the usage of all historical reconciled reservations to the current period's usage_by_account. Since the current time now is already within the active period, advance_period_if_rolled_over will not trigger a rollover, meaning the old usage from previous periods will incorrectly inflate the current period's tally. This can lead to premature quota exhaustion and false-positive LimitExceeded errors.
Recommendation:
Only include reconciled reservations in the tally if they fall within the current active budget period for that account (i.e., after the start of the current period as defined by the account's period anchor/bounds), or avoid clearing and rebuilding the historical usage_by_account from the raw reservations map if the snapshot already contains the correct rolled-over tallies.
| if let Some(store) = self | ||
| .row_turn_stores | ||
| .lock() | ||
| .map_err(|_| { | ||
| OperationFailure::new( | ||
| "turn_store_lock_poisoned", | ||
| "turn_store", | ||
| "row turn-store cache lock poisoned", | ||
| ) | ||
| })? | ||
| .get(&key) | ||
| .cloned() | ||
| { | ||
| return Ok(store); | ||
| } | ||
|
|
||
| let view = user_turn_mount_view(&self.run_id, &resource_scope) | ||
| .map_err(|error| OperationFailure::invalid_request("turn_store", error))?; | ||
| let scoped = Arc::new(ScopedFilesystem::with_fixed_view( | ||
| Arc::clone(&self.root), | ||
| view, | ||
| )); | ||
| let store = Arc::new( | ||
| FilesystemTurnStateRowStore::new(scoped).with_limits(self.turn_state_limits), | ||
| ) as Arc<dyn StressTurnStore>; | ||
| self.row_turn_stores | ||
| .lock() | ||
| .map_err(|_| { | ||
| OperationFailure::new( | ||
| "turn_store_lock_poisoned", | ||
| "turn_store", | ||
| "row turn-store cache lock poisoned", | ||
| ) | ||
| })? | ||
| .insert(key, Arc::clone(&store)); | ||
| Ok(store) |
There was a problem hiding this comment.
Concurrency Bug: Race condition in turn_store_for_context can create duplicate FilesystemTurnStateRowStore instances for the same user
The TurnStateBackend::FilesystemRow branch of turn_store_for_context releases the lock on self.row_turn_stores after checking for an existing store, and then re-acquires it to insert the newly created store. This creates a classic check-then-act race condition. If two concurrent threads call turn_store_for_context for the same user/tenant, both can see a cache miss, initialize separate FilesystemTurnStateRowStore instances, and insert them into the map (with one overwriting the other).
As a result, the two threads will operate on different in-memory hot caches (snapshot_state), completely breaking the in-process cache consistency and leading to stale reads, lost updates, or transaction conflicts.
Recommendation:
Since creating a FilesystemTurnStateRowStore is extremely fast and does not perform any blocking I/O, you should hold the lock for the entire duration of the check-and-insert operation.
let mut stores = self
.row_turn_stores
.lock()
.map_err(|_|
OperationFailure::new(
"turn_store_lock_poisoned",
"turn_store",
"row turn-store cache lock poisoned",
)
)?;
if let Some(store) = stores.get(&key).cloned() {
return Ok(store);
}
let view = user_turn_mount_view(&self.run_id, &resource_scope)
.map_err(|error| OperationFailure::invalid_request("turn_store", error))?;
let scoped = Arc::new(ScopedFilesystem::with_fixed_view(
Arc::clone(&self.root),
view,
));
let store = Arc::new(
FilesystemTurnStateRowStore::new(scoped).with_limits(self.turn_state_limits),
) as Arc<dyn StressTurnStore>;
stores.insert(key, Arc::clone(&store));
Ok(store)References
- Ensure atomicity of read-modify-write operations on shared state by holding a lock across the entire sequence (read, validate, and write). All related operations must share the same lock to prevent race conditions and inconsistent interleaving in concurrent environments.
| async fn read_delta_log(&self) -> Result<Vec<SnapshotDelta>, TurnError> { | ||
| let path = delta_log_path()?; | ||
| let records = match self | ||
| .filesystem | ||
| .tail(&ResourceScope::system(), &path, SeqNo::ZERO) | ||
| .await |
There was a problem hiding this comment.
Performance/Scalability Issue: O(N) complexity in read_delta_log due to full scan from SeqNo::ZERO
The read_delta_log method reads the entire delta log from the beginning of time (SeqNo::ZERO):
let records = match self.filesystem.tail(&ResourceScope::system(), &path, SeqNo::ZERO).awaitSince there is currently no materialized row compactor or truncation mechanism for the delta log, this log will grow indefinitely. As a result, load_snapshot_from_rows (called at startup or when the cache is cleared) will take progressively longer to complete, directly impacting startup latency.
Recommendation:
Because this snapshot-adapter pattern necessitates loading the full state, we should defer immediate optimization of this read path. Instead, please document the requirement for targeted read paths or a background compaction/checkpointing mechanism as a follow-up task.
References
- If a specific design pattern (e.g., snapshot-adapter) necessitates an inefficient implementation (e.g., loading full state), defer optimization and document the requirement for targeted read paths as a follow-up task.
Reborn integration-tier coverageLine coverage (Reborn crates): 28.15% — 49476 / 175731 lines Per-crate breakdown (62 crates, lowest-covered first)
This signal is informational: coverage never gates the PR — not the percentage, not the per-crate holes, not the 0-coverage callout. Exemptions (0 file(s) excluded from the accounting above)No exemptions configured. |
There was a problem hiding this comment.
Actionable comments posted: 13
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
crates/ironclaw_triggers/src/postgres.rs (1)
162-178: 🚀 Performance & Scalability | 🔵 Trivial | ⚡ Quick winRoute
get_triggerthrough the cached query helper
This lookup still callsquery_optdirectly, so it bypassesprepare_cachedwhile the rest of the read paths use the shared wrapper. Switching it over keeps trigger fetches consistent and avoids an extra prepare on a hot path.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/ironclaw_triggers/src/postgres.rs` around lines 162 - 178, The get_trigger read path still bypasses the shared cached query wrapper by calling query_opt directly. Update the get_trigger logic in postgres.rs to route this lookup through the existing cached query helper used by the other read paths, so it benefits from prepare_cached and stays consistent with the rest of the trigger queries. Keep the row_to_record mapping and backend_error handling intact while swapping the direct query call to the shared helper.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/ironclaw_resources/src/filesystem_governor.rs`:
- Around line 692-706: The restoration path in filesystem_governor.rs is
rebuilding v3 tallies from historical reservations via
rebuild_tallies_from_reservations(), which can resurrect expired spend after
rollover. Update the snapshot recovery logic to keep snapshot.state as-is for v3
and only replay deltas after journal_seq in the delta log processing flow around
rebuild_tallies_from_reservations, delta.apply_to, and the records loop. If
legacy repair still needs rebuilding, isolate it behind a separate explicit path
instead of using it during normal snapshot replay.
- Around line 173-193: The resource-governor mutation paths are publishing
in-memory state before the durable journal append is confirmed, then releasing
serialization too early, which can let replay observe operations out of order.
Update the relevant mutation flows around persist_delta(),
write_accounts_from_state(), and the reserve/reconcile/account snapshot handlers
to stage changes first, append the delta while the affected
accounts/reservations are still locked, and only publish the state after the
append succeeds. If needed, funnel all mutations through a single journaled
command pipeline so ordering is preserved across callers. Add a concurrent
reserve/reconcile/restart test that proves replay cannot see B-before-A.
- Line 25: The background compactor/journal diagnostics in filesystem_governor
should use debug! instead of warn! to avoid breaking the REPL/TUI logging
invariant. Update the tracing::warn import and replace the warn! calls in the
internal background paths with debug! in the filesystem governor routines that
emit these messages, including the affected compactor/journal code paths.
- Around line 694-699: Treat FilesystemError::Unsupported in the filesystem.tail
handling as a hard failure instead of returning an empty Vec, because delta-log
tailing cannot be proven safe in that case. Update the match in the tail path to
propagate Unsupported through fs_error(error) like other unexpected errors,
while keeping the NotFound case as an empty result. Use the existing
filesystem.tail and fs_error flow in filesystem_governor to locate the branch
and preserve the current error context propagation.
In `@crates/ironclaw_secrets/src/postgres_store.rs`:
- Around line 691-695: The error mapping in `postgres_error` (and the matching
`serde_to_store_error` path) is leaking raw backend `Display` text into
`SecretStoreError::StoreUnavailable.reason`. Update the fallible paths in
`postgres_store` so they log the original `tokio_postgres::Error` /
`serde_json::Error` internally at debug/error level, but return a fixed
sanitized reason string without backend details. Apply the same pattern to the
store operations that call these helpers (`put`, `read`, `delete`, `lease_once`,
`consume`, `revoke`, `migrate`) and check `FilesystemSecretStore`’s
`fs_to_secret_store_error` for the same treatment.
- Around line 20-24: Extend the secret-store contract coverage to
PostgresSecretStore by adding caller-level tests that exercise the persistence
side effects in lease_once, consume, and revoke. Use the PostgresSecretStore
type and its related methods to verify FOR UPDATE contention behavior and expiry
write-back against a real Postgres-backed setup, following the contract-test
style already expected for secret stores.
In `@crates/ironclaw_threads/src/filesystem_service.rs`:
- Around line 516-531: The thread record update in `FilesystemService` should
not fail immediately on `VersionMismatch` from the `txn.put` for `thread_entry`;
this is a CAS conflict that can happen when another writer advances
`thread.json`. Update the write path to treat this as a retryable conflict in
the same message-write flow, so the outer retry loop re-reads and re-applies the
legacy-counter update instead of returning an error. Use the
`stored.next_sequence` / `thread_version` update block in `FilesystemService` as
the place to handle this conflict before calling `absent_put_error`.
In `@crates/ironclaw_turns/src/filesystem_store/row_store.rs`:
- Around line 744-822: The `apply` flow in `row_store.rs` publishes the updated
snapshot cache before `await_delta_ack` confirms the durable write, allowing a
later caller to build on unconfirmed state. Update `apply` (and the matching
`apply_with_targeted_delta` path referenced in the review) so
`snapshot_state`/`guard` remains protected until the ack succeeds, or add a
fencing/version validation that prevents a write from committing on a baseline
that was not durably confirmed. Ensure the cache is only updated after the
durable delta is acknowledged and clear the cache on any ack failure.
- Around line 391-404: get_run_state and get_loop_checkpoint are still doing
linear scans over the Vec-backed TurnPersistenceSnapshot instead of using the
existing O(1) indexes. Update the lookup path used by
projection::run_state_parts and projection::loop_checkpoint so it reads through
RowSnapshotState/RowSnapshotIndexes or the InMemoryTurnStateStore accessors
(run_record/turn_record) rather than relying on with_cached_snapshot alone. Keep
with_cached_snapshot as the cache entry point, but change the callers in the
read path to use the indexed lookup helpers already exercised by
submit_turn_targeted_delta and apply_with_targeted_delta’s
RunnerLeaseOverlay::Run branch.
- Around line 410-461: The delta-log replay in
load_snapshot_from_rows/replay_deltas is unbounded because it always tails from
SeqNo::ZERO. Update the replay path to start from the current retention floor or
the latest durable checkpoint metadata instead of the beginning, using the
existing RowSnapshotState/TurnPersistenceSnapshot flow to seed the starting
SeqNo. Also make sure read_delta_log, read_run_state_from_durable_rows, and
read_turn_events_from_durable_rows use the same bounded starting point so
durable reads do not scan the full history.
In `@crates/ironclaw_turns/src/memory/mod.rs`:
- Around line 479-506: overlay_runner_lease_record duplicates the same
lease-fencing and stale-heartbeat checks already used by
runner_lease::apply_runner_lease_overlay and run_can_use_external_lease. Extract
a shared predicate/helper over the common record fields (status, runner_id,
lease_token, last_heartbeat_at) and have overlay_runner_lease_record call it
instead of re-implementing the logic, so both paths stay in sync.
In `@harness/latency/lint.sh`:
- Around line 24-33: The temporary file handling in the latency lint script is
using a predictable, shell-PID-based path, which is racy and unsafe. Update the
logic in the lint script to create the scratch file with mktemp, store that path
in a variable, and use the same variable for both rg reads and cleanup. Keep the
rest of the flow intact around the rg checks and the final rm, but ensure the
temp file name is unpredictable and atomically created.
In `@tools/ironclaw_stress/src/user_turn.rs`:
- Around line 757-820: The TurnLifecycleChurn branch in user_turn.rs is
discarding the ClaimedTurnRun returned by claim_next_run() and then calling
complete_run() with the submit-time run_id instead of the claimed state. Update
this path to capture the claimed value from turn_store.claim_next_run in the
existing time_stage call, and use claimed.state.run_id, claimed.runner_id, and
claimed.lease_token when completing the run, matching the other call sites and
preserving the actual lease ownership.
---
Outside diff comments:
In `@crates/ironclaw_triggers/src/postgres.rs`:
- Around line 162-178: The get_trigger read path still bypasses the shared
cached query wrapper by calling query_opt directly. Update the get_trigger logic
in postgres.rs to route this lookup through the existing cached query helper
used by the other read paths, so it benefits from prepare_cached and stays
consistent with the rest of the trigger queries. Keep the row_to_record mapping
and backend_error handling intact while swapping the direct query call to the
shared helper.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: ASSERTIVE
Plan: Pro Plus
Run ID: 30de446b-ce32-48f9-8ec9-fed1a0fe80af
⛔ Files ignored due to path filters (2)
Cargo.lockis excluded by!**/*.lock,!**/Cargo.lockharness/latency/runner/Cargo.lockis excluded by!**/*.lock,!**/Cargo.lock
📒 Files selected for processing (56)
LOG.mdcrates/ironclaw_filesystem/src/backend.rscrates/ironclaw_filesystem/src/libsql.rscrates/ironclaw_filesystem/src/postgres.rscrates/ironclaw_filesystem/src/scoped.rscrates/ironclaw_host_runtime/src/services.rscrates/ironclaw_host_runtime/src/services/builder.rscrates/ironclaw_reborn_composition/src/factory.rscrates/ironclaw_reborn_composition/src/input.rscrates/ironclaw_reborn_composition/src/lib.rscrates/ironclaw_reborn_composition/src/observability/budget.rscrates/ironclaw_reborn_composition/src/outbound/mod.rscrates/ironclaw_reborn_event_store/src/lib.rscrates/ironclaw_resources/Cargo.tomlcrates/ironclaw_resources/src/cas_snapshot.rscrates/ironclaw_resources/src/filesystem_governor.rscrates/ironclaw_resources/src/filesystem_store.rscrates/ironclaw_resources/src/lib.rscrates/ironclaw_resources/tests/resource_governor_contract.rscrates/ironclaw_secrets/Cargo.tomlcrates/ironclaw_secrets/src/filesystem_store.rscrates/ironclaw_secrets/src/lib.rscrates/ironclaw_secrets/src/postgres_store.rscrates/ironclaw_threads/src/filesystem_service.rscrates/ironclaw_triggers/src/libsql.rscrates/ironclaw_triggers/src/postgres.rscrates/ironclaw_turns/src/filesystem_store.rscrates/ironclaw_turns/src/filesystem_store/projection.rscrates/ironclaw_turns/src/filesystem_store/row_store.rscrates/ironclaw_turns/src/filesystem_store/runner_lease.rscrates/ironclaw_turns/src/lib.rscrates/ironclaw_turns/src/memory/mod.rscrates/ironclaw_turns/tests/filesystem_turn_state_contract.rscrates/ironclaw_turns/tests/loop_checkpoint_store_contract.rsdocs/reborn/contracts/resources.mdgoal.mdharness/latency/README.mdharness/latency/lint.shharness/latency/probe.shharness/latency/runner/.gitignoreharness/latency/runner/Cargo.tomlharness/latency/runner/src/main.rsharness/latency/score.shharness/latency/status.shspec.mdtools/ironclaw_stress/Cargo.tomltools/ironclaw_stress/README.mdtools/ironclaw_stress/src/analysis.rstools/ironclaw_stress/src/human.rstools/ironclaw_stress/src/main.rstools/ironclaw_stress/src/process_pressure.rstools/ironclaw_stress/src/report.rstools/ironclaw_stress/src/resource_ops.rstools/ironclaw_stress/src/summary.rstools/ironclaw_stress/src/tests.rstools/ironclaw_stress/src/user_turn.rs
| use ironclaw_filesystem::{FilesystemError, RootFilesystem, ScopedFilesystem, SeqNo}; | ||
| use ironclaw_host_api::{ReservationStatus, ResourceReservationId, ResourceScope, ScopedPath}; | ||
| use serde::{Deserialize, Serialize}; | ||
| use tracing::warn; |
There was a problem hiding this comment.
📐 Maintainability & Code Quality | 🟡 Minor | ⚡ Quick win
Use debug! for background diagnostics.
These compactor/journal messages are internal background diagnostics; warn! can corrupt REPL/TUI paths under the repo logging invariant.
Proposed fix
-use tracing::warn;
+use tracing::debug;
...
- warn!(reason = %error, "resource governor compaction write failed");
+ debug!(reason = %error, "resource governor compaction write failed");
...
- warn!(reason = %error, "resource governor delta journal thread failed to start");
+ debug!(reason = %error, "resource governor delta journal thread failed to start");As per path instructions, “REPL/TUI logging: info!/warn! corrupt the terminal UI — internal diagnostics use debug!.”
Also applies to: 145-155, 559-563
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/ironclaw_resources/src/filesystem_governor.rs` at line 25, The
background compactor/journal diagnostics in filesystem_governor should use
debug! instead of warn! to avoid breaking the REPL/TUI logging invariant. Update
the tracing::warn import and replace the warn! calls in the internal background
paths with debug! in the filesystem governor routines that emit these messages,
including the affected compactor/journal code paths.
Source: Path instructions
| let (tally, changed) = { | ||
| let mut locked = authority.lock_accounts(std::slice::from_ref(account))?; | ||
| let before = locked.account_parts(account); | ||
| let mut state = | ||
| locked.state_for_accounts(std::slice::from_ref(account), HashMap::new()); | ||
| advance_period_if_rolled_over(&mut state, account, now); | ||
| let tally = state | ||
| .reserved_by_account | ||
| .get(account) | ||
| .cloned() | ||
| .unwrap_or_default(); | ||
| locked.write_accounts_from_state(std::slice::from_ref(account), &state); | ||
| let after = locked.account_parts(account); | ||
| (tally, before != after) | ||
| }; | ||
| if changed { | ||
| let delta = ResourceGovernorDelta::AccountSnapshot { | ||
| account: account.clone(), | ||
| at: now, | ||
| }; | ||
| if let Err(error) = self.persist_delta(&authority, delta) { |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🔴 Critical | 🏗️ Heavy lift
Serialize authority mutation with the durable journal append.
These paths publish in-memory changes before the delta append is acknowledged, then release the account/reservation locks before persist_delta(). A second operation can observe A’s state, append B first, and make restart replay see B-before-A—for example Reconcile before the corresponding Reserve.
Stage the state, append while the affected accounts/reservation are still serialized, then publish the state only after the append succeeds; or route all mutations through one journaled command pipeline. Add a concurrent reserve/reconcile/restart caller test. This violates the resource-governor fail-closed durability contract documented for production persistence.
Also applies to: 204-224, 244-257, 282-323, 345-383, 400-437, 454-469
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/ironclaw_resources/src/filesystem_governor.rs` around lines 173 - 193,
The resource-governor mutation paths are publishing in-memory state before the
durable journal append is confirmed, then releasing serialization too early,
which can let replay observe operations out of order. Update the relevant
mutation flows around persist_delta(), write_accounts_from_state(), and the
reserve/reconcile/account snapshot handlers to stage changes first, append the
delta while the affected accounts/reservations are still locked, and only
publish the state after the append succeeds. If needed, funnel all mutations
through a single journaled command pipeline so ordering is preserved across
callers. Add a concurrent reserve/reconcile/restart test that proves replay
cannot see B-before-A.
| rebuild_tallies_from_reservations(&mut state); | ||
| let path = delta_log_path()?; | ||
| let records = match filesystem.tail(&ResourceScope::system(), &path, from).await { | ||
| Ok(records) => records, | ||
| Err(FilesystemError::NotFound { .. }) | Err(FilesystemError::Unsupported { .. }) => { | ||
| Vec::new() | ||
| } | ||
| Err(error) => return Err(fs_error(error)), | ||
| }; | ||
| let mut latest = from; | ||
| for record in records { | ||
| latest = record.seq; | ||
| let delta: ResourceGovernorDelta = serde_json::from_slice(&record.payload) | ||
| .map_err(|error| storage_error(format!("decode resource governor delta: {error}")))?; | ||
| delta.apply_to(&mut state)?; |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick win
Do not rebuild v3 period ledgers from historical reservations.
rebuild_tallies_from_reservations() clears the snapshot’s authoritative ledgers and re-adds every reconciled reservation. After a periodic budget rolls over and compaction stores the cleared ledger with an advanced cursor, restart will resurrect old spend because ReservationRecord has no period/reconcile timestamp to filter by.
For v3 snapshots, replay from snapshot.state as-is and only apply deltas after journal_seq; reserve rebuilds for an explicit legacy repair path if still needed.
Suggested direction
async fn replay_journal<F>(
filesystem: Arc<ScopedFilesystem<F>>,
mut state: ResourceState,
from: SeqNo,
) -> Result<(ResourceState, SeqNo), ResourceError>
where
F: RootFilesystem,
{
- rebuild_tallies_from_reservations(&mut state);
let path = delta_log_path()?;Also applies to: 711-736
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/ironclaw_resources/src/filesystem_governor.rs` around lines 692 - 706,
The restoration path in filesystem_governor.rs is rebuilding v3 tallies from
historical reservations via rebuild_tallies_from_reservations(), which can
resurrect expired spend after rollover. Update the snapshot recovery logic to
keep snapshot.state as-is for v3 and only replay deltas after journal_seq in the
delta log processing flow around rebuild_tallies_from_reservations,
delta.apply_to, and the records loop. If legacy repair still needs rebuilding,
isolate it behind a separate explicit path instead of using it during normal
snapshot replay.
| let records = match filesystem.tail(&ResourceScope::system(), &path, from).await { | ||
| Ok(records) => records, | ||
| Err(FilesystemError::NotFound { .. }) | Err(FilesystemError::Unsupported { .. }) => { | ||
| Vec::new() | ||
| } | ||
| Err(error) => return Err(fs_error(error)), |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🟠 Major | ⚡ Quick win
Fail closed when delta-log tail is unsupported.
Unsupported means recovery cannot prove whether deltas after the snapshot cursor exist. Treating it as an empty log can admit costed work from stale quota state.
Proposed fix
let records = match filesystem.tail(&ResourceScope::system(), &path, from).await {
Ok(records) => records,
- Err(FilesystemError::NotFound { .. }) | Err(FilesystemError::Unsupported { .. }) => {
- Vec::new()
- }
+ Err(FilesystemError::NotFound { .. }) => Vec::new(),
Err(error) => return Err(fs_error(error)),
};As per path instructions, “Fail loud: flag silent-failure patterns … Errors propagate with ? into thiserror types with context.”
📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| let records = match filesystem.tail(&ResourceScope::system(), &path, from).await { | |
| Ok(records) => records, | |
| Err(FilesystemError::NotFound { .. }) | Err(FilesystemError::Unsupported { .. }) => { | |
| Vec::new() | |
| } | |
| Err(error) => return Err(fs_error(error)), | |
| let records = match filesystem.tail(&ResourceScope::system(), &path, from).await { | |
| Ok(records) => records, | |
| Err(FilesystemError::NotFound { .. }) => Vec::new(), | |
| Err(error) => return Err(fs_error(error)), |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/ironclaw_resources/src/filesystem_governor.rs` around lines 694 - 699,
Treat FilesystemError::Unsupported in the filesystem.tail handling as a hard
failure instead of returning an empty Vec, because delta-log tailing cannot be
proven safe in that case. Update the match in the tail path to propagate
Unsupported through fs_error(error) like other unexpected errors, while keeping
the NotFound case as an empty result. Use the existing filesystem.tail and
fs_error flow in filesystem_governor to locate the branch and preserve the
current error context propagation.
Source: Path instructions
| pub struct PostgresSecretStore { | ||
| pool: Pool, | ||
| crypto: Arc<SecretsCrypto>, | ||
| lease_ttl: Duration, | ||
| } |
There was a problem hiding this comment.
🚀 Performance & Scalability | 🟡 Minor | ⚡ Quick win
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
fd . crates/ironclaw_secrets -e rs | xargs rg -ln "PostgresSecretStore"Repository: nearai/ironclaw
Length of output: 234
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '== files ==\n'
git ls-files crates/ironclaw_secrets | sed 's#^`#-` #'
printf '\n== outline: postgres_store.rs ==\n'
ast-grep outline crates/ironclaw_secrets/src/postgres_store.rs --view expanded || true
printf '\n== outline: lib.rs ==\n'
ast-grep outline crates/ironclaw_secrets/src/lib.rs --view expanded || true
printf '\n== tests mentioning PostgresSecretStore ==\n'
rg -n "PostgresSecretStore|SecretStore|lease|FOR UPDATE|expiry" crates/ironclaw_secrets -g '*test*' -g '*.rs' || true
printf '\n== file sizes ==\n'
wc -l crates/ironclaw_secrets/src/postgres_store.rs crates/ironclaw_secrets/src/lib.rsRepository: nearai/ironclaw
Length of output: 50371
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '== postgres store references in tests ==\n'
rg -n "PostgresSecretStore|run_migrations|with_lease_ttl|lease_once\\(|consume\\(|revoke\\(" crates/ironclaw_secrets/tests crates/ironclaw_secrets/src -g '*.rs' || true
printf '\n== secret_store_contract.rs (relevant sections) ==\n'
sed -n '1,260p' crates/ironclaw_secrets/tests/secret_store_contract.rs
printf '\n== boundary_contract.rs (relevant sections) ==\n'
sed -n '1,220p' crates/ironclaw_secrets/tests/boundary_contract.rs
printf '\n== AGENTS and CLAUDE for crates/ironclaw_secrets ==\n'
sed -n '1,220p' crates/ironclaw_secrets/AGENTS.md
printf '\n---\n'
sed -n '1,220p' crates/ironclaw_secrets/CLAUDE.mdRepository: nearai/ironclaw
Length of output: 23406
Add Postgres-backed contract tests
crates/ironclaw_secrets/AGENTS.md asks for caller-level tests when a helper gates persistence side effects. Extend the secret-store contract coverage to PostgresSecretStore, especially lease_once/consume/revoke, FOR UPDATE contention, and expiry write-back.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/ironclaw_secrets/src/postgres_store.rs` around lines 20 - 24, Extend
the secret-store contract coverage to PostgresSecretStore by adding caller-level
tests that exercise the persistence side effects in lease_once, consume, and
revoke. Use the PostgresSecretStore type and its related methods to verify FOR
UPDATE contention behavior and expiry write-back against a real Postgres-backed
setup, following the contract-test style already expected for secret stores.
Source: Path instructions
| async fn load_snapshot_from_rows(&self) -> Result<RowSnapshotState, TurnError> { | ||
| let meta = self.read_meta().await?; | ||
| let turns = self.read_row_collection(TURN_ROWS).await?; | ||
| let runs = self.read_row_collection(RUN_ROWS).await?; | ||
| let active_locks = self.read_row_collection(ACTIVE_LOCK_ROWS).await?; | ||
| let checkpoints = self.read_row_collection(CHECKPOINT_ROWS).await?; | ||
| let loop_checkpoints = self.read_row_collection(LOOP_CHECKPOINT_ROWS).await?; | ||
| let idempotency_records = self.read_row_collection(IDEMPOTENCY_ROWS).await?; | ||
| let events = self.read_row_collection(EVENT_ROWS).await?; | ||
| let admission_reservations = self.read_row_collection(ADMISSION_RESERVATION_ROWS).await?; | ||
| let spawn_tree_reservations = self | ||
| .read_row_collection(SPAWN_TREE_RESERVATION_ROWS) | ||
| .await?; | ||
|
|
||
| let mut snapshot = TurnPersistenceSnapshot { | ||
| turns, | ||
| runs, | ||
| active_locks, | ||
| checkpoints, | ||
| loop_checkpoints, | ||
| idempotency_records, | ||
| events, | ||
| event_retention_floor: meta.event_retention_floor, | ||
| admission_reservations, | ||
| spawn_tree_reservations, | ||
| }; | ||
| self.replay_deltas(&mut snapshot).await?; | ||
| let snapshot = row_store_hot_cache_snapshot(snapshot, self.limits); | ||
| let store = self.build_in_memory_store(snapshot)?; | ||
| let snapshot = store.persistence_snapshot(); | ||
| RowSnapshotState::new(snapshot, Arc::new(store)) | ||
| } | ||
|
|
||
| async fn replay_deltas(&self, snapshot: &mut TurnPersistenceSnapshot) -> Result<(), TurnError> { | ||
| let path = delta_log_path()?; | ||
| let records = match self | ||
| .filesystem | ||
| .tail(&ResourceScope::system(), &path, SeqNo::ZERO) | ||
| .await | ||
| { | ||
| Ok(records) => records, | ||
| Err(FilesystemError::NotFound { .. }) | Err(FilesystemError::Unsupported { .. }) => { | ||
| Vec::new() | ||
| } | ||
| Err(error) => return Err(fs_error(error)), | ||
| }; | ||
| for record in records { | ||
| let delta: SnapshotDelta = deserialize_row(&record.payload, "turn-state delta")?; | ||
| apply_delta(snapshot, delta)?; | ||
| } | ||
| Ok(()) | ||
| } |
There was a problem hiding this comment.
🚀 Performance & Scalability | 🟠 Major | 🏗️ Heavy lift
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Map the file and inspect the relevant regions first.
ast-grep outline crates/ironclaw_turns/src/filesystem_store/row_store.rs --view expanded || true
echo "---- lines around load/replay ----"
sed -n '360,560p' crates/ironclaw_turns/src/filesystem_store/row_store.rs | cat -n
echo "---- search for meta_path writes and delta-log truncation/compaction ----"
rg -n "meta_path\\(|event_retention_floor|write_meta|read_meta|delta_log|tail\\(&ResourceScope::system\\(\\), &path, SeqNo::ZERO|SeqNo::ZERO|truncate|compact|retention_floor" crates/ironclaw_turns/src/filesystem_store/row_store.rs crates/ironclaw_turns/src/filesystem_store -SRepository: nearai/ironclaw
Length of output: 28437
🏁 Script executed:
#!/bin/bash
set -euo pipefail
echo "---- write path around the delta journal ----"
sed -n '240,360p' crates/ironclaw_turns/src/filesystem_store/row_store.rs | cat -n
sed -n '820,900p' crates/ironclaw_turns/src/filesystem_store/row_store.rs | cat -n
echo "---- delta-log / meta-path writes across the crate ----"
rg -n "meta_path\\(|delta_log_path\\(|META_FILE|DELTA_LOG|write\\(|append\\(|truncate\\(|compact\\(|remove\\(|delete\\(" crates/ironclaw_turns/src/filesystem_store -S
echo "---- broader repo search for meta_path writes ----"
rg -n "meta_path\\(" crates src tests -SRepository: nearai/ironclaw
Length of output: 11414
🏁 Script executed:
#!/bin/bash
set -euo pipefail
echo "---- durable-delta shaping and retention-floor handling ----"
sed -n '2128,2410p' crates/ironclaw_turns/src/filesystem_store/row_store.rs | cat -n
echo "---- durable query paths ----"
sed -n '540,640p' crates/ironclaw_turns/src/filesystem_store/row_store.rs | cat -n
echo "---- all row_store references to delta-log / meta write APIs ----"
rg -n "\.append_batch\(|\.append\(|\.truncate\(|\.delete\(|\.remove\(|meta_path\(|delta_log_path\(|read_meta\(|read_delta_log\(" crates/ironclaw_turns/src/filesystem_store/row_store.rs -SRepository: nearai/ironclaw
Length of output: 16914
Delta-log replay is unbounded on cold load and durable reads
load_snapshot_from_rows, read_delta_log, read_run_state_from_durable_rows, and read_turn_events_from_durable_rows all tail deltas/log from SeqNo::ZERO. row_store_durable_delta() strips deletes and retention-floor updates, and meta_path() is never written in this file, so replay cost grows with the full write history of the store.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/ironclaw_turns/src/filesystem_store/row_store.rs` around lines 410 -
461, The delta-log replay in load_snapshot_from_rows/replay_deltas is unbounded
because it always tails from SeqNo::ZERO. Update the replay path to start from
the current retention floor or the latest durable checkpoint metadata instead of
the beginning, using the existing RowSnapshotState/TurnPersistenceSnapshot flow
to seed the starting SeqNo. Also make sure read_delta_log,
read_run_state_from_durable_rows, and read_turn_events_from_durable_rows use the
same bounded starting point so durable reads do not scan the full history.
| async fn apply<T, A, Fut>( | ||
| &self, | ||
| overlay: RunnerLeaseOverlay, | ||
| mut apply: A, | ||
| ) -> Result<T, TurnError> | ||
| where | ||
| A: FnMut(Arc<InMemoryTurnStateStore>) -> Fut + Send, | ||
| Fut: std::future::Future<Output = Result<T, TurnError>> + Send, | ||
| T: Send, | ||
| { | ||
| let critical = async { | ||
| let mut guard = self.snapshot_state.lock().await; | ||
| if guard.is_none() { | ||
| *guard = Some(self.load_snapshot_from_rows().await?); | ||
| } | ||
| let store = match (overlay, guard.as_ref()) { | ||
| (RunnerLeaseOverlay::None, Some(state)) => Arc::clone(&state.store), | ||
| (_, Some(state)) => { | ||
| let snapshot = state.snapshot.clone(); | ||
| let (overlaid_snapshot, _) = self | ||
| .runner_lease_store() | ||
| .overlay((snapshot, None), overlay) | ||
| .await?; | ||
| Arc::new(self.build_in_memory_store(overlaid_snapshot)?) | ||
| } | ||
| (_, None) => unreachable!("row snapshot cache is initialized above"), | ||
| }; | ||
| let baseline = guard | ||
| .as_ref() | ||
| .map(|state| state.snapshot.clone()) | ||
| .unwrap_or_default(); | ||
| let outcome = apply(Arc::clone(&store)).await; | ||
| let mut new_snapshot = store.persistence_snapshot(); | ||
| preserve_loop_checkpoints(&baseline, &mut new_snapshot); | ||
| let value = match outcome { | ||
| Ok(value) => value, | ||
| Err(error) => { | ||
| *guard = None; | ||
| return Err(error); | ||
| } | ||
| }; | ||
| if new_snapshot == baseline { | ||
| return Ok((None, value)); | ||
| } | ||
|
|
||
| let delta = match snapshot_delta(&baseline, &new_snapshot) { | ||
| Ok(delta) => delta, | ||
| Err(RowPersistError::Turn(error)) => { | ||
| *guard = None; | ||
| return Err(error); | ||
| } | ||
| }; | ||
| let persist_delta = row_store_durable_delta(delta); | ||
| let ack = match self.enqueue_delta(persist_delta) { | ||
| Ok(ack) => ack, | ||
| Err(RowPersistError::Turn(error)) => { | ||
| *guard = None; | ||
| return Err(error); | ||
| } | ||
| }; | ||
| *guard = Some(RowSnapshotState::new(new_snapshot, store)?); | ||
| Ok((ack, value)) | ||
| }; | ||
|
|
||
| let (ack, value) = match tokio::time::timeout(self.apply_timeout, critical).await { | ||
| Ok(result) => result?, | ||
| Err(_) => { | ||
| self.clear_snapshot_cache().await; | ||
| return Err(TurnError::Unavailable { | ||
| reason: "turn state row-store apply timed out".to_string(), | ||
| }); | ||
| } | ||
| }; | ||
| if let Err(error) = self.await_delta_ack(ack).await { | ||
| self.clear_snapshot_cache().await; | ||
| return Err(error.into_turn()); | ||
| } | ||
| Ok(value) | ||
| } |
There was a problem hiding this comment.
🗄️ Data Integrity & Integration | 🔴 Critical | 🏗️ Heavy lift
Cache is published before the durable delta is acknowledged — a concurrent writer can build on unconfirmed state.
In both apply and apply_with_targeted_delta, guard (the snapshot_state lock) is local to the critical async block and is dropped when critical returns — i.e. *guard = Some(new_snapshot) (or state.store = store) is visible to the next caller before self.await_delta_ack(ack).await runs. If a second write (B) acquires the lock in that window, it computes its delta against A's not-yet-durable snapshot. If A's delta later fails to persist, A calls clear_snapshot_cache(), but B's delta — computed on top of A's phantom state — may already be enqueued/acked in a separate journal batch and persisted durably. The append log then contains B without the A it logically depends on: permanent corruption on the next replay, and a direct violation of goal.md's own hard-fail bar ("skipped durable writes, lost state transitions"). This is squarely in the blast radius of the new turn-lifecycle-churn stress scenario (concurrent submit/claim/complete against one shared store instance).
Hold the lock across the ack (accepting serialized durable writes per store), or add a fencing/version check that rejects a write whose baseline was never confirmed durable.
🔒 Sketch: keep the guard held until the ack resolves
- let (ack, value) = match tokio::time::timeout(self.apply_timeout, critical).await {
- Ok(result) => result?,
- Err(_) => {
- self.clear_snapshot_cache().await;
- return Err(TurnError::Unavailable {
- reason: "turn state row-store apply timed out".to_string(),
- });
- }
- };
- if let Err(error) = self.await_delta_ack(ack).await {
- self.clear_snapshot_cache().await;
- return Err(error.into_turn());
- }
- Ok(value)
+ // Move the ack-await inside `critical`, still holding `guard`, so no
+ // other writer can observe `new_snapshot` before it is durably
+ // confirmed. On ack failure, do not commit `*guard` at all.Also applies to: 852-940
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/ironclaw_turns/src/filesystem_store/row_store.rs` around lines 744 -
822, The `apply` flow in `row_store.rs` publishes the updated snapshot cache
before `await_delta_ack` confirms the durable write, allowing a later caller to
build on unconfirmed state. Update `apply` (and the matching
`apply_with_targeted_delta` path referenced in the review) so
`snapshot_state`/`guard` remains protected until the ack succeeds, or add a
fencing/version validation that prevents a write from committing on a baseline
that was not durably confirmed. Ensure the cache is only updated after the
durable delta is acknowledged and clear the cache on any ack failure.
| pub(crate) fn overlay_runner_lease_record( | ||
| &self, | ||
| overlaid: TurnRunRecord, | ||
| ) -> Result<(), TurnError> { | ||
| let mut inner = self.lock_inner()?; | ||
| let Some(record) = inner.records.get_mut(&overlaid.run_id) else { | ||
| return Err(TurnError::ScopeNotFound); | ||
| }; | ||
| if !matches!( | ||
| record.status.get(), | ||
| TurnStatus::Running | TurnStatus::CancelRequested | ||
| ) || record.runner_id != overlaid.runner_id | ||
| || record.lease_token != overlaid.lease_token | ||
| || record.runner_id.is_none() | ||
| || record.lease_token.is_none() | ||
| { | ||
| return Ok(()); | ||
| } | ||
| if let (Some(current), Some(incoming)) = | ||
| (record.last_heartbeat_at, overlaid.last_heartbeat_at) | ||
| && incoming < current | ||
| { | ||
| return Ok(()); | ||
| } | ||
| record.last_heartbeat_at = overlaid.last_heartbeat_at; | ||
| record.lease_expires_at = overlaid.lease_expires_at; | ||
| Ok(()) | ||
| } |
There was a problem hiding this comment.
📐 Maintainability & Code Quality | 🔵 Trivial | ⚡ Quick win
Duplicates runner_lease.rs::apply_runner_lease_overlay staleness logic.
overlay_runner_lease_record re-implements the same "status must be Running/CancelRequested + runner/lease-token match + reject stale heartbeat" check that already lives in runner_lease::apply_runner_lease_overlay / run_can_use_external_lease. Two independent copies of a concurrency-safety fencing check will drift silently. Extract a shared predicate (operating on the common fields) both sites call.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/ironclaw_turns/src/memory/mod.rs` around lines 479 - 506,
overlay_runner_lease_record duplicates the same lease-fencing and
stale-heartbeat checks already used by runner_lease::apply_runner_lease_overlay
and run_can_use_external_lease. Extract a shared predicate/helper over the
common record fields (status, runner_id, lease_token, last_heartbeat_at) and
have overlay_runner_lease_record call it instead of re-implementing the logic,
so both paths stay in sync.
| if rg -n "LATENCY_|latency|benchmark|bench" crates src \ | ||
| -g '*.rs' >/tmp/ironclaw-latency-lint.$$ 2>/dev/null; then | ||
| if rg -n "sleep|tokio::time::sleep|std::thread::sleep|mock readiness|fast path|fast-path" \ | ||
| /tmp/ironclaw-latency-lint.$$ >/dev/null 2>&1; then | ||
| rm -f /tmp/ironclaw-latency-lint.$$ | ||
| echo "VOID: constraint violation" | ||
| exit 1 | ||
| fi | ||
| fi | ||
| rm -f /tmp/ironclaw-latency-lint.$$ |
There was a problem hiding this comment.
🔒 Security & Privacy | 🟡 Minor | ⚡ Quick win
Predictable temp file path — use mktemp.
/tmp/ironclaw-latency-lint.$$ is guessable/racy; a local attacker can pre-create/symlink it. Use mktemp for an atomically-created, unpredictable path.
🔒️ Proposed fix
+tmpfile="$(mktemp)"
+trap 'rm -f "$tmpfile"' EXIT
if rg -n "LATENCY_|latency|benchmark|bench" crates src \
- -g '*.rs' >/tmp/ironclaw-latency-lint.$$ 2>/dev/null; then
+ -g '*.rs' >"$tmpfile" 2>/dev/null; then
if rg -n "sleep|tokio::time::sleep|std::thread::sleep|mock readiness|fast path|fast-path" \
- /tmp/ironclaw-latency-lint.$$ >/dev/null 2>&1; then
- rm -f /tmp/ironclaw-latency-lint.$$
+ "$tmpfile" >/dev/null 2>&1; then
echo "VOID: constraint violation"
exit 1
fi
fi
-rm -f /tmp/ironclaw-latency-lint.$$📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| if rg -n "LATENCY_|latency|benchmark|bench" crates src \ | |
| -g '*.rs' >/tmp/ironclaw-latency-lint.$$ 2>/dev/null; then | |
| if rg -n "sleep|tokio::time::sleep|std::thread::sleep|mock readiness|fast path|fast-path" \ | |
| /tmp/ironclaw-latency-lint.$$ >/dev/null 2>&1; then | |
| rm -f /tmp/ironclaw-latency-lint.$$ | |
| echo "VOID: constraint violation" | |
| exit 1 | |
| fi | |
| fi | |
| rm -f /tmp/ironclaw-latency-lint.$$ | |
| tmpfile="$(mktemp)" | |
| trap 'rm -f "$tmpfile"' EXIT | |
| if rg -n "LATENCY_|latency|benchmark|bench" crates src \ | |
| -g '*.rs' >"$tmpfile" 2>/dev/null; then | |
| if rg -n "sleep|tokio::time::sleep|std::thread::sleep|mock readiness|fast path|fast-path" \ | |
| "$tmpfile" >/dev/null 2>&1; then | |
| echo "VOID: constraint violation" | |
| exit 1 | |
| fi | |
| fi |
🧰 Tools
🪛 ast-grep (0.44.1)
[warning] 24-24: Building a temp file path in a world-writable directory from the PID ($$) or `` is predictable and racy: an attacker can pre-create or guess the name and win a symlink/race attack. Use mktemp (e.g. `f=$(mktemp)` or `f=$(mktemp /tmp/myapp.XXXXXX)`) so the kernel atomically creates a unique, unpredictable file.
Context: /tmp/ironclaw-latency-lint.$$
Note: [CWE-377] Insecure Temporary File.
(tmp-file-pid-name-bash)
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@harness/latency/lint.sh` around lines 24 - 33, The temporary file handling in
the latency lint script is using a predictable, shell-PID-based path, which is
racy and unsafe. Update the logic in the lint script to create the scratch file
with mktemp, store that path in a variable, and use the same variable for both
rg reads and cleanup. Keep the rest of the flow intact around the rg checks and
the final rm, but ensure the temp file name is unpredictable and atomically
created.
Source: Linters/SAST tools
| if matches!(args.scenario, Scenario::TurnLifecycleChurn) { | ||
| let operation_ref = turn_operation_ref(args, worker_index, operation_index, 0, 1); | ||
| let SubmitTurnResponse::Accepted { run_id, .. } = time_stage( | ||
| &mut stages.submit_turn, | ||
| turn_coordinator.submit_turn(SubmitTurnRequest { | ||
| scope: context.turn_scope.clone(), | ||
| actor: TurnActor::new(context.user_id.clone()), | ||
| accepted_message_ref: AcceptedMessageRef::new(format!( | ||
| "message:{operation_ref}" | ||
| )) | ||
| .map_err(|error| OperationFailure::invalid_request("submit_turn", error))?, | ||
| source_binding_ref: SourceBindingRef::new(source_binding) | ||
| .map_err(|error| OperationFailure::invalid_request("submit_turn", error))?, | ||
| reply_target_binding_ref: ReplyTargetBindingRef::new(reply_target) | ||
| .map_err(|error| OperationFailure::invalid_request("submit_turn", error))?, | ||
| requested_run_profile: None, | ||
| idempotency_key: IdempotencyKey::new(format!( | ||
| "ironclaw-stress:{operation_ref}" | ||
| )) | ||
| .map_err(|error| OperationFailure::invalid_request("submit_turn", error))?, | ||
| received_at: Utc::now(), | ||
| requested_run_id: None, | ||
| parent_run_id: None, | ||
| subagent_depth: 0, | ||
| spawn_tree_root_run_id: None, | ||
| product_context: None, | ||
| }), | ||
| ) | ||
| .await | ||
| .map_err(|error| turn_failure("submit_turn", error))?; | ||
|
|
||
| let runner_id = TurnRunnerId::new(); | ||
| let lease_token = TurnLeaseToken::new(); | ||
| time_stage( | ||
| &mut stages.claim_run, | ||
| turn_store.claim_next_run(ClaimRunRequest { | ||
| runner_id, | ||
| lease_token, | ||
| scope_filter: Some(context.turn_scope.clone()), | ||
| }), | ||
| ) | ||
| .await | ||
| .map_err(|error| turn_failure("claim_run", error))? | ||
| .ok_or_else(|| { | ||
| OperationFailure::new( | ||
| "turn_claim_miss", | ||
| "claim_run", | ||
| "submitted run was not claimable", | ||
| ) | ||
| })?; | ||
|
|
||
| time_stage( | ||
| &mut stages.complete_run, | ||
| turn_store.complete_run(CompleteRunRequest { | ||
| run_id, | ||
| runner_id, | ||
| lease_token, | ||
| }), | ||
| ) | ||
| .await | ||
| .map_err(|error| turn_failure("complete_run", error))?; | ||
|
|
||
| return Ok(()); | ||
| } |
There was a problem hiding this comment.
🎯 Functional Correctness | 🟠 Major | ⚡ Quick win
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '## user_turn.rs around TurnLifecycleChurn\n'
sed -n '720,840p' tools/ironclaw_stress/src/user_turn.rs
printf '\n## Find claim/complete request and response types\n'
rg -n "struct ClaimRunRequest|struct CompleteRunRequest|enum SubmitTurnResponse|struct ClaimedTurnRun|claim_next_run|complete_run" -S crates tools -g '!target'
printf '\n## Comparable call sites in user_turn.rs\n'
rg -n "claimed\.state\.run_id|claimed\.runner_id|claimed\.lease_token|complete_run\(CompleteRunRequest" tools/ironclaw_stress/src/user_turn.rs -n -A3 -B3Repository: nearai/ironclaw
Length of output: 50372
🏁 Script executed:
#!/bin/bash
set -euo pipefail
sed -n '1,260p' tools/ironclaw_stress/src/user_turn.rsRepository: nearai/ironclaw
Length of output: 10461
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '## runner request/response types\n'
sed -n '1,220p' crates/ironclaw_turns/src/runner.rs
printf '\n## claim_next_run implementation contract in memory store\n'
sed -n '1741,1845p' crates/ironclaw_turns/src/memory/mod.rs
printf '\n## claim_next_run / complete_run in filesystem row store\n'
sed -n '1361,1495p' crates/ironclaw_turns/src/filesystem_store/row_store.rs
printf '\n## TurnLifecycleChurn scenario wiring / operation counts\n'
rg -n "TurnLifecycleChurn|turn-lifecycle-churn|threads_per_owner|operations_per_thread|turns_per_operation" tools/ironclaw_stress/src/user_turn.rs tools/ironclaw_stress/src/*.rs -A4 -B4Repository: nearai/ironclaw
Length of output: 50371
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '## synthetic user turn context / scope derivation\n'
sed -n '1,220p' tools/ironclaw_stress/src/synthetic.rs
printf '\n## turn scope and thread ownership references\n'
rg -n "turn_scope|thread_scope|thread_id|thread_owner|threads_per_owner|scope_filter" tools/ironclaw_stress/src/synthetic.rs tools/ironclaw_stress/src/user_turn.rs tools/ironclaw_stress/src/main.rs -A4 -B4Repository: nearai/ironclaw
Length of output: 49021
🏁 Script executed:
#!/bin/bash
set -euo pipefail
rg -n "fn user_turn_context|turn_scope|thread_scope|thread_owner_user_id|scope_filter" tools/ironclaw_stress/src/synthetic.rs -A30 -B10Repository: nearai/ironclaw
Length of output: 6944
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '## TurnScope definition and resource conversion\n'
rg -n "struct TurnScope|impl TurnScope|to_resource_scope|new_with_owner" crates/ironclaw_turns crates/ironclaw_threads -A40 -B10
printf '\n## user_turn_mount_view / row-store key inputs\n'
rg -n "user_turn_mount_view|row_turn_store_key|to_resource_scope\\(" tools/ironclaw_stress/src/user_turn.rs crates/ironclaw_turns/src -A20 -B10Repository: nearai/ironclaw
Length of output: 50372
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '## TurnScope struct and state-path helpers\n'
rg -n "pub struct TurnScope|impl TurnScope|thread_root_string|turn_root_string|state.json|turns/state" crates/ironclaw_turns crates/ironclaw_threads -A40 -B10
printf '\n## claim/submit scope filtering helpers\n'
rg -n "pop_matching_queued_run|scope_filter|new_with_owner|TurnScope::new_with_owner" crates/ironclaw_turns/src -A30 -B10Repository: nearai/ironclaw
Length of output: 50374
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '## TurnScope definition\n'
sed -n '1,120p' crates/ironclaw_turns/src/contract.rs
printf '\n## partitioned_worker_index helper\n'
sed -n '220,320p' tools/ironclaw_stress/src/synthetic.rsRepository: nearai/ironclaw
Length of output: 258
🏁 Script executed:
#!/bin/bash
set -euo pipefail
sed -n '220,320p' tools/ironclaw_stress/src/synthetic.rsRepository: nearai/ironclaw
Length of output: 336
🏁 Script executed:
#!/bin/bash
set -euo pipefail
printf '## validate_args checks for active_thread_count / threads_per_owner / TurnLifecycleChurn\n'
rg -n "active_thread_count|threads_per_owner|TurnLifecycleChurn|turn-lifecycle-churn|requires --processes 1|one thread per owner" tools/ironclaw_stress/src/main.rs tools/ironclaw_stress/src/tests.rs -A6 -B6
printf '\n## user_turn_task loop / concurrency shape\n'
rg -n "run_user_turn_tasks|operations_per_thread|for operation_index|worker_index" tools/ironclaw_stress/src/user_turn.rs tools/ironclaw_stress/src/main.rs -A10 -B10Repository: nearai/ironclaw
Length of output: 50371
Use the claimed run here. claim_next_run() returns a ClaimedTurnRun, but this branch drops it and completes the submit-time run_id. In the shared-hot-thread mode this workload supports (--active-thread-count 1), multiple workers can queue runs under the same turn_scope, so complete_run() can hit the wrong lease and leave the actual claim running. Use claimed.state.run_id / claimed.runner_id / claimed.lease_token like the other call sites.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@tools/ironclaw_stress/src/user_turn.rs` around lines 757 - 820, The
TurnLifecycleChurn branch in user_turn.rs is discarding the ClaimedTurnRun
returned by claim_next_run() and then calling complete_run() with the
submit-time run_id instead of the claimed state. Update this path to capture the
claimed value from turn_store.claim_next_run in the existing time_stage call,
and use claimed.state.run_id, claimed.runner_id, and claimed.lease_token when
completing the run, matching the other call sites and preserving the actual
lease ownership.
Human comment: I will chop this up.
Summary
This PR is the hosted-single-tenant Postgres latency cycle. It moves the hot paths away from blob-style persistence and per-reservation Postgres transactions toward RootFilesystem-backed append/row stores with in-process authority where the deployment contract allows it.
Main changes:
append_batch.turns,runs,events) and bounded hot-cache state; terminal runs/events are evicted from memory without deleting durable rows.ironclaw_stress --scenario turn-lifecycle-churnfor c100 submit/claim/complete churn and RSS/cache validation.Journey / starting point
Representative starting signals from the loop:
Main bottlenecks found:
Ending point
Latest local gates on this branch:
Notes:
Verification
Ran locally:
cargo check -p ironclaw_turnscargo test -p ironclaw_turns --test filesystem_turn_state_contractcargo test -p ironclaw_stresscargo test -p ironclaw_filesystem --features libsql,postgres --test db_root_filesystem_contractcargo test -p ironclaw_reborn_composition --features libsql,postgres --test libsql_substrate --test postgres_substrateharness/latency/score.sh --devplus focused c1/c4 reruns during the cycleironclaw_stressc100 mixed-flow and turn-lifecycle-churn runs listed above