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

Commit b2fd876

Browse files
z23ccclaude
andcommitted
feat(service+daemon): complete RuntimeState migration to json_store
Service lifecycle.rs: zero DB calls for task/epic/runtime. All status updates, evidence, and runtime state now write to .state/tasks/*.json. Extended TaskState with baseline_rev, final_rev, retry_count, completed_at. Updated parity_test to read from json_store instead of DB. Updated daemon server tests to create data via json_store. 2 daemon tests still fail (flow_dir() in handlers uses cwd, not state). Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent fe218ac commit b2fd876

4 files changed

Lines changed: 137 additions & 130 deletions

File tree

flowctl/crates/flowctl-cli/tests/parity_test.rs

Lines changed: 6 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -367,29 +367,17 @@ fn setup_task(prefix: &str) -> (tempfile::TempDir, String) {
367367

368368
/// Read task status from the DB directly via async libSQL.
369369
#[allow(dead_code)]
370-
fn db_task_status(work_dir: &Path, task_id: &str) -> String {
371-
let rt = tokio::runtime::Builder::new_current_thread()
372-
.enable_all()
373-
.build()
374-
.unwrap();
375-
rt.block_on(async {
376-
let db = flowctl_db::open_async(work_dir).await.expect("open db");
377-
let conn = db.connect().expect("connect");
378-
let repo = flowctl_db::TaskRepo::new(conn);
379-
let task = repo.get(task_id).await.expect("get task");
380-
task.status.to_string()
381-
})
370+
fn json_task_status(work_dir: &Path, task_id: &str) -> String {
371+
let flow_dir = work_dir.join(".flow");
372+
let task = flowctl_core::json_store::task_read(&flow_dir, task_id).expect("read task");
373+
task.status.to_string()
382374
}
383375

384-
// Removed: rusqlite parity tests (fn-19 migration complete). The service
385-
// layer is now async libSQL end-to-end; the original parity placeholders
386-
// have been deleted.
387-
388376
#[test]
389377
fn parity_service_round_trip() {
390378
// Smoke test: create an epic+task via the CLI, then read it back via
391-
// the async libsql repo. Mirrors what the old parity tests checked.
379+
// json_store. Verifies CLI writes JSON files correctly.
392380
let (dir, task_id) = setup_task("parity-rt");
393-
let status = db_task_status(dir.path(), &task_id);
381+
let status = json_task_status(dir.path(), &task_id);
394382
assert_eq!(status, "todo", "newly created task should be todo");
395383
}

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

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -43,11 +43,19 @@ pub struct TaskState {
4343
#[serde(default, skip_serializing_if = "Option::is_none")]
4444
pub claimed_at: Option<DateTime<Utc>>,
4545
#[serde(default, skip_serializing_if = "Option::is_none")]
46+
pub completed_at: Option<DateTime<Utc>>,
47+
#[serde(default, skip_serializing_if = "Option::is_none")]
4648
pub evidence: Option<Evidence>,
4749
#[serde(default, skip_serializing_if = "Option::is_none")]
4850
pub blocked_reason: Option<String>,
4951
#[serde(default, skip_serializing_if = "Option::is_none")]
5052
pub duration_seconds: Option<u64>,
53+
#[serde(default, skip_serializing_if = "Option::is_none")]
54+
pub baseline_rev: Option<String>,
55+
#[serde(default, skip_serializing_if = "Option::is_none")]
56+
pub final_rev: Option<String>,
57+
#[serde(default)]
58+
pub retry_count: u32,
5159
#[serde(default = "Utc::now")]
5260
pub updated_at: DateTime<Utc>,
5361
}
@@ -58,9 +66,13 @@ impl Default for TaskState {
5866
status: Status::Todo,
5967
assignee: None,
6068
claimed_at: None,
69+
completed_at: None,
6170
evidence: None,
6271
blocked_reason: None,
6372
duration_seconds: None,
73+
baseline_rev: None,
74+
final_rev: None,
75+
retry_count: 0,
6476
updated_at: Utc::now(),
6577
}
6678
}

flowctl/crates/flowctl-daemon/src/server.rs

Lines changed: 41 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -378,18 +378,29 @@ mod tests {
378378

379379
#[tokio::test]
380380
async fn start_task_validates_transition() {
381-
// Setup: create epic + task in todo state, then start it (should succeed),
382-
// then try to start again from in_progress (should fail with CONFLICT).
383381
let (_tmp, runtime, event_bus) = test_setup();
384382
let (state, _cancel) = create_state(runtime, event_bus).await.unwrap();
385-
state.db.execute(
386-
"INSERT INTO epics (id, title, status, file_path, created_at, updated_at) VALUES ('fn-1', 'E', 'open', 'e.md', '2025-01-01T00:00:00Z', '2025-01-01T00:00:00Z')",
387-
(),
388-
).await.unwrap();
389-
state.db.execute(
390-
"INSERT INTO tasks (id, epic_id, title, status, domain, file_path, created_at, updated_at) VALUES ('fn-1.1', 'fn-1', 'T', 'todo', 'general', 't.md', '2025-01-01T00:00:00Z', '2025-01-01T00:00:00Z')",
391-
(),
392-
).await.unwrap();
383+
// Create epic + task via json_store
384+
use chrono::Utc;
385+
let epic = flowctl_core::types::Epic {
386+
schema_version: 1, id: "fn-1".into(), title: "E".into(),
387+
status: flowctl_core::types::EpicStatus::Open,
388+
branch_name: None, plan_review: flowctl_core::types::ReviewStatus::Unknown,
389+
completion_review: flowctl_core::types::ReviewStatus::Unknown,
390+
depends_on_epics: vec![], default_impl: None, default_review: None,
391+
default_sync: None, auto_execute_pending: None, auto_execute_set_at: None,
392+
archived: false, file_path: None, created_at: Utc::now(), updated_at: Utc::now(),
393+
};
394+
flowctl_core::json_store::epic_write(&state.flow_dir, &epic).unwrap();
395+
let task = flowctl_core::types::Task {
396+
schema_version: 1, id: "fn-1.1".into(), epic: "fn-1".into(), title: "T".into(),
397+
status: flowctl_core::Status::Todo, priority: None,
398+
domain: flowctl_core::types::Domain::General, depends_on: vec![],
399+
files: vec![], r#impl: None, review: None, sync: None, file_path: None,
400+
created_at: Utc::now(), updated_at: Utc::now(),
401+
};
402+
flowctl_core::json_store::task_write_definition(&state.flow_dir, &task).unwrap();
403+
flowctl_core::json_store::state_write(&state.flow_dir, "fn-1.1", &flowctl_core::json_store::TaskState::default()).unwrap();
393404
let app = build_router(state.clone());
394405

395406
// Start: todo → in_progress (should succeed)
@@ -429,14 +440,26 @@ mod tests {
429440
async fn done_task_rejects_from_todo() {
430441
let (_tmp, runtime, event_bus) = test_setup();
431442
let (state, _cancel) = create_state(runtime, event_bus).await.unwrap();
432-
state.db.execute(
433-
"INSERT INTO epics (id, title, status, file_path, created_at, updated_at) VALUES ('fn-2', 'E', 'open', 'e.md', '2025-01-01T00:00:00Z', '2025-01-01T00:00:00Z')",
434-
(),
435-
).await.unwrap();
436-
state.db.execute(
437-
"INSERT INTO tasks (id, epic_id, title, status, domain, file_path, created_at, updated_at) VALUES ('fn-2.1', 'fn-2', 'T', 'todo', 'general', 't.md', '2025-01-01T00:00:00Z', '2025-01-01T00:00:00Z')",
438-
(),
439-
).await.unwrap();
443+
use chrono::Utc;
444+
let epic = flowctl_core::types::Epic {
445+
schema_version: 1, id: "fn-2".into(), title: "E".into(),
446+
status: flowctl_core::types::EpicStatus::Open,
447+
branch_name: None, plan_review: flowctl_core::types::ReviewStatus::Unknown,
448+
completion_review: flowctl_core::types::ReviewStatus::Unknown,
449+
depends_on_epics: vec![], default_impl: None, default_review: None,
450+
default_sync: None, auto_execute_pending: None, auto_execute_set_at: None,
451+
archived: false, file_path: None, created_at: Utc::now(), updated_at: Utc::now(),
452+
};
453+
flowctl_core::json_store::epic_write(&state.flow_dir, &epic).unwrap();
454+
let task = flowctl_core::types::Task {
455+
schema_version: 1, id: "fn-2.1".into(), epic: "fn-2".into(), title: "T".into(),
456+
status: flowctl_core::Status::Todo, priority: None,
457+
domain: flowctl_core::types::Domain::General, depends_on: vec![],
458+
files: vec![], r#impl: None, review: None, sync: None, file_path: None,
459+
created_at: Utc::now(), updated_at: Utc::now(),
460+
};
461+
flowctl_core::json_store::task_write_definition(&state.flow_dir, &task).unwrap();
462+
flowctl_core::json_store::state_write(&state.flow_dir, "fn-2.1", &flowctl_core::json_store::TaskState::default()).unwrap();
440463
let app = build_router(state);
441464

442465
// done from todo → should be rejected

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

Lines changed: 78 additions & 94 deletions
Original file line numberDiff line numberDiff line change
@@ -118,10 +118,20 @@ async fn load_epic(_conn: Option<&Connection>, flow_dir: &Path, epic_id: &str) -
118118
flowctl_core::json_store::epic_read(flow_dir, epic_id).ok()
119119
}
120120

121-
async fn get_runtime(conn: Option<&Connection>, _flow_dir: &Path, task_id: &str) -> Option<RuntimeState> {
122-
let conn = conn?;
123-
let repo = flowctl_db::RuntimeRepo::new(conn.clone());
124-
repo.get(task_id).await.ok().flatten()
121+
async fn get_runtime(_conn: Option<&Connection>, flow_dir: &Path, task_id: &str) -> Option<RuntimeState> {
122+
// Read from JSON state and convert to RuntimeState for compatibility
123+
let state = flowctl_core::json_store::state_read(flow_dir, task_id).ok()?;
124+
Some(RuntimeState {
125+
task_id: task_id.to_string(),
126+
assignee: state.assignee,
127+
claimed_at: state.claimed_at,
128+
completed_at: state.completed_at,
129+
duration_secs: state.duration_seconds,
130+
blocked_reason: state.blocked_reason,
131+
baseline_rev: state.baseline_rev,
132+
final_rev: state.final_rev,
133+
retry_count: state.retry_count,
134+
})
125135
}
126136

127137
/// Load all tasks for an epic from JSON files.
@@ -221,10 +231,17 @@ async fn propagate_upstream_failure(
221231
continue;
222232
}
223233

224-
// Update DB (sole source of truth)
225-
if let Some(conn) = conn {
226-
let task_repo = flowctl_db::TaskRepo::new(conn.clone());
227-
let _ = task_repo.update_status(tid, Status::UpstreamFailed).await;
234+
// Update JSON state
235+
if let Ok(mut state) = flowctl_core::json_store::state_read(flow_dir, tid) {
236+
state.status = Status::UpstreamFailed;
237+
state.updated_at = Utc::now();
238+
let _ = flowctl_core::json_store::state_write(flow_dir, tid, &state);
239+
} else {
240+
let state = flowctl_core::json_store::TaskState {
241+
status: Status::UpstreamFailed,
242+
..Default::default()
243+
};
244+
let _ = flowctl_core::json_store::state_write(flow_dir, tid, &state);
228245
}
229246

230247
affected.push(tid.clone());
@@ -246,31 +263,28 @@ async fn handle_task_failure(
246263
if max_retries > 0 && current_retry_count < max_retries {
247264
let new_retry_count = current_retry_count + 1;
248265

249-
if let Some(conn) = conn {
250-
let task_repo = flowctl_db::TaskRepo::new(conn.clone());
251-
let _ = task_repo.update_status(task_id, Status::UpForRetry).await;
252-
253-
let runtime_repo = flowctl_db::RuntimeRepo::new(conn.clone());
254-
let rt = RuntimeState {
255-
task_id: task_id.to_string(),
256-
assignee: runtime.as_ref().and_then(|r| r.assignee.clone()),
257-
claimed_at: None,
258-
completed_at: None,
259-
duration_secs: None,
260-
blocked_reason: None,
261-
baseline_rev: runtime.as_ref().and_then(|r| r.baseline_rev.clone()),
262-
final_rev: None,
263-
retry_count: new_retry_count,
264-
};
265-
let _ = runtime_repo.upsert(&rt).await;
266-
}
266+
let task_state = flowctl_core::json_store::TaskState {
267+
status: Status::UpForRetry,
268+
assignee: runtime.as_ref().and_then(|r| r.assignee.clone()),
269+
claimed_at: None,
270+
completed_at: None,
271+
evidence: None,
272+
blocked_reason: None,
273+
duration_seconds: None,
274+
baseline_rev: runtime.as_ref().and_then(|r| r.baseline_rev.clone()),
275+
final_rev: None,
276+
retry_count: new_retry_count,
277+
updated_at: Utc::now(),
278+
};
279+
let _ = flowctl_core::json_store::state_write(flow_dir, task_id, &task_state);
267280

268281
(Status::UpForRetry, Vec::new())
269282
} else {
270-
if let Some(conn) = conn {
271-
let task_repo = flowctl_db::TaskRepo::new(conn.clone());
272-
let _ = task_repo.update_status(task_id, Status::Failed).await;
273-
}
283+
let task_state = flowctl_core::json_store::TaskState {
284+
status: Status::Failed,
285+
..Default::default()
286+
};
287+
let _ = flowctl_core::json_store::state_write(flow_dir, task_id, &task_state);
274288

275289
let affected = propagate_upstream_failure(conn, flow_dir, task_id).await;
276290
(Status::Failed, affected)
@@ -373,13 +387,14 @@ pub async fn start_task(
373387
Some(now)
374388
};
375389

376-
let runtime_state = RuntimeState {
377-
task_id: req.task_id.clone(),
390+
let task_state = flowctl_core::json_store::TaskState {
391+
status: Status::InProgress,
378392
assignee: Some(new_assignee),
379393
claimed_at,
380394
completed_at: None,
381-
duration_secs: None,
395+
evidence: None,
382396
blocked_reason: None,
397+
duration_seconds: None,
383398
baseline_rev: existing_rt
384399
.as_ref()
385400
.and_then(|rt| rt.baseline_rev.clone()),
@@ -388,21 +403,12 @@ pub async fn start_task(
388403
.as_ref()
389404
.map(|rt| rt.retry_count)
390405
.unwrap_or(0),
406+
updated_at: Utc::now(),
391407
};
392408

393-
// Write SQLite (authoritative)
394-
if let Some(conn) = conn {
395-
let task_repo = flowctl_db::TaskRepo::new(conn.clone());
396-
task_repo
397-
.update_status(&req.task_id, Status::InProgress)
398-
.await
399-
.map_err(ServiceError::from)?;
400-
let runtime_repo = flowctl_db::RuntimeRepo::new(conn.clone());
401-
runtime_repo
402-
.upsert(&runtime_state)
403-
.await
404-
.map_err(ServiceError::from)?;
405-
}
409+
// Write to JSON state file
410+
flowctl_core::json_store::state_write(flow_dir, &req.task_id, &task_state)
411+
.map_err(|e| ServiceError::IoError(std::io::Error::new(std::io::ErrorKind::Other, e.to_string())))?;
406412

407413
Ok(StartTaskResponse {
408414
task_id: req.task_id,
@@ -547,34 +553,29 @@ pub async fn done_task(
547553
let tests = to_list(evidence_obj.get("tests"));
548554
let prs = to_list(evidence_obj.get("prs"));
549555

550-
// Write SQLite (sole source of truth)
551-
if let Some(conn) = conn {
552-
let task_repo = flowctl_db::TaskRepo::new(conn.clone());
553-
let _ = task_repo.update_status(&req.task_id, Status::Done).await;
554-
555-
let runtime_repo = flowctl_db::RuntimeRepo::new(conn.clone());
556+
// Write to JSON state file
557+
{
556558
let now = Utc::now();
557-
let rt = RuntimeState {
558-
task_id: req.task_id.clone(),
559+
let ev = Evidence {
560+
commits: commits.clone(),
561+
tests: tests.clone(),
562+
prs: prs.clone(),
563+
..Evidence::default()
564+
};
565+
let task_state = flowctl_core::json_store::TaskState {
566+
status: Status::Done,
559567
assignee: runtime.as_ref().and_then(|r| r.assignee.clone()),
560568
claimed_at: runtime.as_ref().and_then(|r| r.claimed_at),
561569
completed_at: Some(now),
562-
duration_secs: duration_seconds,
570+
evidence: Some(ev),
563571
blocked_reason: None,
572+
duration_seconds: duration_seconds.map(|d| d as u64),
564573
baseline_rev: runtime.as_ref().and_then(|r| r.baseline_rev.clone()),
565574
final_rev: runtime.as_ref().and_then(|r| r.final_rev.clone()),
566575
retry_count: runtime.as_ref().map(|r| r.retry_count).unwrap_or(0),
576+
updated_at: now,
567577
};
568-
let _ = runtime_repo.upsert(&rt).await;
569-
570-
let ev = Evidence {
571-
commits: commits.clone(),
572-
tests: tests.clone(),
573-
prs: prs.clone(),
574-
..Evidence::default()
575-
};
576-
let evidence_repo = flowctl_db::EvidenceRepo::new(conn.clone());
577-
let _ = evidence_repo.upsert(&req.task_id, &ev).await;
578+
let _ = flowctl_core::json_store::state_write(flow_dir, &req.task_id, &task_state);
578579
}
579580

580581
// Archive review receipt if present
@@ -632,25 +633,23 @@ pub async fn block_task(
632633
));
633634
}
634635

635-
// Write SQLite (authoritative)
636-
if let Some(conn) = conn {
637-
let task_repo = flowctl_db::TaskRepo::new(conn.clone());
638-
let _ = task_repo.update_status(&req.task_id, Status::Blocked).await;
639-
640-
let runtime_repo = flowctl_db::RuntimeRepo::new(conn.clone());
641-
let existing = runtime_repo.get(&req.task_id).await.ok().flatten();
642-
let rt = RuntimeState {
643-
task_id: req.task_id.clone(),
636+
// Write to JSON state file
637+
{
638+
let existing = flowctl_core::json_store::state_read(flow_dir, &req.task_id).ok();
639+
let task_state = flowctl_core::json_store::TaskState {
640+
status: Status::Blocked,
644641
assignee: existing.as_ref().and_then(|r| r.assignee.clone()),
645642
claimed_at: existing.as_ref().and_then(|r| r.claimed_at),
646643
completed_at: None,
647-
duration_secs: None,
644+
evidence: existing.as_ref().and_then(|r| r.evidence.clone()),
648645
blocked_reason: Some(reason.clone()),
646+
duration_seconds: None,
649647
baseline_rev: existing.as_ref().and_then(|r| r.baseline_rev.clone()),
650648
final_rev: None,
651649
retry_count: existing.as_ref().map(|r| r.retry_count).unwrap_or(0),
650+
updated_at: Utc::now(),
652651
};
653-
let _ = runtime_repo.upsert(&rt).await;
652+
let _ = flowctl_core::json_store::state_write(flow_dir, &req.task_id, &task_state);
654653
}
655654

656655
Ok(BlockTaskResponse {
@@ -788,24 +787,9 @@ pub async fn restart_task(
788787
// Execute reset
789788
let mut reset_ids = Vec::new();
790789
for tid in &to_reset {
791-
if let Some(conn) = conn {
792-
let task_repo = flowctl_db::TaskRepo::new(conn.clone());
793-
let _ = task_repo.update_status(tid, Status::Todo).await;
794-
795-
let runtime_repo = flowctl_db::RuntimeRepo::new(conn.clone());
796-
let rt = RuntimeState {
797-
task_id: tid.clone(),
798-
assignee: None,
799-
claimed_at: None,
800-
completed_at: None,
801-
duration_secs: None,
802-
blocked_reason: None,
803-
baseline_rev: None,
804-
final_rev: None,
805-
retry_count: 0,
806-
};
807-
let _ = runtime_repo.upsert(&rt).await;
808-
}
790+
// Reset to blank JSON state
791+
let blank = flowctl_core::json_store::TaskState::default();
792+
let _ = flowctl_core::json_store::state_write(flow_dir, tid, &blank);
809793

810794
reset_ids.push(tid.clone());
811795
}

0 commit comments

Comments
 (0)