Skip to content
This repository was archived by the owner on Apr 11, 2026. It is now read-only.

Commit 69fa893

Browse files
z23ccclaude
andcommitted
fix(flowctl): harden state writes with atomic rename and error propagation
[fn-110-workflow-reliability-and-performance.3] - state_write() now uses write-to-temp + fs::rename() for POSIX atomic writes, preventing partial/corrupt state files on crash - All 7 `let _ = state_write(...)` sites in lifecycle.rs now either propagate errors via `?` (done_task, block_task, restart_task, handle_task_failure) or log warnings (propagate_upstream_failure loop) - Added test verifying atomic write leaves no temp file and fully replaces prior content Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent 85b501e commit 69fa893

2 files changed

Lines changed: 62 additions & 12 deletions

File tree

flowctl/crates/flowctl-core/src/json_store.rs

Lines changed: 41 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -340,11 +340,18 @@ pub fn state_read(flow_dir: &Path, task_id: &str) -> Result<TaskState> {
340340
}
341341

342342
/// Write task runtime state to `.state/tasks/<id>.state.json`.
343+
///
344+
/// Uses write-to-temp + `fs::rename()` for atomic writes on POSIX,
345+
/// preventing partial/corrupt state files on crash or power loss.
343346
pub fn state_write(flow_dir: &Path, task_id: &str, state: &TaskState) -> Result<()> {
344347
ensure_dir(&state_tasks_dir(flow_dir))?;
345348
let path = task_state_path(flow_dir, task_id);
346349
let content = serde_json::to_string_pretty(state)?;
347-
fs::write(&path, content)?;
350+
351+
// Write to a temporary file in the same directory, then atomic rename
352+
let tmp_path = path.with_extension("state.json.tmp");
353+
fs::write(&tmp_path, &content)?;
354+
fs::rename(&tmp_path, &path)?;
348355
Ok(())
349356
}
350357

@@ -628,6 +635,39 @@ mod tests {
628635
assert!(read_back[1].resolved);
629636
}
630637

638+
#[test]
639+
fn test_state_write_atomic_no_corrupt() {
640+
// Verify that state_write uses atomic rename: if the original file
641+
// exists, a second write should fully replace it (no partial content).
642+
let tmp = TempDir::new().unwrap();
643+
let flow_dir = tmp.path();
644+
645+
// Write initial state
646+
let state1 = TaskState {
647+
status: Status::InProgress,
648+
assignee: Some("worker-1".to_string()),
649+
..Default::default()
650+
};
651+
state_write(flow_dir, "fn-1-test.1", &state1).unwrap();
652+
653+
// Write a second state that overwrites the first
654+
let state2 = TaskState {
655+
status: Status::Done,
656+
assignee: Some("worker-2".to_string()),
657+
..Default::default()
658+
};
659+
state_write(flow_dir, "fn-1-test.1", &state2).unwrap();
660+
661+
// Read back and verify the file contains only the second write
662+
let read_back = state_read(flow_dir, "fn-1-test.1").unwrap();
663+
assert_eq!(read_back.status, Status::Done);
664+
assert_eq!(read_back.assignee.as_deref(), Some("worker-2"));
665+
666+
// Verify no leftover .tmp file exists
667+
let tmp_path = task_state_path(flow_dir, "fn-1-test.1").with_extension("state.json.tmp");
668+
assert!(!tmp_path.exists(), "temporary file should be cleaned up by rename");
669+
}
670+
631671
#[test]
632672
fn test_not_found_errors() {
633673
let tmp = TempDir::new().unwrap();

flowctl/crates/flowctl-service/src/lifecycle.rs

Lines changed: 21 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -235,13 +235,17 @@ async fn propagate_upstream_failure(
235235
if let Ok(mut state) = flowctl_core::json_store::state_read(flow_dir, tid) {
236236
state.status = Status::UpstreamFailed;
237237
state.updated_at = Utc::now();
238-
let _ = flowctl_core::json_store::state_write(flow_dir, tid, &state);
238+
if let Err(e) = flowctl_core::json_store::state_write(flow_dir, tid, &state) {
239+
eprintln!("warning: failed to write upstream_failed state for {tid}: {e}");
240+
}
239241
} else {
240242
let state = flowctl_core::json_store::TaskState {
241243
status: Status::UpstreamFailed,
242244
..Default::default()
243245
};
244-
let _ = flowctl_core::json_store::state_write(flow_dir, tid, &state);
246+
if let Err(e) = flowctl_core::json_store::state_write(flow_dir, tid, &state) {
247+
eprintln!("warning: failed to write upstream_failed state for {tid}: {e}");
248+
}
245249
}
246250

247251
affected.push(tid.clone());
@@ -256,7 +260,7 @@ async fn handle_task_failure(
256260
flow_dir: &Path,
257261
task_id: &str,
258262
runtime: &Option<RuntimeState>,
259-
) -> (Status, Vec<String>) {
263+
) -> std::io::Result<(Status, Vec<String>)> {
260264
let max_retries = get_max_retries(flow_dir);
261265
let current_retry_count = runtime.as_ref().map(|r| r.retry_count).unwrap_or(0);
262266

@@ -276,18 +280,20 @@ async fn handle_task_failure(
276280
retry_count: new_retry_count,
277281
updated_at: Utc::now(),
278282
};
279-
let _ = flowctl_core::json_store::state_write(flow_dir, task_id, &task_state);
283+
flowctl_core::json_store::state_write(flow_dir, task_id, &task_state)
284+
.map_err(|e| std::io::Error::other(format!("failed to write retry state for {task_id}: {e}")))?;
280285

281-
(Status::UpForRetry, Vec::new())
286+
Ok((Status::UpForRetry, Vec::new()))
282287
} else {
283288
let task_state = flowctl_core::json_store::TaskState {
284289
status: Status::Failed,
285290
..Default::default()
286291
};
287-
let _ = flowctl_core::json_store::state_write(flow_dir, task_id, &task_state);
292+
flowctl_core::json_store::state_write(flow_dir, task_id, &task_state)
293+
.map_err(|e| std::io::Error::other(format!("failed to write failed state for {task_id}: {e}")))?;
288294

289295
let affected = propagate_upstream_failure(conn, flow_dir, task_id).await;
290-
(Status::Failed, affected)
296+
Ok((Status::Failed, affected))
291297
}
292298
}
293299

@@ -575,7 +581,8 @@ pub async fn done_task(
575581
retry_count: runtime.as_ref().map(|r| r.retry_count).unwrap_or(0),
576582
updated_at: now,
577583
};
578-
let _ = flowctl_core::json_store::state_write(flow_dir, &req.task_id, &task_state);
584+
flowctl_core::json_store::state_write(flow_dir, &req.task_id, &task_state)
585+
.map_err(|e| ServiceError::IoError(std::io::Error::other(e.to_string())))?;
579586
}
580587

581588
// Archive review receipt if present
@@ -649,7 +656,8 @@ pub async fn block_task(
649656
retry_count: existing.as_ref().map(|r| r.retry_count).unwrap_or(0),
650657
updated_at: Utc::now(),
651658
};
652-
let _ = flowctl_core::json_store::state_write(flow_dir, &req.task_id, &task_state);
659+
flowctl_core::json_store::state_write(flow_dir, &req.task_id, &task_state)
660+
.map_err(|e| ServiceError::IoError(std::io::Error::other(e.to_string())))?;
653661
}
654662

655663
Ok(BlockTaskResponse {
@@ -681,7 +689,8 @@ pub async fn fail_task(
681689
let reason_text = req.reason.unwrap_or_else(|| "Task failed".to_string());
682690

683691
let (final_status, upstream_failed_ids) =
684-
handle_task_failure(conn, flow_dir, &req.task_id, &runtime).await;
692+
handle_task_failure(conn, flow_dir, &req.task_id, &runtime).await
693+
.map_err(ServiceError::IoError)?;
685694

686695
let max_retries = get_max_retries(flow_dir);
687696
let retry_count = if final_status == Status::UpForRetry {
@@ -789,7 +798,8 @@ pub async fn restart_task(
789798
for tid in &to_reset {
790799
// Reset to blank JSON state
791800
let blank = flowctl_core::json_store::TaskState::default();
792-
let _ = flowctl_core::json_store::state_write(flow_dir, tid, &blank);
801+
flowctl_core::json_store::state_write(flow_dir, tid, &blank)
802+
.map_err(|e| ServiceError::IoError(std::io::Error::other(e.to_string())))?;
793803

794804
reset_ids.push(tid.clone());
795805
}

0 commit comments

Comments
 (0)