checkpoint

⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
Fabro 2026-05-27 19:33:27 -04:00
parent 12280825a2
commit 39b7e5a648
5 changed files with 1071 additions and 408 deletions

1119
run.json

File diff suppressed because one or more lines are too long

View file

@ -0,0 +1,321 @@
diff --git a/lib/crates/fabro-cli/src/commands/run/runner.rs b/lib/crates/fabro-cli/src/commands/run/runner.rs
index 17c96751d..ac1a2bbd0 100644
--- a/lib/crates/fabro-cli/src/commands/run/runner.rs
+++ b/lib/crates/fabro-cli/src/commands/run/runner.rs
@@ -1,3 +1,4 @@
+use std::collections::{HashSet, VecDeque};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::time::Duration;
@@ -289,6 +290,38 @@ fn load_worker_vault(storage_dir: Option<&Path>) -> Result<Option<Arc<AsyncRwLoc
const WORKER_CONTROL_RECONNECT_INITIAL_BACKOFF: Duration = Duration::from_millis(100);
const WORKER_CONTROL_RECONNECT_MAX_BACKOFF: Duration = Duration::from_secs(5);
+const WORKER_CONTROL_APPLIED_ID_DEDUPE_CAPACITY: usize = 2048;
+
+#[derive(Default)]
+struct AppliedWorkerControlDeliveryIds {
+ last: Option<String>,
+ order: VecDeque<String>,
+ recent: HashSet<String>,
+}
+
+impl AppliedWorkerControlDeliveryIds {
+ fn last_applied_id(&self) -> Option<&str> {
+ self.last.as_deref()
+ }
+
+ fn contains(&self, id: &str) -> bool {
+ self.recent.contains(id)
+ }
+
+ fn record(&mut self, id: String) {
+ if !self.recent.insert(id.clone()) {
+ self.last = Some(id);
+ return;
+ }
+ self.order.push_back(id.clone());
+ self.last = Some(id);
+ while self.order.len() > WORKER_CONTROL_APPLIED_ID_DEDUPE_CAPACITY {
+ if let Some(evicted) = self.order.pop_front() {
+ self.recent.remove(&evicted);
+ }
+ }
+ }
+}
struct WorkerControlManagerHandle {
first_connection: Option<oneshot::Receiver<Result<()>>>,
@@ -439,14 +472,14 @@ async fn run_worker_control_manager(
let mut first_tx = Some(first_tx);
let mut fatal_tx = Some(fatal_tx);
let mut backoff = WORKER_CONTROL_RECONNECT_INITIAL_BACKOFF;
- let mut last_applied_id: Option<String> = None;
+ let mut applied_ids = AppliedWorkerControlDeliveryIds::default();
while !done.is_cancelled() {
let request = match build_worker_control_stream_request(
&target,
&run_id,
&worker_token,
- last_applied_id.as_deref(),
+ applied_ids.last_applied_id(),
) {
Ok(request) => request,
Err(err) => {
@@ -474,7 +507,7 @@ async fn run_worker_control_manager(
&cancel_token,
&steering_hub,
&run_control,
- &mut last_applied_id,
+ &mut applied_ids,
&done,
)
.await
@@ -645,7 +678,7 @@ async fn handle_worker_control_socket(
cancel_token: &CancellationToken,
steering_hub: &fabro_workflow::SteeringHub,
run_control: &RunControlState,
- last_applied_id: &mut Option<String>,
+ applied_ids: &mut AppliedWorkerControlDeliveryIds,
done: &CancellationToken,
) -> Result<(), WorkerControlConnectError> {
let mut ping_interval = time::interval(WORKER_CONTROL_WS_PING_INTERVAL);
@@ -692,7 +725,7 @@ async fn handle_worker_control_socket(
cancel_token,
steering_hub,
run_control,
- last_applied_id,
+ applied_ids,
frame,
)
.await;
@@ -730,16 +763,16 @@ async fn apply_worker_control_delivery_frame(
cancel_token: &CancellationToken,
steering_hub: &fabro_workflow::SteeringHub,
run_control: &RunControlState,
- last_applied_id: &mut Option<String>,
+ applied_ids: &mut AppliedWorkerControlDeliveryIds,
frame: WorkerControlDeliveryFrame,
) -> bool {
// Duplicate ids cannot reach us under normal operation: the server replays
- // strictly after `last_applied_id`. Guard against a server-side bug by
- // ignoring any frame whose id is not strictly newer than what we last
- // applied.
- if last_applied_id.as_deref() == Some(frame.id.as_str()) {
+ // strictly after the last applied id. Guard against a server-side bug or
+ // reconnect race by ignoring recently-applied delivery ids.
+ if applied_ids.contains(&frame.id) {
return false;
}
+ let frame_id = frame.id;
apply_worker_control_message(
interviewer,
cancel_token,
@@ -748,7 +781,7 @@ async fn apply_worker_control_delivery_frame(
frame.envelope,
)
.await;
- *last_applied_id = Some(frame.id);
+ applied_ids.record(frame_id);
true
}
@@ -1202,8 +1235,8 @@ mod tests {
use tokio_util::sync::CancellationToken;
use super::{
- WorkerControlConnectError, WorkerControlSocket, WorkerTitlePhase,
- apply_worker_control_delivery_frame, apply_worker_control_message,
+ AppliedWorkerControlDeliveryIds, WorkerControlConnectError, WorkerControlSocket,
+ WorkerTitlePhase, apply_worker_control_delivery_frame, apply_worker_control_message,
build_worker_control_stream_request, connect_worker_control_stream,
handle_worker_control_socket, initial_worker_title_phase, load_worker_vault,
next_worker_control_reconnect_backoff, stamp_system_worker, worker_title,
@@ -1538,7 +1571,7 @@ mod tests {
let cancel_token = CancellationToken::new();
let run_control = RunControlState::new();
let hub = test_steering_hub();
- let mut last_applied_id: Option<String> = None;
+ let mut applied_ids = AppliedWorkerControlDeliveryIds::default();
let frame = fabro_interview::WorkerControlDeliveryFrame {
id: "local:1".to_string(),
envelope: WorkerControlEnvelope::pause_run(),
@@ -1550,7 +1583,7 @@ mod tests {
&cancel_token,
&hub,
&run_control,
- &mut last_applied_id,
+ &mut applied_ids,
frame.clone(),
)
.await
@@ -1561,13 +1594,13 @@ mod tests {
&cancel_token,
&hub,
&run_control,
- &mut last_applied_id,
+ &mut applied_ids,
frame,
)
.await
);
- assert_eq!(last_applied_id, Some("local:1".to_string()));
+ assert_eq!(applied_ids.last_applied_id(), Some("local:1"));
}
#[test]
@@ -1648,7 +1681,7 @@ mod tests {
let cancel_token = CancellationToken::new();
let hub = test_steering_hub();
let run_control = RunControlState::new();
- let mut last_applied_id: Option<String> = None;
+ let mut applied_ids = AppliedWorkerControlDeliveryIds::default();
let done = CancellationToken::new();
let task = tokio::spawn(async move {
@@ -1658,7 +1691,7 @@ mod tests {
&cancel_token,
&hub,
&run_control,
- &mut last_applied_id,
+ &mut applied_ids,
&done,
)
.await
diff --git a/lib/crates/fabro-server/src/server/handler/worker_control.rs b/lib/crates/fabro-server/src/server/handler/worker_control.rs
index 682107748..d643a904e 100644
--- a/lib/crates/fabro-server/src/server/handler/worker_control.rs
+++ b/lib/crates/fabro-server/src/server/handler/worker_control.rs
@@ -1,6 +1,8 @@
use std::sync::Arc;
-use axum::extract::ws::{CloseFrame, Message as WsMessage, WebSocket, WebSocketUpgrade, close_code};
+use axum::extract::ws::{
+ CloseFrame, Message as WsMessage, WebSocket, WebSocketUpgrade, close_code,
+};
use fabro_interview::{
WORKER_CONTROL_INVALID_CURSOR_REASON, WORKER_CONTROL_PONG_TIMEOUT_REASON,
WORKER_CONTROL_WS_LIVENESS_TIMEOUT, WORKER_CONTROL_WS_PING_INTERVAL,
diff --git a/lib/crates/fabro-server/src/worker_control/local.rs b/lib/crates/fabro-server/src/worker_control/local.rs
index 51ea1421f..219ff18ec 100644
--- a/lib/crates/fabro-server/src/worker_control/local.rs
+++ b/lib/crates/fabro-server/src/worker_control/local.rs
@@ -23,8 +23,9 @@ pub(crate) struct LocalWorkerControlBus {
}
struct LocalRunControlStream {
- messages: VecDeque<LocalMessage>,
- notify: Arc<Notify>,
+ messages: VecDeque<LocalMessage>,
+ notify: Arc<Notify>,
+ has_trimmed: bool,
}
#[derive(Clone)]
@@ -60,9 +61,17 @@ impl LocalWorkerControlBus {
.streams
.lock()
.expect("worker control streams poisoned");
- let stream = streams
- .entry(run_id)
- .or_insert_with(LocalRunControlStream::new);
+ let stream = match cursor {
+ WorkerControlCursor::Start => streams
+ .entry(run_id)
+ .or_insert_with(LocalRunControlStream::new),
+ WorkerControlCursor::After(id) => streams.get_mut(&run_id).ok_or_else(|| {
+ WorkerControlBusError::invalid_cursor(
+ id.as_str(),
+ "message id is not retained for this run",
+ )
+ })?,
+ };
let next_sequence = stream.next_sequence_for_cursor(cursor)?;
(next_sequence, Arc::clone(&stream.notify))
};
@@ -90,8 +99,9 @@ impl LocalWorkerControlBus {
impl LocalRunControlStream {
fn new() -> Self {
Self {
- messages: VecDeque::new(),
- notify: Arc::new(Notify::new()),
+ messages: VecDeque::new(),
+ notify: Arc::new(Notify::new()),
+ has_trimmed: false,
}
}
@@ -100,6 +110,9 @@ impl LocalRunControlStream {
cursor: &WorkerControlCursor,
) -> Result<Option<u64>, WorkerControlBusError> {
match cursor {
+ WorkerControlCursor::Start if self.has_trimmed => Err(
+ WorkerControlBusError::invalid_cursor("start", "retained local stream is trimmed"),
+ ),
WorkerControlCursor::Start => Ok(self.messages.front().map(|message| message.sequence)),
WorkerControlCursor::After(id) => {
let requested_sequence = parse_local_sequence(id)?;
@@ -125,6 +138,7 @@ impl LocalRunControlStream {
fn trim_retained(&mut self) {
while self.messages.len() > LOCAL_WORKER_CONTROL_RETAINED_MESSAGES_PER_RUN {
self.messages.pop_front();
+ self.has_trimmed = true;
}
}
}
@@ -163,12 +177,19 @@ async fn local_subscription_task(
Some(stream) => {
// A `Start` subscriber that joined before any publish lazily
// adopts the first retained message as its cursor.
- if next_sequence.is_none() {
- next_sequence = stream.messages.front().map(|message| message.sequence);
- }
- match next_sequence {
- None => Some(Ok(Vec::new())),
- Some(next) => Some(collect_messages_from(&stream.messages, next)),
+ if next_sequence.is_none() && stream.has_trimmed {
+ Some(Err(WorkerControlBusError::invalid_cursor(
+ "start",
+ "retained local stream is trimmed",
+ )))
+ } else {
+ if next_sequence.is_none() {
+ next_sequence = stream.messages.front().map(|message| message.sequence);
+ }
+ match next_sequence {
+ None => Some(Ok(Vec::new())),
+ Some(next) => Some(collect_messages_from(&stream.messages, next)),
+ }
}
}
}
@@ -399,6 +420,27 @@ mod tests {
.await
.unwrap_err();
assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. }));
+
+ let err = bus
+ .subscribe(fixtures::RUN_1, WorkerControlCursor::Start)
+ .await
+ .unwrap_err();
+ assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. }));
+ }
+
+ #[tokio::test]
+ async fn after_cursor_for_unknown_run_does_not_create_stream() {
+ let bus = LocalWorkerControlBus::new();
+ let err = bus
+ .subscribe(
+ fixtures::RUN_1,
+ WorkerControlCursor::After(WorkerControlMessageId::new("local:1")),
+ )
+ .await
+ .unwrap_err();
+
+ assert!(matches!(err, WorkerControlBusError::InvalidCursor { .. }));
+ assert_eq!(bus.retained_len(&fixtures::RUN_1), 0);
}
#[tokio::test]

View file

@ -0,0 +1,28 @@
Implemented review cleanup after running the three parallel review agents.
Fixed:
- **Local bus replay correctness**
- `Start` subscriptions now fail with `InvalidCursor` if the local stream has already trimmed messages, rather than silently replaying a truncated stream.
- Active `Start` subscribers that joined before first publish also fail if messages were trimmed before they could establish a cursor.
- `After(id)` subscriptions for unknown runs no longer create empty run stream entries.
- Added tests for these cursor/cleanup cases.
- **Worker-side delivery-id dedupe**
- Replaced the single `last_applied_id` duplicate check with a bounded recent-id dedupe set.
- Uses FIFO + `HashSet`, capped at `2048` delivery ids.
- Still uses the last applied id for reconnect `?after=...`.
- Prevents repeated non-adjacent duplicate delivery ids from being applied twice without unbounded memory growth.
- **Formatting**
- `worker_control.rs` had rustfmt-only import formatting changes.
Already clean / previously addressed in current tree:
- Worker control route already checks invalid cursor before WebSocket upgrade and returns HTTP `410 Gone`.
- Invalid-cursor close reason and ping/pong constants were already shared through `fabro-interview`.
- Terminal-run check already uses `is_terminal()`.
Validation run:
- `cargo +nightly-2026-04-14 fmt --check --all` ✅
- `cargo nextest run -p fabro-server worker_control` ✅
- `cargo nextest run -p fabro-cli runner` ✅
- `cargo +nightly-2026-04-14 clippy -q -p fabro-server -p fabro-cli --all-targets -- -D warnings` ✅

View file

@ -0,0 +1,6 @@
{
"outcome": "succeeded",
"notes": "Stage completed: simplify_gpt",
"failure_reason": null,
"timestamp": "2026-05-27T23:29:21.638800Z"
}

View file

@ -0,0 +1,5 @@
{
"script": "git fetch origin main 2>&1 && git merge --no-edit --no-stat origin/main 2>&1 && cargo +nightly-2026-04-14 fmt --all 2>&1 && cargo dev docs refresh 2>&1 && cargo +nightly-2026-04-14 fmt --check --all 2>&1 && { command -v rg >/dev/null 2>&1 || { echo 'rg is required for verify'; exit 127; }; } && ! rg -n 'AuthMode::Disabled|RunAuthMethod|RunSubjectProvenance|\\bActorRef\\b|\\bActorKind\\b|AuthenticatedSubject|AuthenticatedService|AuthorizeRunScoped|AuthorizeRunBlob|AuthorizeStageArtifact|AuthorizeCommandLog|auth_method\\s*==\\s*\"disabled\"' lib/crates apps lib/packages docs/public/api-reference/fabro-api.yaml 2>&1 && cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings 2>&1 && cargo nextest run --workspace --status-level slow --profile ci 2>&1 && cargo dev docs check 2>&1 && bun install --frozen-lockfile 2>&1 && (cd apps/fabro-web && bun run typecheck) 2>&1 && (cd apps/fabro-web && bun run test) 2>&1 && (cd lib/packages/fabro-api-client && bun run typecheck) 2>&1 && cargo dev build -- -p fabro-cli --release 2>&1",
"command": "exec 2>&1\ngit fetch origin main 2>&1 && git merge --no-edit --no-stat origin/main 2>&1 && cargo +nightly-2026-04-14 fmt --all 2>&1 && cargo dev docs refresh 2>&1 && cargo +nightly-2026-04-14 fmt --check --all 2>&1 && { command -v rg >/dev/null 2>&1 || { echo 'rg is required for verify'; exit 127; }; } && ! rg -n 'AuthMode::Disabled|RunAuthMethod|RunSubjectProvenance|\\bActorRef\\b|\\bActorKind\\b|AuthenticatedSubject|AuthenticatedService|AuthorizeRunScoped|AuthorizeRunBlob|AuthorizeStageArtifact|AuthorizeCommandLog|auth_method\\s*==\\s*\"disabled\"' lib/crates apps lib/packages docs/public/api-reference/fabro-api.yaml 2>&1 && cargo +nightly-2026-04-14 clippy --workspace --all-targets -- -D warnings 2>&1 && cargo nextest run --workspace --status-level slow --profile ci 2>&1 && cargo dev docs check 2>&1 && bun install --frozen-lockfile 2>&1 && (cd apps/fabro-web && bun run typecheck) 2>&1 && (cd apps/fabro-web && bun run test) 2>&1 && (cd lib/packages/fabro-api-client && bun run typecheck) 2>&1 && cargo dev build -- -p fabro-cli --release 2>&1",
"language": "shell"
}