Skip to content

Await writer acknowledgement before reporting send success - #2746

Open
strawgate wants to merge 2 commits into
mainfrom
codex/f18c-2736-pipelined-ack
Open

strawgate wants to merge 2 commits into
mainfrom
codex/f18c-2736-pipelined-ack

Conversation

@strawgate

@strawgate strawgate commented Apr 29, 2026 •

Copy link
Copy Markdown
Owner

The pipelined sink's send_batch previously returned success before the writer thread confirmed the OS write, causing failed writes to be incorrectly tallied as delivered. This change makes the async worker block until the writer thread acknowledges each write.

  • Replaced deferred error polling with a oneshot channel; send_batch now awaits writer acknowledgement before returning
  • Moved byte/record counters after confirmed writes so only durable writes are tallied
  • Replaced unsafe libc calls with safe rustix wrappers on Linux for file hints
  • Added async and property-based tests covering ack semantics, error propagation, and buffer reuse

Note

Await writer acknowledgement before reporting send success in pipelined sink

  • Changes PipelinedSink so that send_batch only returns Ok after the writer thread acknowledges the OS write via a oneshot channel, replacing the previous deferred error-polling approach.
  • Buffer acquisition is now non-blocking, returning WouldBlock when no buffers are available; serialization failures retain the buffer as a spare for retry instead of cycling it through the writer.
  • Replaces libc calls with rustix equivalents for fallocate and fadvise on Linux in pipelined.rs.
  • Adds async unit and property tests covering ack flow, error propagation, backpressure, and serialization retry.
  • Behavioral Change: callers that previously received Ok before the write completed will now block until the writer acknowledges; WouldBlock is surfaced as a transient error and retried by the worker pool.

Macroscope summarized 9dc5235.

Copilot AI review requested due to automatic review settings April 29, 2026 05:42
@dosubot dosubot Bot added the size:L This PR changes 100-499 lines, ignoring generated files. label Apr 29, 2026
@pr-rocket

pr-rocket Bot commented Apr 29, 2026 •

Copy link
Copy Markdown

⚡ Two genuine bugs: (0) blocking empty_tx.send at line 382 creates a circular wait deadlock when empty_tx fills up while ack awaits, and (3) blocking send on filled_tx at line 289 stalls the tokio worker thread when the writer is slow.

  • 🐛 crates/ffwd-output/src/pipelined.rs:380 — Blocking send on empty_tx can deadlock when buffer pool exhausted and writer waits for async ack. — Blocking empty_tx.send creates circular wait with filled_tx.send and ack await
  • ⚠️ crates/ffwd-output/src/pipelined.rs:194 — Blocking send on filled_tx can stall async worker thread if writer thread is slow. — Blocking send on filled_tx stalls tokio worker when writer thread is slow
3 dismissed findings
  • crates/ffwd-output/src/pipelined.rs:286 — Spare buffer consumed on serialization failure but never recycled, causing pool starvation on repeated errors. — Spare buffer stored in self.spare_buf and reused on next acquire_buffer call
  • crates/ffwd-output/src/pipelined.rs:191 — Empty match arm returns WouldBlock for exhausted spare buffer, semantically incorrect error kind. — WouldBlock correctly signals transient buffer pool exhaustion to caller
  • crates/ffwd-output/src/pipelined.rs:482 — advise_dontneed returns early for len=0; original code was a no-op but now returns Ok(()), changing semantics. — Early return for len=0 is correct; fadvise with zero length is undefined

Last activity: Auto-review completed


Configuration (4 of 11 enabled)

Auto-title · Auto-body · Auto-review · Auto-labels

/rocket enable auto-pilot · /rocket enable <feature> · /rocket disable <feature>

@coderabbitai

coderabbitai Bot commented Apr 29, 2026 •

Copy link
Copy Markdown

Note

Reviews paused

It looks like this branch is under active development. To avoid overwhelming you with review comments due to an influx of new commits, CodeRabbit has automatically paused this review. You can configure this behavior by changing the reviews.auto_review.auto_pause_after_reviewed_commits setting.

Use the following commands to manage reviews:

  • @coderabbitai resume to resume automatic reviews.
  • @coderabbitai review to trigger a single review.

Use the checkboxes below for quick actions:

  • ▶️ Resume reviews
  • 🔍 Trigger review

Walkthrough

Workspace and output crate dependencies are refactored: rustix is added to workspace dependencies with fs feature, while libc is removed from ffwd-output's always-on dependencies and replaced with Linux-specific rustix usage. The pipelined sink's error handling is reworked to provide per-buffer synchronous acknowledgement via dedicated channels instead of queued errors, decoupling serialization failures from writer thread health. Platform hinting in FileWriter is migrated from macOS-specific flags to rustix-based Linux syscalls (fallocate, fadvise). Multiple profiling recipes are updated to include the io-bench feature flag. Comprehensive new tests validate the ack mechanism and error handling.

Possibly related PRs


Caution

Pre-merge checks failed

Please resolve all errors before merging. Addressing warnings is optional.

  • Ignore

❌ Failed checks (1 error, 2 warnings)

Check name Status Explanation Resolution
Formal Verification Coverage ❌ Error PR summary claims proptest-based randomized testing and formal verification for async pipeline refactoring, but code inspection reveals no proptest imports, strategies, or test cases implemented. No #[cfg(kani)] verification module present for new public APIs. Add proptest test cases covering arbitrary event ordering, out-of-order acks, and drain-while-in-flight scenarios. Implement #[cfg(kani)] verification module for critical state transition functions. Update dev-docs/VERIFICATION.md documenting proof coverage.
Documentation Thoroughly Updated ⚠️ Warning PR implements significant architectural changes (synchronous acknowledgement mechanism, buffer lifecycle refactoring) without updating canonical documentation. Add ADR to dev-docs/DESIGN.md for synchronous acknowledgement design choice; update DEVELOPING.md with buffer lifecycle lesson; update dev-docs/ARCHITECTURE.md for pipelined sink semantics.
Maintainer Fitness ⚠️ Warning PR claims measurable latency regression in hot-path code but provides no criterion benchmark results, baselines, or performance deltas as required by project policy. Run just bench on baseline and branch, capture criterion outputs with percent deltas, add results and baseline commit hash to PR description.
✅ Passed checks (4 passed)
Check name Status Explanation
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
High-Quality Rust Practices ✅ Passed PR demonstrates high-quality Rust practices with proper error handling, no unsafe code outside tests, safe rustix wrappers, and complete documentation.
Crate Boundary And Dependency Integrity ✅ Passed PR satisfies all crate boundary and dependency integrity requirements: ffwd-core isolated with forbid(unsafe_code), dependency direction strictly downward, binary crate contains only orchestration, single new dependency (rustix v1 with fs feature) properly gated for Linux.

Comment @coderabbitai help to get the list of available commands and usage tips.

@pr-rocket pr-rocket Bot changed the title fix: make pipelined file writes acknowledged Make pipelined file writes synchronous before returning success Apr 29, 2026
@pr-rocket pr-rocket Bot added bug Something isn't working refactor test labels Apr 29, 2026
@pr-rocket

pr-rocket Bot commented Apr 29, 2026

Copy link
Copy Markdown

⚡ Quick Review

The change correctly makes send_batch wait for the writer to complete before returning, fixing the acknowledged-write guarantee. However, there are two reliability issues worth addressing before merging.

  • Potential hang on writer panic/crash: In writer_thread_loop, after empty_tx.try_send(recycled) succeeds but before ack.send(result), if the writer panics the buffer is already recycled but the ack is never sent, causing the caller to hang indefinitely on ack_rx.recv(). Consider sending the ack first or using SendResult that can signal a closed writer.

  • Buffer leak on serialization error when writer is gone: In send_batch, on serialization failure the code calls recycle_buffer(buf) which sends through filled_tx. If the writer thread has already exited (channel closed), this silently drops the buffer — it was acquired from the pool but never returned, causing pool exhaustion on repeated errors. The buffer should be returned directly to empty_rx in this case.

Minor observations: the justfile changes appear unrelated to the core fix. The removal of macOS F_NOCACHE/F_PREALLOCATE hints is a behavior change for that platform — worth a separate note in the commit message. The new tests cover the happy path and error path well.


Review by PR Rocket

@gitar-bot gitar-bot Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Gitar has auto-approved this PR (configure)

@macroscopeapp

macroscopeapp Bot commented Apr 29, 2026 •

Copy link
Copy Markdown

Approvability

Verdict: Needs human review

1 blocking correctness issue found. This PR changes write delivery semantics from fire-and-forget to synchronous acknowledgement, fundamentally altering when write errors are reported to callers. Such behavioral changes to delivery guarantees warrant human review. Additionally, there are unresolved comments identifying a bug in the find_java() helper.

You can customize Macroscope's approvability policy. Learn more.

@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 591df8d to 34a24a6 Compare April 29, 2026 05:47

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 2

🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@crates/ffwd-output/src/pipelined.rs`:
- Around line 199-201: The hot path currently allocates a new sync_channel per
batch (ack_tx/ack_rx) before sending WriterMsg::Data via self.filled_tx; remove
that per-batch allocation by introducing a reusable acknowledgement mechanism
(e.g., a pre-created per-worker ack sender/receiver pair stored on the struct
like self.reusable_ack or a small pool/Arc-wrapped oneshot pool) and use that
instead of creating ack_tx/ack_rx inside the loop; update the send site that
constructs WriterMsg::Data to pass the reusable sender (or an index/ID into the
pool) and adjust the receiver-side handling to read from the shared/reused
channel, ensuring no per-batch sync_channel creation remains in the filled_tx
send path.
- Around line 615-674: The new async pipelined sink lacks proptest-based
property tests for arbitrary event ordering, out-of-order acknowledgements, and
draining while writes are in-flight; add a new proptest module that fuzzes
sequences of send_batch calls, simulated writer ACKs (including reordering and
duplicates) and drain() calls against PipelinedSink (constructed with
PipelineConfig and a test writer that can be driven by the strategy), and assert
invariants: all successful writes are observed in RecordingWriter (or
corresponding writes list), SendResult variants follow expectations (IoError vs
Ok), and ComponentStats (lines_total, bytes_total) reflect only successful acks;
reuse/extend the existing test helpers (FixedSerializer, RecordingWriter,
FailingWriter, empty_batch, metadata) and create strategies that produce
arbitrary ordering, out-of-order ACK delivery, and interleaved drain events to
validate the pipeline 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: Repository YAML (base), Organization UI (inherited)

Review profile: ASSERTIVE

Plan: Pro Plus

Run ID: e4351b40-3574-4f3a-9d7a-aeb0a50c7c37

📥 Commits

Reviewing files that changed from the base of the PR and between 8723f55 and 591df8d.

⛔ Files ignored due to path filters (1)
  • Cargo.lock is excluded by !**/*.lock
📒 Files selected for processing (3)
  • crates/ffwd-output/Cargo.toml
  • crates/ffwd-output/src/pipelined.rs
  • justfile

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
Comment thread crates/ffwd-output/src/pipelined.rs
@strawgate

Copy link
Copy Markdown
Owner Author

Addressed the buffer recycling concern in 34a24a6: serialization failures now keep the caller-owned buffer as a local spare instead of routing it through the writer thread, and serialization_error_keeps_buffer_available_for_retry covers the one-buffer retry case.

For the potential writer-panic hang: the ack sender is owned by the writer message, so if the writer thread panics or exits before sending the ack, ack_rx.recv() disconnects and returns BrokenPipe; the caller does not wait indefinitely.

@pr-rocket

pr-rocket Bot commented Apr 29, 2026

Copy link
Copy Markdown

⚡ Quick Review

Makes writes synchronous with a blocking ack channel, removing the asynchronous error-surfacing pattern. The behavioral shift is substantial: send_batch now blocks until bytes are confirmed on disk rather than overlapping writes and returning before I/O completes. This trades throughput for stronger durability guarantees — acceptable if that's the documented contract, but worth calling out in the changelog.

  • Blocking ack risk: write_filled waits on ack_rx.recv() indefinitely if the writer thread crashes or deadlocks before sending the ack. The let _ = ack.send(result) silently drops failures, which could leave callers stuck if the sender is dropped. Consider a timeout or at least documenting this invariant.
  • Buffer pool contention under slow I/O: With num_buffers: 1 and synchronous writes, the serialization thread will block waiting for the buffer to return. The pool sizing becomes more critical now that writes are serial with serialization.
  • macOS support removed: apply_macos_hints is gone; only Linux optimizations remain. If macOS is a supported target, this needs clarification or a follow-up.
  • spare_buf logic is sound: Caller-owned fallback on serialization failure is cleaner than routing through the writer thread.
  • Test coverage is good: The three new tests cover the key paths (success, writer failure, serialization retry). No test for the case where ack send fails in the writer thread — worth adding.
  • Dependency swap (libc → rustix) looks safe; rustix provides a safe API surface for the same Linux syscalls.

Review by PR Rocket

@pr-rocket pr-rocket Bot added the enhancement New feature or request label Apr 29, 2026

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Note

Copilot was unable to run its full agentic suite in this review.

This PR updates ffwd-output’s pipelined sink to only acknowledge a batch as delivered once the writer thread has persisted it, and adjusts related tooling and platform-specific file I/O behavior.

Changes:

  • Add per-batch writer acknowledgements so send_batch returns Ok only after the writer finishes the write.
  • Fix accounting/buffer handling so success counters and buffer recycling reflect actual persisted writes and serialization failures.
  • Remove unsafe file allocation hints in favor of safe Linux rustix APIs; update just benchmark/profile recipes to enable io-bench.

Reviewed changes

Copilot reviewed 3 out of 4 changed files in this pull request and generated 3 comments.

File Description
justfile Enables io-bench feature for profiling/benchmark recipes so gated bins can run.
crates/ffwd-output/src/pipelined.rs Introduces write acknowledgements, adjusts buffer lifecycle, switches Linux file hints to rustix, and adds coverage tests.
crates/ffwd-output/Cargo.toml Drops libc and adds Linux-only rustix dependency for safe file APIs.

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
Comment thread crates/ffwd-output/src/pipelined.rs Outdated
Comment thread crates/ffwd-output/src/pipelined.rs Outdated
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 34a24a6 to d3421db Compare April 29, 2026 05:53
@dosubot dosubot Bot added size:XL This PR changes 500-999 lines, ignoring generated files. and removed size:L This PR changes 100-499 lines, ignoring generated files. labels Apr 29, 2026
@strawgate

Copy link
Copy Markdown
Owner Author

Follow-up pushed in d3421db:

  • replaced per-batch sync_channel creation with a reusable writer acknowledgement channel owned by the sink/writer pair
  • moved rustix into root workspace.dependencies and changed ffwd-output to use workspace = true
  • added a proptest over arbitrary send/flush sequences with scripted serialization and writer success/failure outcomes, asserting writes and stats only reflect acknowledged successes

Out-of-order writer acknowledgements are not representable in the current implementation: Sink::send_batch requires &mut self, the file writer is a single OS thread, and the reusable ack channel is consumed synchronously before the next batch can be sent.

@pr-rocket

pr-rocket Bot commented Apr 29, 2026

Copy link
Copy Markdown

⚡ Quick Review

The change makes send_batch wait synchronously for the writer to complete each write before returning, replacing the previous fire-and-forget approach with per-batch blocking. This fundamentally removes the I/O pipelining overlap (serializing N+1 while writing N) — the PR description's framing is misleading since the architecture is now purely sequential. The rustix migration and macOS hints removal are reasonable cleanups.

Key observations:

  • Pipelining is gone: With synchronous write_filled that blocks on write_ack_rx.recv(), each send_batch call now waits for the previous batch's write to complete before the buffer pool unblocks. This eliminates the main benefit of the triple-buffer architecture and will increase worker thread latency under load.
  • Buffer pool starvation risk: With num_buffers: 1, a send_batch holds its buffer for the entire write duration. A second concurrent send_batch (if any) will block on acquire_buffer() since the pool is empty. The synchronous model effectively serializes at the buffer boundary.
  • macOS optimizations dropped: apply_macos_hints (F_NOCACHE, F_PREALLOCATE) was removed entirely. The updated table shows "standard write_all" for macOS, which may regress performance for macOS users.
  • Pre-allocation safety: PREALLOC_BYTES is u64 but fallocate takes i64; on 32-bit targets this could truncate. Consider an explicit cast or i64::try_from.
  • spare_buf semantics shift: Serialization failures now store the caller's buffer directly rather than sending it through the writer thread as a zero-byte no-op. This is safer (avoids writer-thread dependency) but different enough to warrant a comment or test case covering the exact retry path with a full pool.
  • apply_linux_hints error handling: fallocate and fadvise now silently ignore errors (no tracing::warn); the previous code was similarly silent but this could mask permission issues or filesystem incompatibilities (e.g., tmpfs, network mounts).
  • Test nesting: The proptest case spawns tokio::runtime::Builder::new_current_thread() inside the test body. If the test runner itself uses a single-threaded runtime, this could panic on enter() — worth verifying against the actual CI runtime configuration.

Review by PR Rocket

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Actionable comments posted: 1

♻️ Duplicate comments (2)
crates/ffwd-output/src/pipelined.rs (2)

203-205: ⚠️ Potential issue | 🟠 Major

Avoid allocating a fresh ack channel on every batch.

Line 205 creates a new sync_channel for each send_batch, which keeps allocation/synchronization setup on the steady-state write path. Please reuse/pool the ack mechanism or add a benchmark-backed justification for this overhead.

As per coding guidelines: "Do not introduce allocations in hot paths without benchmarking".

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@crates/ffwd-output/src/pipelined.rs` around lines 203 - 205, The current
write_filled implementation (fn write_filled) creates a new sync_channel for
every batch which allocates on the hot path; change this to reuse an
acknowledgment mechanism instead of allocating per call — e.g., move a single
ack channel (or a small pool of ack senders/receivers) into the writer struct
and reuse that for each send_batch, or replace per-batch channels with a
preallocated oneshot/pool or an atomic/Condvar handshake owned by the writer
thread; update write_filled and the corresponding send_batch/reader logic to use
the reused ack handle (e.g., a field like ack_tx/ack_rx or an AckPool) so no new
sync_channel is created on each batch.

522-722: ⚠️ Potential issue | 🟠 Major

Required property coverage for the pipeline is still missing.

These tests cover concrete success/error/retry cases, but the repository rule for pipeline code also requires proptest coverage for arbitrary event ordering, out-of-order acks, and drain while writes are in-flight. Please add that property suite or document why those states are impossible here.

As per coding guidelines: "New async pipeline code must have proptest coverage for: arbitrary event ordering, acks out of order, drain while in-flight".

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@crates/ffwd-output/src/pipelined.rs` around lines 522 - 722, Add a
proptest-based property suite for the PipelinedSink that generates arbitrary
sequences of events (send_batch calls, simulated writer successes/failures,
out-of-order acknowledgements, and intermittent drain calls) to exercise
concurrency edges; specifically write proptest tests invoking PipelinedSink::new
and repeatedly calling send_batch (observing SendResult), toggling behaviors on
custom BatchSerializer and BatchWriter mocks (e.g., variants of
FixedSerializer/FailsOnceSerializer/FailingWriter) and calling
sink.drain/shutdown while writes are in-flight, then assert invariant
properties: no silent loss of lines/bytes except when an IoError is returned
(check ComponentStats.lines_total and bytes_total), buffers are returned to
availability (no deadlock/hangs), and eventual consistency of writes (either
recorded writes match produced payloads or an IoError was reported). Place tests
in the same tests module and if any state is impossible to reach, add a short
comment documenting why instead of a test.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Inline comments:
In `@crates/ffwd-output/src/pipelined.rs`:
- Around line 278-281: The code currently acknowledges (returns SendResult::Ok)
after BatchWriter::write() via write_filled, but FileWriter::write() only does
write_all() and defers durable persistence to flush()/sync_data()/shutdown(), so
the ack can happen before bytes are crash-safe; update the ack point so
durability is guaranteed by either (a) making BatchWriter::write() perform the
durability boundary (call flush()/sync_data()/shutdown() as appropriate on the
underlying FileWriter) before returning, or (b) move the acknowledgement in
write_filled to occur only after invoking the writer's flush()/sync_data() (or
shutdown()) and confirming success; alternatively, if you intentionally do not
want to wait for durability, change the surrounding comments/contract that
promise "persisted" bytes to reflect that SendResult::Ok does not imply durable
persistence. Ensure you update the code paths referencing write_filled,
BatchWriter::write, FileWriter::write, flush, sync_data, shutdown and the
SendResult::Ok behavior consistently.

---

Duplicate comments:
In `@crates/ffwd-output/src/pipelined.rs`:
- Around line 203-205: The current write_filled implementation (fn write_filled)
creates a new sync_channel for every batch which allocates on the hot path;
change this to reuse an acknowledgment mechanism instead of allocating per call
— e.g., move a single ack channel (or a small pool of ack senders/receivers)
into the writer struct and reuse that for each send_batch, or replace per-batch
channels with a preallocated oneshot/pool or an atomic/Condvar handshake owned
by the writer thread; update write_filled and the corresponding
send_batch/reader logic to use the reused ack handle (e.g., a field like
ack_tx/ack_rx or an AckPool) so no new sync_channel is created on each batch.
- Around line 522-722: Add a proptest-based property suite for the PipelinedSink
that generates arbitrary sequences of events (send_batch calls, simulated writer
successes/failures, out-of-order acknowledgements, and intermittent drain calls)
to exercise concurrency edges; specifically write proptest tests invoking
PipelinedSink::new and repeatedly calling send_batch (observing SendResult),
toggling behaviors on custom BatchSerializer and BatchWriter mocks (e.g.,
variants of FixedSerializer/FailsOnceSerializer/FailingWriter) and calling
sink.drain/shutdown while writes are in-flight, then assert invariant
properties: no silent loss of lines/bytes except when an IoError is returned
(check ComponentStats.lines_total and bytes_total), buffers are returned to
availability (no deadlock/hangs), and eventual consistency of writes (either
recorded writes match produced payloads or an IoError was reported). Place tests
in the same tests module and if any state is impossible to reach, add a short
comment documenting why instead of a test.
🪄 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: Repository YAML (base), Organization UI (inherited)

Review profile: ASSERTIVE

Plan: Pro Plus

Run ID: 25b0a356-8a8b-46c8-9cc8-835055d86fc5

📥 Commits

Reviewing files that changed from the base of the PR and between 591df8d and 34a24a6.

⛔ Files ignored due to path filters (1)
  • Cargo.lock is excluded by !**/*.lock
📒 Files selected for processing (3)
  • crates/ffwd-output/Cargo.toml
  • crates/ffwd-output/src/pipelined.rs
  • justfile

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
@gitar-bot

gitar-bot Bot commented Apr 29, 2026 •

Copy link
Copy Markdown

Note

Your trial team has used its Gitar budget, so automatic reviews are paused. Upgrade now to unlock full capacity. Comment "Gitar review" to trigger a review manually.
Learn more about usage limits

Code Review ✅ Approved

Synchronous file write implementation ensures data persistence before returning success. No issues found.

Auto-approved: The PR correctly addresses the bug by introducing a synchronous acknowledgement mechanism between the writer thread and the sink, replacing unreliable error reporting. The refactoring to use rustix for platform-specific file operations is safer and more idiomatic, and the inclusion of unit tests provides high confidence in the correctness of the changes.

Options

Display: compact → Showing less information.

Comment with these commands to change:

Compact
gitar display:verbose         

Important

Your trial ends in 4 days — upgrade now to keep code review, CI analysis, auto-apply, custom automations, and more.

Was this helpful? React with 👍 / 👎 | Gitar

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: d3421db239

ℹ️ 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".

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from d3421db to 603e76a Compare April 29, 2026 06:03
@pr-rocket pr-rocket Bot changed the title Make pipelined file writes synchronous before returning success Make pipelined writes synchronous before returning success Apr 29, 2026
@pr-rocket

pr-rocket Bot commented Apr 29, 2026

Copy link
Copy Markdown

⚡ Quick Review

This is a well-designed change that converts the pipelined sink from fire-and-forget writes to synchronous acknowledgment before returning. The architecture is sound and the test coverage is thorough. Minor concerns around potential deadlock under extreme backpressure and a removed macOS optimization.

  • Synchronous write acknowledgement: The shift from async error reporting via a separate channel to waiting on write_ack_rx before returning is the right approach — callers now know immediately if a write succeeded, and stats counters only increment on confirmed writes. The ordering in writer_thread_loop (recycle buffer → send ack) is correct.
  • Serialization error handling: Keeping the failed buffer in spare_buf instead of sending it through the writer thread is cleaner and avoids a potential dependency on an unhealthy writer.
  • Buffer pool exhaustion / deadlock risk: With num_buffers buffers in the pool, if all buffers are in-flight awaiting ack and all callers are blocked in acquire_buffer, the writer cannot make progress (no slot to recycle into). This is the intended backpressure behavior, but worth documenting — it means the pool size acts as the maximum concurrency limit.
  • macOS optimization removed: The F_NOCACHE and F_PREALLOCATE hints are gone. If this sink targets macOS in production, the sequential-write performance may regress without fallocate-equivalent pre-allocation.
  • Minor: apply_platform_hints parameter: _file is unused — either remove it or add #[allow(unused_variables)] for clarity.
  • Minor: justfile changes: Feature flag io-bench added consistently to all bench commands — looks correct and unrelated to the core change.

Review by PR Rocket

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

♻️ Duplicate comments (1)
crates/ffwd-output/src/pipelined.rs (1)

278-281: ⚠️ Potential issue | 🟠 Major

SendResult::Ok still overclaims durability.

Line 280 says checkpoints advance only after bytes are "actually persisted", but FileWriter::write() at Lines 451-463 only does write_all(). Crash-safe persistence still happens later in flush()/shutdown() via sync_data() at Lines 466-476, so an acknowledged batch can still be lost on process or host crash. Either move the ack behind that durability boundary or relax the contract to "written to the OS".

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@crates/ffwd-output/src/pipelined.rs` around lines 278 - 281, The
SendResult::Ok acknowledgement currently claims durability ("actually
persisted") but FileWriter::write()/write_filled only calls write_all() (no
sync) while true durability is provided later by flush()/shutdown() via
sync_data(); either move the delivery acknowledgement to after the durability
boundary (i.e., perform sync_data()/fsync in FileWriter before returning
SendResult::Ok from write_filled/write()) or relax the contract and update
SendResult::Ok's documentation/semantics to state it only guarantees the bytes
were written to the OS buffer (write_all()), not fsynced; modify the
implementations and comments around write_filled, FileWriter::write,
FileWriter::flush/shutdown, and the SendResult::Ok docs accordingly so behavior
and docs match.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.

Duplicate comments:
In `@crates/ffwd-output/src/pipelined.rs`:
- Around line 278-281: The SendResult::Ok acknowledgement currently claims
durability ("actually persisted") but FileWriter::write()/write_filled only
calls write_all() (no sync) while true durability is provided later by
flush()/shutdown() via sync_data(); either move the delivery acknowledgement to
after the durability boundary (i.e., perform sync_data()/fsync in FileWriter
before returning SendResult::Ok from write_filled/write()) or relax the contract
and update SendResult::Ok's documentation/semantics to state it only guarantees
the bytes were written to the OS buffer (write_all()), not fsynced; modify the
implementations and comments around write_filled, FileWriter::write,
FileWriter::flush/shutdown, and the SendResult::Ok docs accordingly so behavior
and docs match.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Repository YAML (base), Organization UI (inherited)

Review profile: ASSERTIVE

Plan: Pro Plus

Run ID: 0eceae81-d39c-43e7-bae7-c2dd4d3b5d76

📥 Commits

Reviewing files that changed from the base of the PR and between 34a24a6 and d3421db.

⛔ Files ignored due to path filters (1)
  • Cargo.lock is excluded by !**/*.lock
📒 Files selected for processing (4)
  • Cargo.toml
  • crates/ffwd-output/Cargo.toml
  • crates/ffwd-output/src/pipelined.rs
  • justfile

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 603e76a443

ℹ️ 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".

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 603e76a to 48e94e5 Compare April 29, 2026 06:07
@pr-rocket pr-rocket Bot changed the title Make pipelined writes synchronous before returning success Await write acknowledgement before reporting success Apr 29, 2026
@pr-rocket

pr-rocket Bot commented Apr 29, 2026

Copy link
Copy Markdown

⚡ Quick Review

The change makes pipelined writes synchronous by waiting for the writer thread to confirm successful I/O before returning from send_batch, replacing the previous fire-and-forget pattern with deferred error checking. This is a safer default, though it reduces write throughput by eliminating pipelining overlap.

Key observations:

  • Blocking risk: write_ack_rx.recv().await blocks the async task waiting for the std writer thread. If the writer thread stalls (e.g., slow disk), send_batch will not progress. This trades throughput for the stronger guarantee that success means "written", which seems intentional per the PR title.
  • Serde error path: Storing the failed buffer in self.spare_buf avoids the old behavior of sending a no-op through the writer thread, which is cleaner. However, if send_batch is called twice concurrently (shouldn't happen given the Sink contract, but worth noting), the second call could overwrite the spare buffer before the first retry completes.
  • macOS hints removed: The macOS-specific F_NOCACHE and F_PREALLOCATE code was removed entirely. If macOS support is still intended, those optimizations are now gone.
  • Drop ordering in writer_thread_loop: Changed from try_send to send for buffer recycling — if empty_tx.send() blocks indefinitely, the writer thread could deadlock if the receiver is gone. The old try_send was more resilient to receiver shutdown.
  • Minor: apply_linux_hints(_file: &std::fs::File) parameter is unused (named _file) but the function does use file internally — the underscore prefix is misleading.
  • The new rustix dependency replaces raw libc FFI, which is a net positive for safety.

Review by PR Rocket

@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 48e94e5 to 417ec65 Compare April 29, 2026 06:15
@pr-rocket

pr-rocket Bot commented Apr 29, 2026

Copy link
Copy Markdown

⚡ Quick Review

The PR correctly implements synchronous write acknowledgement — the sink now awaits the writer thread's result before reporting success, preventing premature success reports on I/O failures. Replacing unsafe libc calls with rustix is a solid safety win for Linux. Overall low risk.

Key observations:

  • write_ack_tx.blocking_send ordering: The ack is sent before the empty buffer is recycled. If blocking_send fails (receiver dropped), the buffer is still returned to the pool via empty_tx.send(recycled). This means a caller whose ack receiver is dropped won't deadlock, but also won't see the failure — which is acceptable since the caller is being dropped anyway.
  • Blocking empty_tx.send in writer thread: This changed from try_send. If the receiver is gone during shutdown, this will block indefinitely. However, in normal operation the receiver always drains, and on true shutdown the writer loop exits anyway — so this should be fine in practice.
  • try_send → send for empty buffer: Previously a non-blocking send, now blocking. The writer thread is dedicated with no other senders, so this is safe and removes a silent buffer leak on receiver overload.
  • map_or → map_or_else: Correct fix — avoids evaluating the fallback when metadata is Ok.
  • macOS hints removed: The macOS-specific F_NOCACHE/F_PREALLOCATE optimization was dropped without mention in the summary. This is intentional per the table in the new doc comment, but worth confirming.
  • Tests: Property-based test with proptest covers mixed success/failure sequences well, including the critical case where serialization succeeds but write fails — counters should not be incremented.

Review by PR Rocket

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 417ec65c2e

ℹ️ 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".

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 417ec65 to cfa50f1 Compare April 29, 2026 07:01
@pr-rocket pr-rocket Bot changed the title Await write acknowledgement before reporting success Await writer acknowledgement before reporting send success Apr 29, 2026
@pr-rocket

pr-rocket Bot commented Apr 29, 2026

Copy link
Copy Markdown

⚡ Quick Review

The PR implements synchronous write acknowledgement semantics for the pipelined sink, replacing the deferred error channel with an awaitable ack per batch. This is a meaningful correctness improvement but introduces a few subtle risks worth examining.

Observations:

  • Potential ack-drops-buffer race in writer thread: In writer_thread_loop, empty_tx.send(recycled) completes before write_ack_tx.blocking_send(result). If the ack channel is full (e.g., caller under load) or the receiver is dropped mid-flight, the buffer is already back in the pool but the caller will block indefinitely on write_ack_rx.recv(). Consider reversing the order or using a paired send that only succeeds if both complete.

  • macOS optimizations silently removed: apply_macos_hints (F_NOCACHE, F_PREALLOCATE) was deleted entirely. The stub apply_platform_hints now just ignores the file on non-Linux targets. If macOS support was intentional, this is a regression; if not, the code should document that explicitly.

  • spare_buf is caller-owned but never reclaimed on drop: On serialization failure, the buffer is held in spare_buf for retry. If the sink is dropped while spare_buf is populated (e.g., repeated serialization failures followed by shutdown), the buffer is leaked. Minor since it's bounded by buf_capacity, but worth noting in the struct's doc comment.

  • acquire_buffer silently discards capacity on retry: spare_buf stores whatever buffer the failed serializer left behind. If a subsequent batch requires more capacity than that buffer provides, serialize will re-allocate anyway — this is fine but worth confirming it's the intended behavior.

  • apply_platform_hints dead-branches on non-Linux: The #[cfg(not(target_os = "linux"))] arm is a no-op _ = file; statement that has no effect. This is harmless but misleading; the function's behavior on non-Linux platforms is not well documented.

  • Proptest tests spawn nested tokio runtime: The property tests build a single-threaded runtime inside a tokio test context that may already be multi-threaded. This works but adds unnecessary overhead; consider reusing the existing runtime or marking the test #[tokio::test(flavor = "current_thread")] to match.

The libc → rustix migration is sound; the types and semantics align. The preallocation constant type change from i64 to u64 is safe for the given value.


Review by PR Rocket

@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch 2 times, most recently from 6ac8d9c to b92f312 Compare April 29, 2026 13:15

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: b92f312536

ℹ️ 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".

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
Comment thread crates/ffwd-output/src/pipelined.rs Outdated
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from b92f312 to 9ea8185 Compare April 29, 2026 13:19
@pr-rocket pr-rocket Bot changed the title Await writer acknowledgement before reporting send success Ensure writer acknowledgement before reporting send success Apr 29, 2026
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 9ea8185 to 3fe96b8 Compare April 29, 2026 15:42
@pr-rocket pr-rocket Bot changed the title Ensure writer acknowledgement before reporting send success Await writer acknowledgement before reporting send success Apr 29, 2026

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 3fe96b8cb1

ℹ️ 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".

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 3fe96b8 to 704e957 Compare April 29, 2026 16:09

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 704e9576be

ℹ️ 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".

Comment thread crates/ffwd-output/src/pipelined.rs Outdated
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 704e957 to 75209e5 Compare April 29, 2026 16:17
@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 75209e5 to 82ef6cd Compare April 29, 2026 16:34
Comment thread scripts/verify_tla_coverage.py

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 271771e02f

ℹ️ 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".

Comment on lines +91 to +95
proc = subprocess.run(
[java_bin, "-version"],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Skip nonexistent Java paths when probing runtimes

find_java() now probes hardcoded absolute paths by calling subprocess.run directly, but when /opt/homebrew/opt/openjdk/bin/java (or /usr/local/opt/openjdk/bin/java) is absent—as on most Linux CI hosts—subprocess.run raises FileNotFoundError and aborts immediately instead of falling through to java on PATH. This makes scripts/verify_tla_coverage.py fail even when Java is installed, unlike the justfile fallback loop that checks candidate executables before invoking them.

Useful? React with 👍 / 👎.

@strawgate
strawgate force-pushed the codex/f18c-2736-pipelined-ack branch from 271771e to 9dc5235 Compare April 29, 2026 18:07

@chatgpt-codex-connector chatgpt-codex-connector Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

💡 Codex Review

Here are some automated review suggestions for this pull request.

Reviewed commit: 9dc5235a56

ℹ️ 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".

Comment on lines +95 to +99
proc = subprocess.run(
[java_bin, "-version"],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Skip non-executable Java candidates during probing

find_java() now checks os.path.exists() and then runs [java_bin, "-version"], but existing paths that are not executable (for example a stale JAVA_BIN pointing to a regular file) raise PermissionError and abort the script before it can fall back to later candidates like java on PATH. I reproduced this by setting JAVA_BIN to a non-executable temp file, which crashes in subprocess.run instead of continuing, so CI environments with misconfigured Java env vars can fail even when a valid runtime is installed.

Useful? React with 👍 / 👎.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

bug Something isn't working enhancement New feature or request refactor size:XL This PR changes 500-999 lines, ignoring generated files. test

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants