use std::collections::BTreeMap;
use std::fs;
use std::path::Path;
use serde_json::{Map, Value};
use super::super::{
Attempt, Board, Comment, Dependency, Handoff, Lane, Review, Task, Verdict, Workflow, Workspace,
WorkspaceKind,
};
use crate::error::Result;
use crate::ontology::Residue;
use crate::orchestration::codec::sqlite::{read_rows, table_exists};
use crate::sidecar::ms_to_rfc3339;
type Row = Map<String, Value>;
const TASK_COLUMNS: &[&str] = &[
"id",
"title",
"body",
"assignee",
"status",
"priority",
"created_by",
"created_at",
"started_at",
"completed_at",
"workspace_kind",
"workspace_path",
"branch_name",
"tenant",
"idempotency_key",
"result",
"skills",
"model_override",
"provider_override",
];
fn text(row: &Row, key: &str) -> Option<String> {
match row.get(key) {
Some(Value::String(s)) if !s.is_empty() => Some(s.clone()),
Some(Value::Number(n)) => Some(n.to_string()),
_ => None,
}
}
fn int(row: &Row, key: &str) -> Option<i64> {
match row.get(key) {
Some(Value::Number(n)) => n.as_i64(),
Some(Value::String(s)) => s.parse().ok(),
_ => None,
}
}
fn at(row: &Row, key: &str) -> Option<String> {
int(row, key).map(|secs| ms_to_rfc3339(secs.saturating_mul(1000)))
}
fn json_text(row: &Row, key: &str) -> Option<Value> {
text(row, key).and_then(|s| serde_json::from_str(&s).ok())
}
fn skills(row: &Row) -> Vec<String> {
match json_text(row, "skills") {
Some(Value::Array(items)) => items
.into_iter()
.filter_map(|v| v.as_str().map(str::to_string))
.collect(),
_ => text(row, "skills")
.map(|s| {
s.split(',')
.map(|p| p.trim().to_string())
.filter(|p| !p.is_empty())
.collect()
})
.unwrap_or_default(),
}
}
fn workspace(row: &Row) -> Workspace {
let kind = match text(row, "workspace_kind").as_deref() {
Some("scratch") | None => WorkspaceKind::Scratch,
Some("dir") => WorkspaceKind::Dir,
Some("worktree") => WorkspaceKind::Worktree,
Some(_) => WorkspaceKind::Unknown,
};
Workspace {
kind,
path: text(row, "workspace_path"),
branch: text(row, "branch_name"),
}
}
fn task(row: &Row) -> Option<Task> {
let id = text(row, "id")?;
let status = text(row, "status").unwrap_or_default();
let lane = Lane::parse(&status);
let mut residue = Residue::default();
for (key, value) in row {
if !TASK_COLUMNS.contains(&key.as_str()) && !value.is_null() {
residue.keep(key.clone(), value.clone());
}
}
if lane == Lane::Unknown {
residue.keep("status", Value::String(status));
}
let ws = workspace(row);
if ws.kind == WorkspaceKind::Unknown {
residue.keep("workspace_kind", row["workspace_kind"].clone());
}
Some(Task {
id,
title: text(row, "title").unwrap_or_default(),
body: text(row, "body"),
assignee: text(row, "assignee"),
lane,
priority: int(row, "priority").unwrap_or(0),
tenant: text(row, "tenant"),
idempotency_key: text(row, "idempotency_key"),
workspace: ws,
skills: skills(row),
model: text(row, "model_override"),
provider: text(row, "provider_override"),
created_by: text(row, "created_by"),
created_at: at(row, "created_at"),
started_at: at(row, "started_at"),
completed_at: at(row, "completed_at"),
result: text(row, "result"),
attempts: Vec::new(),
reviews: Vec::new(),
comments: Vec::new(),
residue,
})
}
fn attempt(row: &Row) -> Option<Attempt> {
let summary = text(row, "summary");
let metadata = json_text(row, "metadata");
Some(Attempt {
id: text(row, "id")?,
profile: text(row, "profile"),
step: text(row, "step_key"),
status: text(row, "status").unwrap_or_default(),
started_at: at(row, "started_at"),
ended_at: at(row, "ended_at"),
outcome: text(row, "outcome"),
handoff: (summary.is_some() || metadata.is_some()).then_some(Handoff { summary, metadata }),
error: text(row, "error"),
})
}
fn review(row: &Row) -> Option<Review> {
let verdict = match text(row, "kind")?.as_str() {
"review_requested" => Verdict::Requested,
"approved" | "review_approved" => Verdict::Approved,
"changes_requested" => Verdict::ChangesRequested,
"escalated" => Verdict::Escalated,
_ => return None,
};
let payload = json_text(row, "payload").unwrap_or(Value::Null);
let field = |k: &str| payload.get(k).and_then(Value::as_str).map(str::to_string);
Some(Review {
verdict,
by: field("reviewer")
.or_else(|| field("by"))
.or_else(|| field("profile")),
reason: field("reason"),
at: at(row, "created_at"),
})
}
fn rows(db: &Path, table: &str, sql: &str) -> Result<Vec<Row>> {
if !table_exists(db, table) {
return Ok(Vec::new());
}
Ok(read_rows(db, sql, &[])?.unwrap_or_default())
}
fn read_board(slug: &str, dir: &Path) -> Result<Board> {
let db = dir.join("kanban.db");
let mut tasks: BTreeMap<String, Task> = rows(&db, "tasks", "SELECT * FROM tasks")?
.iter()
.filter_map(task)
.map(|t| (t.id.clone(), t))
.collect();
for row in rows(&db, "task_runs", "SELECT * FROM task_runs ORDER BY id")? {
if let (Some(task_id), Some(a)) = (text(&row, "task_id"), attempt(&row)) {
if let Some(t) = tasks.get_mut(&task_id) {
t.attempts.push(a);
}
}
}
let events = "SELECT task_id, kind, payload, created_at FROM task_events ORDER BY id";
for row in rows(&db, "task_events", events)? {
if let (Some(task_id), Some(r)) = (text(&row, "task_id"), review(&row)) {
if let Some(t) = tasks.get_mut(&task_id) {
t.reviews.push(r);
}
}
}
let comments = "SELECT task_id, author, body, created_at FROM task_comments ORDER BY id";
for row in rows(&db, "task_comments", comments)? {
if let (Some(task_id), Some(author), Some(body)) = (
text(&row, "task_id"),
text(&row, "author"),
text(&row, "body"),
) {
if let Some(t) = tasks.get_mut(&task_id) {
t.comments.push(Comment {
author,
body,
at: at(&row, "created_at"),
});
}
}
}
let links = "SELECT parent_id, child_id FROM task_links ORDER BY parent_id, child_id";
let dependencies = rows(&db, "task_links", links)?
.iter()
.filter_map(|row| {
Some(Dependency {
parent: text(row, "parent_id")?,
child: text(row, "child_id")?,
})
})
.collect();
Ok(Board {
slug: slug.to_string(),
name: None,
root: dir.to_path_buf(),
tasks,
dependencies,
})
}
pub fn from_hermes(home: &Path) -> Result<Workflow> {
let mut boards = BTreeMap::new();
if home.join("kanban.db").is_file() {
boards.insert("default".to_string(), read_board("default", home)?);
}
if let Ok(entries) = fs::read_dir(home.join("kanban").join("boards")) {
let mut dirs: Vec<_> = entries.flatten().map(|e| e.path()).collect();
dirs.sort();
for dir in dirs {
let slug = dir
.file_name()
.and_then(|n| n.to_str())
.unwrap_or_default()
.to_string();
if slug.is_empty() || slug.starts_with('_') || !dir.join("kanban.db").is_file() {
continue;
}
boards.insert(slug.clone(), read_board(&slug, &dir)?);
}
}
Ok(Workflow {
root: home.to_path_buf(),
boards,
})
}