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

Commit f8ad320

Browse files
z23ccclaude
andcommitted
feat(db): add require_db, gaps table, epic fields, and init auto-import
- Replace try_open_db() with require_db() in epic.rs and task/mod.rs (DB unavailable = hard error, not silent degradation) - Add gaps table to schema.sql with full CRUD via GapRepo - Add auto_execute_pending, auto_execute_set_at, archived columns to epics table and Epic struct - Add max_epic_num() and max_task_num() DB queries in db_shim - Make flowctl init create flow.db and auto-import from existing MD files - Update all Epic struct constructors across crates for new fields Task: fn-2-db-as-sole-truth-db.1 Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent d727f3c commit f8ad320

15 files changed

Lines changed: 441 additions & 45 deletions

File tree

flowctl/crates/flowctl-cli/src/commands/admin/init.rs

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -60,6 +60,42 @@ pub fn cmd_init(json: bool) {
6060
}
6161
}
6262

63+
// Create/open flow.db (runs migrations automatically)
64+
let cwd = std::env::current_dir()
65+
.unwrap_or_else(|e| error_exit(&format!("Cannot get current dir: {e}")));
66+
match crate::commands::db_shim::open(&cwd) {
67+
Ok(conn) => {
68+
actions.push("flow.db ready".to_string());
69+
70+
// Auto-import from existing MD files if epics exist
71+
let epics_dir = flow_dir.join(EPICS_DIR);
72+
if epics_dir.is_dir() {
73+
let has_md_files = fs::read_dir(&epics_dir)
74+
.map(|entries| entries.flatten().any(|e| {
75+
e.file_name().to_string_lossy().ends_with(".md")
76+
}))
77+
.unwrap_or(false);
78+
79+
if has_md_files {
80+
match crate::commands::db_shim::reindex(&conn, &flow_dir, None) {
81+
Ok(result) => {
82+
actions.push(format!(
83+
"auto-imported {} epics, {} tasks from MD",
84+
result.epics_indexed, result.tasks_indexed
85+
));
86+
}
87+
Err(e) => {
88+
eprintln!("warning: auto-import failed: {e}");
89+
}
90+
}
91+
}
92+
}
93+
}
94+
Err(e) => {
95+
eprintln!("warning: DB creation failed: {e}");
96+
}
97+
}
98+
6399
// Build output
64100
let message = if actions.is_empty() {
65101
".flow/ already up to date".to_string()

flowctl/crates/flowctl-cli/src/commands/db_shim.rs

Lines changed: 79 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@
1414

1515
use std::path::{Path, PathBuf};
1616

17-
pub use flowctl_db::{DbError, ReindexResult};
17+
pub use flowctl_db::{DbError, GapRow, ReindexResult};
1818
pub use flowctl_db::metrics::{
1919
Bottleneck, DoraMetrics, EpicStats, Summary, TokenBreakdown, WeeklyTrend,
2020
};
@@ -61,10 +61,29 @@ pub fn open(working_dir: &Path) -> Result<Connection, DbError> {
6161
})
6262
}
6363

64+
/// Open DB connection with hard error on failure (DB must be available).
65+
/// This is the preferred entry point — replaces `try_open_db()` which
66+
/// returned `Option<Connection>` and silently degraded to MD fallback.
67+
pub fn require_db() -> Result<Connection, DbError> {
68+
let cwd = std::env::current_dir()
69+
.map_err(|e| DbError::StateDir(format!("cannot get current dir: {e}")))?;
70+
open(&cwd)
71+
}
72+
6473
pub fn cleanup(conn: &Connection) -> Result<u64, DbError> {
6574
block_on(flowctl_db::cleanup(&conn.inner()))
6675
}
6776

77+
/// Get the maximum epic number from DB.
78+
pub fn max_epic_num(conn: &Connection) -> Result<i64, DbError> {
79+
block_on(flowctl_db::max_epic_num(&conn.inner()))
80+
}
81+
82+
/// Get the maximum task number for an epic from DB.
83+
pub fn max_task_num(conn: &Connection, epic_id: &str) -> Result<i64, DbError> {
84+
block_on(flowctl_db::max_task_num(&conn.inner(), epic_id))
85+
}
86+
6887
pub fn reindex(
6988
conn: &Connection,
7089
flow_dir: &Path,
@@ -321,6 +340,65 @@ impl PhaseProgressRepo {
321340
}
322341
}
323342

343+
// ── Gap repository ────────────────────────────────────────────────
344+
345+
pub struct GapRepo(libsql::Connection);
346+
347+
impl GapRepo {
348+
pub fn new(conn: &Connection) -> Self {
349+
Self(conn.inner())
350+
}
351+
352+
pub fn add(
353+
&self,
354+
epic_id: &str,
355+
capability: &str,
356+
priority: &str,
357+
source: Option<&str>,
358+
task_id: Option<&str>,
359+
) -> Result<i64, DbError> {
360+
block_on(
361+
flowctl_db::GapRepo::new(self.0.clone())
362+
.add(epic_id, capability, priority, source, task_id),
363+
)
364+
}
365+
366+
pub fn list(
367+
&self,
368+
epic_id: &str,
369+
status: Option<&str>,
370+
) -> Result<Vec<GapRow>, DbError> {
371+
block_on(
372+
flowctl_db::GapRepo::new(self.0.clone())
373+
.list(epic_id, status),
374+
)
375+
}
376+
377+
pub fn remove(&self, id: i64) -> Result<(), DbError> {
378+
block_on(flowctl_db::GapRepo::new(self.0.clone()).remove(id))
379+
}
380+
381+
pub fn remove_all(&self, epic_id: &str) -> Result<u64, DbError> {
382+
block_on(flowctl_db::GapRepo::new(self.0.clone()).remove_all(epic_id))
383+
}
384+
385+
pub fn resolve(&self, id: i64, evidence: &str) -> Result<(), DbError> {
386+
block_on(flowctl_db::GapRepo::new(self.0.clone()).resolve(id, evidence))
387+
}
388+
389+
pub fn resolve_by_capability(
390+
&self,
391+
epic_id: &str,
392+
capability: &str,
393+
evidence: &str,
394+
) -> Result<(), DbError> {
395+
block_on(
396+
flowctl_db::GapRepo::new(self.0.clone())
397+
.resolve_by_capability(epic_id, capability, evidence),
398+
)
399+
}
400+
}
401+
324402
// ── Stats query ────────────────────────────────────────────────────
325403

326404
pub struct StatsQuery(libsql::Connection);

flowctl/crates/flowctl-cli/src/commands/epic.rs

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -170,7 +170,7 @@ fn validate_epic_id(id: &str) {
170170
/// Load epic document: DB first, markdown fallback.
171171
fn load_epic(epic_path: &Path, id: &str) -> frontmatter::Document<Epic> {
172172
// Try DB first.
173-
if let Some(conn) = try_open_db() {
173+
if let Ok(conn) = require_db() {
174174
let repo = crate::commands::db_shim::EpicRepo::new(&conn);
175175
if let Ok((epic, body)) = repo.get_with_body(id) {
176176
return frontmatter::Document { frontmatter: epic, body };
@@ -189,7 +189,7 @@ fn load_epic(epic_path: &Path, id: &str) -> frontmatter::Document<Epic> {
189189
/// Write an epic document: DB first, then export markdown.
190190
fn save_epic(epic_path: &Path, doc: &frontmatter::Document<Epic>) {
191191
// Write to DB.
192-
if let Some(conn) = try_open_db() {
192+
if let Ok(conn) = require_db() {
193193
let repo = crate::commands::db_shim::EpicRepo::new(&conn);
194194
if let Err(e) = repo.upsert_with_body(&doc.frontmatter, &doc.body) {
195195
eprintln!("warning: DB write failed for {}: {e}", doc.frontmatter.id);
@@ -205,15 +205,14 @@ fn save_epic(epic_path: &Path, doc: &frontmatter::Document<Epic>) {
205205
.unwrap_or_else(|e| error_exit(&format!("Failed to write {}: {e}", epic_path.display())));
206206
}
207207

208-
/// Try to open DB connection for SQLite dual-write.
209-
fn try_open_db() -> Option<crate::commands::db_shim::Connection> {
210-
let cwd = env::current_dir().ok()?;
211-
crate::commands::db_shim::open(&cwd).ok()
208+
/// Open DB connection (hard error if unavailable).
209+
fn require_db() -> Result<crate::commands::db_shim::Connection, crate::commands::db_shim::DbError> {
210+
crate::commands::db_shim::require_db()
212211
}
213212

214-
/// Upsert epic into SQLite if DB is available.
213+
/// Upsert epic into SQLite (hard error if DB unavailable).
215214
fn db_upsert_epic(epic: &Epic) {
216-
if let Some(conn) = try_open_db() {
215+
if let Ok(conn) = require_db() {
217216
let repo = crate::commands::db_shim::EpicRepo::new(&conn);
218217
let _ = repo.upsert(epic);
219218
}
@@ -363,6 +362,9 @@ fn cmd_create(title: &str, branch: &Option<String>, json_mode: bool) {
363362
default_impl: None,
364363
default_review: None,
365364
default_sync: None,
365+
auto_execute_pending: None,
366+
auto_execute_set_at: None,
367+
archived: false,
366368
file_path: Some(format!("epics/{epic_id}.md")),
367369
created_at: now,
368370
updated_at: now,
@@ -682,7 +684,7 @@ fn cmd_set_title(id: &str, new_title: &str, json_mode: bool) {
682684
let _ = fs::write(&task_path, serialized);
683685
}
684686
// SQLite update
685-
if let Some(conn) = try_open_db() {
687+
if let Ok(conn) = require_db() {
686688
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
687689
let _ = repo.upsert(&task_doc.frontmatter);
688690
}

flowctl/crates/flowctl-cli/src/commands/task/create.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,7 @@ use flowctl_core::types::{Domain, Task, EPICS_DIR, FLOW_DIR, TASKS_DIR};
1111

1212
use super::{
1313
create_task_spec, ensure_flow_exists, parse_domain, read_file_or_stdin, scan_max_task_id,
14-
try_open_db, write_task_doc,
14+
require_db, write_task_doc,
1515
};
1616

1717
#[allow(clippy::too_many_arguments)]
@@ -127,7 +127,7 @@ pub(super) fn cmd_task_create(
127127
write_task_doc(&flow_dir, &task_id, &doc);
128128

129129
// Upsert into SQLite if DB available
130-
if let Some(conn) = try_open_db() {
130+
if let Ok(conn) = require_db() {
131131
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
132132
let _ = repo.upsert(&task);
133133
}

flowctl/crates/flowctl-cli/src/commands/task/mod.rs

Lines changed: 7 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@ mod create;
55
mod mutate;
66
mod query;
77

8-
use std::env;
98
use std::fs;
109
use std::io::{self, Read as _};
1110
use std::path::{Path, PathBuf};
@@ -130,10 +129,9 @@ fn ensure_flow_exists() -> PathBuf {
130129
flow_dir
131130
}
132131

133-
/// Try to open a DB connection.
134-
fn try_open_db() -> Option<crate::commands::db_shim::Connection> {
135-
let cwd = env::current_dir().ok()?;
136-
crate::commands::db_shim::open(&cwd).ok()
132+
/// Open DB connection (hard error if unavailable).
133+
fn require_db() -> Result<crate::commands::db_shim::Connection, crate::commands::db_shim::DbError> {
134+
crate::commands::db_shim::require_db()
137135
}
138136

139137
/// Read file content, or read from stdin if path is "-".
@@ -199,7 +197,7 @@ fn create_task_spec(id: &str, title: &str, acceptance: Option<&str>) -> String {
199197

200198
/// Load a task: DB first, markdown fallback.
201199
fn load_task_md(_flow_dir: &Path, task_id: &str) -> Task {
202-
if let Some(conn) = try_open_db() {
200+
if let Ok(conn) = require_db() {
203201
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
204202
if let Ok(task) = repo.get(task_id) {
205203
return task;
@@ -220,7 +218,7 @@ fn load_task_md(_flow_dir: &Path, task_id: &str) -> Task {
220218

221219
/// Load an epic: DB first, markdown fallback.
222220
fn load_epic_md(_flow_dir: &Path, epic_id: &str) -> Option<Epic> {
223-
if let Some(conn) = try_open_db() {
221+
if let Ok(conn) = require_db() {
224222
let repo = crate::commands::db_shim::EpicRepo::new(&conn);
225223
if let Ok(epic) = repo.get(epic_id) {
226224
return Some(epic);
@@ -238,7 +236,7 @@ fn load_epic_md(_flow_dir: &Path, epic_id: &str) -> Option<Epic> {
238236

239237
/// Load task's full document (frontmatter + body): DB first, markdown fallback.
240238
fn load_task_doc(flow_dir: &Path, task_id: &str) -> frontmatter::Document<Task> {
241-
if let Some(conn) = try_open_db() {
239+
if let Ok(conn) = require_db() {
242240
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
243241
if let Ok((task, body)) = repo.get_with_body(task_id) {
244242
return frontmatter::Document {
@@ -261,7 +259,7 @@ fn load_task_doc(flow_dir: &Path, task_id: &str) -> frontmatter::Document<Task>
261259
/// Write a task document: DB first, then export markdown.
262260
fn write_task_doc(flow_dir: &Path, task_id: &str, doc: &frontmatter::Document<Task>) {
263261
// Write to DB.
264-
if let Some(conn) = try_open_db() {
262+
if let Ok(conn) = require_db() {
265263
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
266264
if let Err(e) = repo.upsert_with_body(&doc.frontmatter, &doc.body) {
267265
eprintln!("warning: DB write failed for {task_id}: {e}");

flowctl/crates/flowctl-cli/src/commands/task/mutate.rs

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -14,7 +14,7 @@ use flowctl_core::types::{Task, FLOW_DIR, TASKS_DIR};
1414

1515
use super::{
1616
clear_evidence_in_body, create_task_spec, ensure_flow_exists, find_dependents, load_epic_md,
17-
load_task_doc, scan_max_task_id, try_open_db, write_task_doc,
17+
load_task_doc, scan_max_task_id, require_db, write_task_doc,
1818
};
1919

2020
pub(super) fn cmd_task_reset(json_mode: bool, task_id: &str, cascade: bool) {
@@ -65,7 +65,7 @@ pub(super) fn cmd_task_reset(json_mode: bool, task_id: &str, cascade: bool) {
6565
write_task_doc(&flow_dir, task_id, &doc);
6666

6767
// Update DB
68-
if let Some(conn) = try_open_db() {
68+
if let Ok(conn) = require_db() {
6969
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
7070
let _ = repo.update_status(task_id, Status::Todo);
7171
// Clear runtime state by upserting a blank state
@@ -99,7 +99,7 @@ pub(super) fn cmd_task_reset(json_mode: bool, task_id: &str, cascade: bool) {
9999
dep_doc.body = clear_evidence_in_body(&dep_doc.body);
100100
write_task_doc(&flow_dir, dep_id, &dep_doc);
101101

102-
if let Some(conn) = try_open_db() {
102+
if let Ok(conn) = require_db() {
103103
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
104104
let _ = repo.update_status(dep_id, Status::Todo);
105105
let runtime_repo = crate::commands::db_shim::RuntimeRepo::new(&conn);
@@ -141,7 +141,7 @@ pub(super) fn cmd_task_skip(json_mode: bool, task_id: &str, reason: Option<&str>
141141
write_task_doc(&flow_dir, task_id, &doc);
142142

143143
// Update DB
144-
if let Some(conn) = try_open_db() {
144+
if let Ok(conn) = require_db() {
145145
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
146146
let _ = repo.update_status(task_id, Status::Skipped);
147147
}
@@ -238,7 +238,7 @@ pub(super) fn cmd_task_split(json_mode: bool, task_id: &str, titles: &str, chain
238238
};
239239
write_task_doc(&flow_dir, &sub_id, &sub_doc);
240240

241-
if let Some(conn) = try_open_db() {
241+
if let Ok(conn) = require_db() {
242242
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
243243
let _ = repo.upsert(&sub_task);
244244
}
@@ -252,7 +252,7 @@ pub(super) fn cmd_task_split(json_mode: bool, task_id: &str, titles: &str, chain
252252
orig_doc.frontmatter.updated_at = now;
253253
write_task_doc(&flow_dir, task_id, &orig_doc);
254254

255-
if let Some(conn) = try_open_db() {
255+
if let Ok(conn) = require_db() {
256256
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
257257
let _ = repo.update_status(task_id, Status::Skipped);
258258
}
@@ -379,7 +379,7 @@ pub(super) fn cmd_task_set_deps(json_mode: bool, task_id: &str, deps: &str) {
379379
doc.frontmatter.updated_at = Utc::now();
380380
write_task_doc(&flow_dir, task_id, &doc);
381381

382-
if let Some(conn) = try_open_db() {
382+
if let Ok(conn) = require_db() {
383383
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
384384
let _ = repo.upsert(&doc.frontmatter);
385385
}

flowctl/crates/flowctl-cli/src/commands/task/query.rs

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -9,7 +9,7 @@ use flowctl_core::id::is_task_id;
99

1010
use super::{
1111
ensure_flow_exists, load_epic_md, load_task_doc, load_task_md, patch_body_section,
12-
read_file_or_stdin, try_open_db, write_task_doc,
12+
read_file_or_stdin, require_db, write_task_doc,
1313
};
1414

1515
pub(super) fn cmd_task_set_spec(
@@ -41,7 +41,7 @@ pub(super) fn cmd_task_set_spec(
4141
doc.frontmatter.updated_at = Utc::now();
4242
write_task_doc(&flow_dir, task_id, &doc);
4343

44-
if let Some(conn) = try_open_db() {
44+
if let Ok(conn) = require_db() {
4545
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
4646
let _ = repo.upsert(&doc.frontmatter);
4747
}
@@ -75,7 +75,7 @@ pub(super) fn cmd_task_set_spec(
7575
doc.frontmatter.updated_at = Utc::now();
7676
write_task_doc(&flow_dir, task_id, &doc);
7777

78-
if let Some(conn) = try_open_db() {
78+
if let Ok(conn) = require_db() {
7979
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
8080
let _ = repo.upsert(&doc.frontmatter);
8181
}
@@ -133,7 +133,7 @@ pub(super) fn cmd_task_set_backend(
133133
doc.frontmatter.updated_at = Utc::now();
134134
write_task_doc(&flow_dir, task_id, &doc);
135135

136-
if let Some(conn) = try_open_db() {
136+
if let Ok(conn) = require_db() {
137137
let repo = crate::commands::db_shim::TaskRepo::new(&conn);
138138
let _ = repo.upsert(&doc.frontmatter);
139139
}

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

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,9 @@ fn make_test_epic(id: &str, title: &str) -> Epic {
2525
default_impl: None,
2626
default_review: None,
2727
default_sync: None,
28+
auto_execute_pending: None,
29+
auto_execute_set_at: None,
30+
archived: false,
2831
file_path: Some(format!("epics/{id}.md")),
2932
created_at: chrono::Utc::now(),
3033
updated_at: chrono::Utc::now(),

0 commit comments

Comments
 (0)