Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -128,6 +128,7 @@ xxhash-rust = { version = "0.8.15", features = ["xxh32", "xxh64"] }
rustls = { version = "0.23", default-features = false, features = [
"aws_lc_rs",
] }
rustix = { version = "1", features = ["fs"] }
rcgen = "0.14"

[workspace.lints.clippy]
Expand Down
7 changes: 6 additions & 1 deletion crates/ffwd-output/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@ ffwd-core = { version = "0.1.0", path = "../ffwd-core" }
ffwd-otap-proto = { version = "0.1.0", path = "../ffwd-otap-proto" }
ffwd-types = { version = "0.1.0", path = "../ffwd-types" }
flate2 = "1"
libc = { workspace = true }
httpdate = "1"
itoa = "1"
memchr = "2"
Expand Down Expand Up @@ -45,6 +44,12 @@ tracing = { workspace = true }
ureq = { version = "3", default-features = false, features = ["rustls"] }
zstd = "0.13"

[target.'cfg(target_os = "linux")'.dependencies]
rustix = { workspace = true }

[target.'cfg(target_os = "macos")'.dependencies]
libc = { workspace = true }

[dev-dependencies]
arrow = { workspace = true, features = ["test_utils"] }
insta = "1"
Expand Down
682 changes: 571 additions & 111 deletions crates/ffwd-output/src/pipelined.rs

Large diffs are not rendered by default.

75 changes: 71 additions & 4 deletions crates/ffwd-runtime/src/worker_pool/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -462,6 +462,7 @@ pub(super) async fn process_item(
mod tests {
use std::collections::BTreeSet;
use std::future::Future;
use std::io;
use std::pin::Pin;
use std::sync::Arc;
use std::time::Duration;
Expand Down Expand Up @@ -525,15 +526,15 @@ mod tests {
Box::pin(async { SendResult::Ok })
}

fn flush(&mut self) -> Pin<Box<dyn Future<Output = std::io::Result<()>> + Send + '_>> {
fn flush(&mut self) -> Pin<Box<dyn Future<Output = io::Result<()>> + Send + '_>> {
Box::pin(async { Ok(()) })
}

fn name(&self) -> &str {
"ok-sink"
}

fn shutdown(&mut self) -> Pin<Box<dyn Future<Output = std::io::Result<()>> + Send + '_>> {
fn shutdown(&mut self) -> Pin<Box<dyn Future<Output = io::Result<()>> + Send + '_>> {
Box::pin(async { Ok(()) })
}
}
Expand All @@ -549,15 +550,51 @@ mod tests {
Box::pin(std::future::pending())
}

fn flush(&mut self) -> Pin<Box<dyn Future<Output = std::io::Result<()>> + Send + '_>> {
fn flush(&mut self) -> Pin<Box<dyn Future<Output = io::Result<()>> + Send + '_>> {
Box::pin(async { Ok(()) })
}

fn name(&self) -> &str {
"hanging-sink"
}

fn shutdown(&mut self) -> Pin<Box<dyn Future<Output = std::io::Result<()>> + Send + '_>> {
fn shutdown(&mut self) -> Pin<Box<dyn Future<Output = io::Result<()>> + Send + '_>> {
Box::pin(async { Ok(()) })
}
}

struct WouldBlockThenOkSink {
calls: usize,
}

impl Sink for WouldBlockThenOkSink {
fn send_batch<'a>(
&'a mut self,
_batch: &'a RecordBatch,
_metadata: &'a BatchMetadata,
) -> Pin<Box<dyn Future<Output = SendResult> + Send + 'a>> {
let result = if self.calls == 0 {
self.calls += 1;
SendResult::IoError(io::Error::new(
io::ErrorKind::WouldBlock,
"buffer pool exhausted",
))
} else {
self.calls += 1;
SendResult::Ok
};
Box::pin(async move { result })
}

fn flush(&mut self) -> Pin<Box<dyn Future<Output = io::Result<()>> + Send + '_>> {
Box::pin(async { Ok(()) })
}

fn name(&self) -> &str {
"would-block-then-ok"
}

fn shutdown(&mut self) -> Pin<Box<dyn Future<Output = io::Result<()>> + Send + '_>> {
Box::pin(async { Ok(()) })
}
}
Expand Down Expand Up @@ -647,6 +684,36 @@ mod tests {
assert_eq!(retries, 0);
}

#[tokio::test]
async fn process_item_treats_would_block_as_transient_backpressure() {
let cancel = CancellationToken::new();
let mut sink = WouldBlockThenOkSink { calls: 0 };
let output_health = Arc::new(OutputHealthTracker::new(vec![]));
let metadata = BatchMetadata {
resource_attrs: Arc::default(),
observed_time_ns: 0,
};

let (outcome, _send_latency_ns, retries) = process_item(
ProcessItemContext {
worker_id: 0,
sink: &mut sink,
output_health: &output_health,
metadata: &metadata,
max_retry_delay: Duration::from_millis(1),
cancel: &cancel,
#[cfg(feature = "turmoil")]
batch_id: 0, // test only
},
make_batch(),
)
.await;

assert_eq!(outcome, DeliveryOutcome::Delivered);
assert_eq!(retries, 1);
assert_eq!(sink.calls, 2);
}

#[test]
fn terminalization_reducer_is_idempotent_for_any_two_step_schedule() {
let actions = [
Expand Down
8 changes: 4 additions & 4 deletions justfile
Original file line number Diff line number Diff line change
Expand Up @@ -689,23 +689,23 @@ bench-source-metadata-fast *ARGS:

# Profile OTLP decode/encode CPU with the normal allocator (flamegraph, per-mode timings).
profile-otlp-io *ARGS:
cargo run -p ffwd-bench --release --features bench-tools --bin otlp_io_profile -- {{ARGS}}
cargo run -p ffwd-bench --release --features bench-tools,io-bench --bin otlp_io_profile -- {{ARGS}}

# Profile OTLP decode/encode allocation counts with stats_alloc instrumentation.
profile-otlp-io-alloc *ARGS:
cargo run -p ffwd-bench --release --features bench-tools,otlp-profile-alloc --bin otlp_io_profile -- {{ARGS}}
cargo run -p ffwd-bench --release --features bench-tools,io-bench,otlp-profile-alloc --bin otlp_io_profile -- {{ARGS}}

# Generate microbenchmark report (markdown)
bench-report:
cargo run -p ffwd-bench --features bench-tools

# Profile FramedInput / format processing overhead and print a markdown report.
bench-framed-input *ARGS:
cargo run -p ffwd-bench --release --features bench-tools --bin framed_input_profile -- {{ARGS}}
cargo run -p ffwd-bench --release --features bench-tools,io-bench --bin framed_input_profile -- {{ARGS}}

# Allocation-focused FramedInput profiling (dhat-backed, slower; no throughput numbers).
bench-framed-input-alloc *ARGS:
cargo run -p ffwd-bench --release --features bench-tools,dhat-heap --bin framed_input_profile -- --alloc-only {{ARGS}}
cargo run -p ffwd-bench --release --features bench-tools,io-bench,dhat-heap --bin framed_input_profile -- --alloc-only {{ARGS}}

# Run sustained-load memory profiler (generator → SQL → null, default 5 minutes).
# Use --quick (30s) for CI or --medium (120s) for quick checks.
Expand Down
35 changes: 34 additions & 1 deletion scripts/verify_tla_coverage.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,10 @@
from __future__ import annotations

import argparse
import os
import pathlib
import re
import shutil
import subprocess
import tempfile
from typing import Sequence
Expand Down Expand Up @@ -72,6 +74,36 @@ def parse_invariants(body_lines: Sequence[str]) -> list[str]:
return invariants


def find_java() -> str:
"""Find a Java runtime, matching the repo's justfile fallback behavior."""
java_home = os.environ.get("JAVA_HOME")
candidates = [
os.environ.get("JAVA_BIN"),
os.path.join(java_home, "bin", "java") if java_home else None,
"/opt/homebrew/opt/openjdk/bin/java",
"/usr/local/opt/openjdk/bin/java",
"java",
]
for candidate in candidates:
if not candidate:
continue
java_bin = shutil.which(candidate) if os.path.basename(candidate) == candidate else candidate
if not java_bin:
continue
if not os.path.exists(java_bin):
continue
proc = subprocess.run(
[java_bin, "-version"],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
)
Comment on lines +95 to +99

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 👍 / 👎.

Comment on lines +95 to +99

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 👍 / 👎.

if proc.returncode == 0:
return java_bin
raise RuntimeError(
"Java runtime not found. Install OpenJDK (e.g. 'brew install openjdk') or set JAVA_BIN."
)
Comment thread
macroscopeapp[bot] marked this conversation as resolved.


def run_one(jar: pathlib.Path, tla_file: pathlib.Path, cfg_lines: Sequence[str], invariant: str) -> None:
before, _body, after = split_cfg_sections(cfg_lines)
rendered = before + ["INVARIANTS", f" {invariant}"] + after
Expand All @@ -81,8 +113,9 @@ def run_one(jar: pathlib.Path, tla_file: pathlib.Path, cfg_lines: Sequence[str],
tmp_cfg = pathlib.Path(fh.name)

try:
java_bin = find_java()
cmd = [
"java",
java_bin,
"-cp",
str(jar),
"tlc2.TLC",
Expand Down
9 changes: 6 additions & 3 deletions tla/DeliveryRetry.tla
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@
* 1. Send batch to sink
* 2. On Ok: terminal success (Delivered)
* 3. On Rejected: terminal permanent failure (Rejected)
* 4. On IoError/RetryAfter/Timeout: exponential backoff, retry forever
* 4. On IoError/RetryAfter/Timeout: transient retry, retry forever
* 5. On cancel (shutdown): terminal (PoolClosed)
*
* Key insight: the code deliberately retries FOREVER for transient
Expand Down Expand Up @@ -162,10 +162,13 @@ ReceiveRejected ==
/\ outcome' = "Rejected"
/\ UNCHANGED <<retryCount, backoffMs, cancelled, sinkState>>

\* ReceiveTransient: transient failure, compute next backoff and wait.
\* ReceiveTransient: transient failure, compute an abstract retry delay and wait.
\* Maps to Ok(SendResult::IoError), Ok(SendResult::RetryAfter), and
\* Err(_elapsed) (timeout) in process_item. All three follow the same
\* pattern: increment retry count, compute backoff, sleep, loop.
\* protocol pattern: increment retry count, sleep, loop. Timing details are
\* intentionally abstracted here: IoError/timeout use exponential backoff,
\* while RetryAfter uses the server-directed delay. Turmoil trace validators
\* cover that concrete timing distinction.
ReceiveTransient ==
/\ state = "Sending"
/\ sinkState = "Transient"
Expand Down
Loading