fix: serialize concurrent inbound-message writes per thread - #6096
henrypark133 wants to merge 2 commits into
Conversation
On the non-transactional fallback path used by the production libSQL backend, accept_inbound_message reserved a sequence number and wrote the message record as two separate, unsynchronized steps. Two concurrent messages to the same thread could commit out of call order, so a later message could become durably visible (and get acted on) before an earlier one. Add a per-(scope, thread_id) in-process lock around just the reserve + write step (not the best-effort index/cache work after it), sharing the weak-referenced lock-registry primitive already used by thread_index_load_lock. Split write_new_message into put_new_message_record (the ordering-sensitive CAS put) and finish_new_message_indexes (the rest) so the lock covers only what needs it. Closes #6047 Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
|
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configurationConfiguration used: Path: .coderabbit.yaml Review profile: ASSERTIVE Plan: Pro Plus Run ID: 📒 Files selected for processing (2)
📝 WalkthroughSummary by CodeRabbit
WalkthroughChangesUnsupported filesystem backends now serialize sequence reservation and message-record visibility per thread. Shared weak-keyed locks and injective thread keys support the path, while index/cache work remains outside the lock. Tests cover ordering, failures, lock scope, and reclamation. Inbound message ordering
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant Client
participant FilesystemSessionThreadService
participant ThreadWriteLock
participant FilesystemBackend
participant IndexAndCache
Client->>FilesystemSessionThreadService: accept_inbound_message
FilesystemSessionThreadService->>ThreadWriteLock: acquire thread-scoped lock
ThreadWriteLock->>FilesystemBackend: reserve sequence
ThreadWriteLock->>FilesystemBackend: CAS-put message record
FilesystemBackend-->>FilesystemSessionThreadService: write result
FilesystemSessionThreadService->>IndexAndCache: finish indexes and invalidate cache
FilesystemSessionThreadService-->>Client: acceptance result
Possibly related issues
Possibly related PRs
Suggested reviewers: 🚥 Pre-merge checks | ✅ 4✅ Passed checks (4 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 addresses Issue #6047 by introducing a process-local, weak-keyed async lock (inbound_message_write_locks) to serialize the sequence reservation and message write steps for non-transactional backends, preventing out-of-order message visibility. It refactors message writing to isolate the ordering-sensitive step and adds robust regression tests. The review feedback suggests improving the robustness of thread_scope_key by using a length-prefixed encoding to prevent separator collisions, and extending the per-thread write lock to other message types (such as assistant drafts and tool results) to ensure complete write ordering across all roles.
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 thread_scope_key(scope: &ThreadScope, thread_id: &ThreadId) -> String { | ||
| format!( | ||
| "{}:{}", | ||
| scope.tenant_id.as_str(), | ||
| thread_root_string(scope, thread_id) | ||
| ) | ||
| } |
There was a problem hiding this comment.
The thread_scope_key function constructs a cache and lock key by simply concatenating scope.tenant_id and thread_root_string with a : delimiter. This introduces a separator collision risk (e.g., if a tenant ID contains a : or if there is any overlap in the dynamic string components), which could lead to cross-tenant cache pollution or lock contention.
To eliminate this risk and ensure complete injectivity, use a length-prefixed encoding for each component as recommended by the general rules.
fn thread_scope_key(scope: &ThreadScope, thread_id: &ThreadId) -> String {
let tenant = scope.tenant_id.as_str();
let root = thread_root_string(scope, thread_id);
format!("{}:{}:{}:{}", tenant.len(), tenant, root.len(), root)
}References
- When generating deterministic identifiers or hashes from multiple string components, use an injective length-prefixed encoding rather than simple concatenation with a delimiter to eliminate the risk of separator-collision attacks or accidental collisions.
There was a problem hiding this comment.
Confirmed real — validate_scope_id forbids / and \\ but not :. Fixed in a86c579: thread_scope_key now length-prefixes each component. One nuance worth recording: I verified the specific collision isn't reachable through the current concrete ThreadId/TenantId shapes (thread_root_string always starts with a fixed /threads prefix, and / is forbidden in every id component, so the delimiter can't actually be repositioned today) — so this was defense-in-depth against a fragile implicit invariant rather than a live exploit, but the fix is correct either way and a regression test now proves the new encoding's injectivity directly rather than relying on that invariant.
| async fn write_new_message( | ||
| &self, | ||
| scope: &ThreadScope, | ||
| thread_id: &ThreadId, | ||
| message: &ThreadMessageRecord, | ||
| description: &'static str, | ||
| ) -> Result<(), SessionThreadError> { | ||
| self.put_new_message_record(scope, thread_id, message, description) | ||
| .await?; | ||
| self.finish_new_message_indexes(scope, thread_id, message) | ||
| .await | ||
| } |
There was a problem hiding this comment.
While this PR successfully serializes user messages in accept_inbound_message, other message types (such as assistant drafts, finalized assistant messages, and tool result references) still use the un-serialized reserve_sequence followed by write_new_message pattern.
If a user message is sent concurrently with an assistant response or tool result write on the same thread, they can still race and be written/displayed out of chronological order. To ensure complete thread-level write ordering and prevent cross-role message reordering, consider acquiring the per-thread write lock around the reserve-then-write sequence for assistant and tool messages as well.
There was a problem hiding this comment.
Confirmed real and in-scope for a follow-up, not this PR — #6047 is specifically about user messages (that's the reported bug and what the regression tests prove), and append_assistant_draft/append_finalized_assistant_message/append_tool_result_reference do share the same unserialized two-step pattern. Filed #6101 to extend the same lock mechanism to those call sites.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: d236be5a6c
ℹ️ About Codex in GitHub
Your team has set up Codex to review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
Codex can also answer questions or update the PR. Try commenting "@codex address that feedback".
| self.finish_new_message_indexes(&scope, &thread_id, &message) | ||
| .await?; |
There was a problem hiding this comment.
Keep one-shot cache updates inside the thread write lock
On non-transactional backends this runs after _write_guard is dropped, but finish_new_message_indexes invalidates the one-shot context cache and the sequence == 1 branch below may re-seed it. If message A (seq 1) stalls in the sequence-index/idempotency work after dropping the lock, message B (seq 2) can complete and invalidate the cache, then A resumes and seeds a one-message context; the next load_context_window takes that stale cache and omits B even though B is already durable. Keep the cache seed/invalidate (or the full finish) serialized with the per-thread write order.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Verified this is real, but it predates this PR — the if sequence == 1 { seed } else { invalidate } block already existed on origin/main unchanged, outside any lock/transaction, for both the transactional and non-transactional write paths. The #6047 fix's lock-narrowing (from the first review round) didn't open this window since that decision point was never covered by any lock before or after. Filed #6100 as its own follow-up rather than folding it into this PR's scope.
There was a problem hiding this comment.
Actionable comments posted: 3
🤖 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_threads/src/filesystem_service.rs`:
- Around line 3096-3112: Amortize stale-lock cleanup in weak_keyed_lock instead
of calling retain on every lock acquisition. Preserve reuse and insertion of
keyed locks, but trigger periodic cleanup based on an established counter or
time interval so accept_inbound_message avoids an O(n) map scan on each call.
- Around line 3096-3112: Fix weak_keyed_lock in
crates/ironclaw_threads/src/filesystem_service.rs:3096-3112 to recover the
poisoned Mutex guard via into_inner instead of returning a fresh unshared lock,
preserving shared mutual exclusion for all keys. No direct change is needed at
thread_index_load_lock in
crates/ironclaw_threads/src/filesystem_service/thread_index.rs:497-499; it is
corrected automatically by the shared helper fix.
- Around line 193-219: Add the repository-required arch-exempt annotation for
the process-local lock used by inbound_message_write_locks and _write_guard
across reserve_sequence and put_new_message_record I/O. Include a valid plan or
follow-up tracking link documenting the intended durable/CAS-based replacement,
while preserving the existing lock scope and behavior.
🪄 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: 40187cfc-68e9-4f41-82ab-a7a33268afc6
📒 Files selected for processing (3)
crates/ironclaw_threads/src/filesystem_service.rscrates/ironclaw_threads/src/filesystem_service/thread_index.rscrates/ironclaw_threads/tests/filesystem_session_thread_contract.rs
Coverage ratchetReborn integration-tier coverageLine coverage (Reborn crates): 85.59% — 303487 / 354600 lines Per-crate breakdown (63 crates, lowest-covered first)
This table itself is informational and never gates the PR on its own — not the percentage, not the per-crate holes, not the 0-coverage callout. A separate coverage ratchet (dry-run until enforce=true; see tests/integration/coverage-floor.toml) can fail the build on specific configured floors. Exemptions (3 entry/entries excluded from the accounting above)
|
|
🚅 Deployed to the ironclaw-pr-6096 environment in ironclaw-ci-preview
|
henrypark133
left a comment
There was a problem hiding this comment.
Code Review (multi-agent)
Intent: Serialize concurrent inbound-message writes per thread to preserve chronological persistence, display, and processing order on non-transactional backends.
Stats: 1 finding (from 1 retained finding) across 1 file; 8 review lenses completed using the exact-head snapshot, with parent-context fallback after the environment rejected parallel reviewer-agent dispatch at the thread limit. Body-only: 0. Existing current-head threads were deduplicated.
Tests
-
Medium Missing sequence-index failure coverage (
crates/ironclaw_threads/src/filesystem_service.rs:490-498, confidence 80) — anchor:crates/ironclaw_threads/src/filesystem_service.rs:496The non-transactional fallback can return a sequence-index write error after the message record is already durable and the write lock has been released. The current regression tests cover message-record failure and best-effort lookup-index failure, but do not exercise this sequence-index failure path or prove that a subsequent same-thread inbound message still proceeds.
Fix: Add
filesystem_accept_inbound_message_sequence_index_failure_releases_lock, injecting a sequence-index failure after the record succeeds, then asserting the next same-thread inbound message completes.
Existing unresolved current-head coverage not duplicated
The live PR already has unresolved threads for the process-local mutex across backend I/O, stale one-shot cache reseeding, delimiter-based thread-key collisions, the O(n) lock-registry sweep, and poisoned-lock fallback behavior. Those were excluded from this review to avoid duplicate feedback.
| /// CAS put — not ordering-sensitive, so callers that need the put itself | ||
| /// under a narrower critical section (the issue #6047 write lock in | ||
| /// `accept_inbound_message`) run this afterward, unlocked. | ||
| async fn finish_new_message_indexes( |
There was a problem hiding this comment.
Medium — Missing sequence-index failure coverage.
The non-transactional fallback propagates finish_new_message_indexes errors after the message record is durable and after the write lock is released. Existing tests cover message-record failure and best-effort lookup-index failure, but none inject a sequence-index write failure to verify the caller-level error and subsequent same-thread call behavior.
Fix: tests::filesystem_accept_inbound_message_sequence_index_failure_releases_lock covering a sequence-index write error after the message record succeeds, followed by a successful same-thread inbound message
There was a problem hiding this comment.
Added filesystem_accept_inbound_message_sequence_index_failure_releases_lock in a86c579, following the same pattern as the existing message-record-failure test — injects a sequence-index write failure after the record succeeds, asserts the caller gets Err, then asserts a same-thread follow-up message completes (proving the lock, already released before this failure point, doesn't wedge).
henrypark133
left a comment
There was a problem hiding this comment.
Follow-up finding
A delayed concurrency review completed after the initial batch and identified one additional issue not covered by the existing current-head threads. The review remains scoped to exact head d236be5a6c10132d6c5bca8c0de75dc2889e8b44.
| // which doesn't affect visibility and doesn't need to wait. | ||
| { | ||
| let write_lock = | ||
| self.inbound_message_write_lock(&thread_scope_key(&scope, &thread_id)); |
There was a problem hiding this comment.
High — Lock is not shared across service instances.
The registry is stored inside each FilesystemSessionThreadService, so two service instances sharing the same ScopedFilesystem or backend have independent locks. Concurrent inbound calls through those instances can still reserve and write sequences out of order, making the per-thread ordering guarantee depend on an undocumented singleton-service assumption.
Fix: Move the per-thread lock registry into shared filesystem/backend state, or otherwise ensure every wrapper over a store shares the same registry.
There was a problem hiding this comment.
This one I can't fully resolve with confidence from static reading — traced the wiring in ironclaw_reborn_composition/src/factory.rs and confirmed it's genuinely not a simple singleton (two separate build functions each construct their own instance, and there's 'reopen' language suggesting a store graph can be rebuilt mid-process). I don't have enough visibility into that reopen path's draining/lifecycle guarantees to say whether overlap is actually possible. Strengthened the doc comment on inbound_message_write_locks in a86c579 to state the assumption explicitly rather than leave it implicit, and filed #6102 to get a real answer on the reopen path before deciding whether the registry needs to move into shared backend state.
- weak_keyed_lock: recover a poisoned std::sync::Mutex instead of fabricating a fresh, unshared lock (was silently disabling mutual exclusion for every key once poisoned — coderabbit). - weak_keyed_lock: skip the retain() sweep on the live-key hit path so repeated lookups for an already-live key stay O(1) instead of rescanning the whole map on every accept_inbound_message call (coderabbit). - thread_scope_key: length-prefix each component instead of naive `:` concatenation, so it's injective regardless of what characters end up inside a tenant/thread-id component (gemini-code-assist; verified validate_scope_id permits `:`, though not reachable through the current concrete ThreadId/TenantId shapes today). - Add filesystem_accept_inbound_message_sequence_index_failure_releases_lock covering the sequence-index-write-failure path, alongside the existing message-record-failure coverage (review feedback from henrypark133). - Document the inbound_message_write_locks singleton-instance assumption more precisely and flag it as an open verification item against the composition factory's "reopen" store-graph path (henrypark133, High). Two items triaged as out-of-scope follow-ups rather than folded in here (reasoning posted on the PR): the one-shot context-window cache's seed/invalidate race, which predates this PR and affects both storage backends; and extending per-thread write serialization to assistant/ tool-result message writes, which is a real but separately-scoped gap (#6047 is specifically about user messages). Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
|
Closing as stale — no activity in over three weeks. The branch is untouched; reopen if this is still needed. |
Summary
Fixes #6047 — two chat messages sent in quick succession to the same thread could be persisted, displayed, and acted on out of chronological order.
Verified bug
Reproduced with a failing test on
origin/main@d68073336(before this fix), and confirmed the fix turns it green:Root cause
FilesystemSessionThreadService::accept_inbound_message(crates/ironclaw_threads/src/filesystem_service.rs) falls back to a non-transactional path on backends without multi-key transaction support — this is the shape of the production libSQL backend (crates/ironclaw_filesystem/src/libsql.rshas nobegin()override). On that path, the function:put.Nothing serialized two concurrent
accept_inbound_messagecalls against the same thread, so a message that reserved a later sequence number could finish its write before a message that reserved an earlier one — making the later message durably visible (and readable by the UI/agent) ahead of the earlier one. The transactional path (backends withbegin()support) is unaffected: its CAS-retry loop commits reserve+write atomically, and a losing writer retries into a fresh, strictly-later sequence number rather than reusing a stale one.Fix
(scope, thread_id)in-process async lock, held only around the ordering-sensitive step (reserve_sequence+ a newput_new_message_record— just the message-record CAS put), not the best-effort sequence-index/lookup-index/cache work that follows it.write_new_messageintoput_new_message_record(the CAS put) andfinish_new_message_indexes(the rest), so both the lock-guarded fallback path andwrite_new_message's three other callers share the same code with no duplication.weak_keyed_lockhelper (get-or-create a weak-referenced per-key async lock) used by both the new lock and the pre-existingthread_index_load_locksingleflight cache-fill lock in the same file, instead of duplicating that logic.one_shot_context_window_cache_key→thread_scope_keysince it's now shared by two call sites, with a doc comment distinguishing it from the differently-shapedthread_index_record_cache_keyin the neighboring module.with_fixed_viewboot-owner model used elsewhere in Reborn), and the pre-existingthread_index_load_lockalready relies on the same mechanism. A durable cross-process ordering guarantee would need the write itself to become one atomic domain-level operation (gap-aware reads with a bounded recovery contract, or transactional support on backends that lack it) — a materially larger change tracked as a follow-up, not folded into this fix.Regression tests
All in
crates/ironclaw_threads/tests/filesystem_session_thread_contract.rs:filesystem_accept_inbound_message_serializes_concurrent_calls_for_same_thread— the core proof: message B (called second) must not reserve a sequence or become visible while message A (called first) is still mid-write.filesystem_accept_inbound_message_write_lock_releases_after_failed_write— the lock releases via RAII even when the guarded write errors, so a subsequent call to the same thread doesn't hang forever.filesystem_accept_inbound_message_does_not_serialize_across_different_threads— the lock is scoped per-thread, not globally/per-tenant; unrelated threads never block each other.weak_keyed_lock_reclaims_entry_after_guard_dropped(inline unit test) — the shared lock registry doesn't leak an entry per distinct key forever.Verification
Review
Went through an 8-agent code review + two thermo-nuclear maintainability passes (plan and final diff). Findings addressed: narrowed the lock's critical section (was originally wider), fixed a lock-accessor signature inconsistency, disambiguated the renamed cache-key helper from a same-purpose-but-different-shape sibling, and added the two missing regression tests above. Security, bugs, maintainability, and approach reviewers returned zero findings.
Deferred (out of scope for this fix)
accept_inbound_message(durable, sequenced) andsubmit_turn's busy-check incrates/ironclaw_product_workflow/src/inbound_turn.rs— a second message can win the "claim the active run" race ahead of an earlier-accepted one. Bigger lift (touchesironclaw_turns'TurnActiveLockKeymachinery); flagging for a follow-up issue rather than folding into this PR.append_assistant_draftandappend_finalized_assistant_messageuse the same reserve-then-write two-step pattern for assistant messages on the same thread — out of scope since Task messages are processed and displayed out of chronological order #6047 specifically reports user-submitted messages, but worth a maintainer's awareness.