use crate::lfd::id::LfdId;
use crate::lfd::store::{ForkRun, ForkRunStatus, StoreError, StoreResult};
use crate::lfd::types::{
ActivationLog, ActivationOutcome, AgentRun, AgentStatus, ChatMemoryBlock, ChatMessage,
LivePrState, LivePullRequestState, PendingActivation, PullRequest, Repo, RepoEdge, RepoId,
Signal, Summary, Trigger, Wave, WaveCron, WaveMode, WaveRun, WaveRunSnapshot,
WaveRunStackStatus, WaveRunStatus, WaveStatus,
};
pub trait StoreRow {
fn text(&self, idx: usize) -> StoreResult<String>;
fn opt_text(&self, idx: usize) -> StoreResult<Option<String>>;
fn int(&self, idx: usize) -> StoreResult<i32>;
fn opt_int(&self, idx: usize) -> StoreResult<Option<i32>>;
fn bigint(&self, idx: usize) -> StoreResult<i64>;
fn opt_bigint(&self, idx: usize) -> StoreResult<Option<i64>>;
}
impl StoreRow for rusqlite::Row<'_> {
fn text(&self, idx: usize) -> StoreResult<String> {
Ok(self.get(idx)?)
}
fn opt_text(&self, idx: usize) -> StoreResult<Option<String>> {
Ok(self.get(idx)?)
}
fn int(&self, idx: usize) -> StoreResult<i32> {
Ok(self.get::<_, i64>(idx)? as i32)
}
fn opt_int(&self, idx: usize) -> StoreResult<Option<i32>> {
Ok(self.get::<_, Option<i64>>(idx)?.map(|v| v as i32))
}
fn bigint(&self, idx: usize) -> StoreResult<i64> {
Ok(self.get(idx)?)
}
fn opt_bigint(&self, idx: usize) -> StoreResult<Option<i64>> {
Ok(self.get(idx)?)
}
}
impl StoreRow for tokio_postgres::Row {
fn text(&self, idx: usize) -> StoreResult<String> {
Ok(self.try_get(idx)?)
}
fn opt_text(&self, idx: usize) -> StoreResult<Option<String>> {
Ok(self.try_get(idx)?)
}
fn int(&self, idx: usize) -> StoreResult<i32> {
Ok(self.try_get(idx)?)
}
fn opt_int(&self, idx: usize) -> StoreResult<Option<i32>> {
Ok(self.try_get(idx)?)
}
fn bigint(&self, idx: usize) -> StoreResult<i64> {
Ok(self.try_get(idx)?)
}
fn opt_bigint(&self, idx: usize) -> StoreResult<Option<i64>> {
Ok(self.try_get(idx)?)
}
}
pub fn unix_to_datetime(seconds: i64) -> time::OffsetDateTime {
time::OffsetDateTime::from_unix_timestamp(seconds).unwrap_or(time::OffsetDateTime::UNIX_EPOCH)
}
pub fn now_unix() -> i64 {
time::OffsetDateTime::now_utc().unix_timestamp()
}
pub fn parse_json_vec(value: &str) -> StoreResult<Vec<String>> {
serde_json::from_str::<Vec<String>>(value).map_err(StoreError::Serde)
}
pub fn serialize_pr(value: &Option<PullRequest>) -> StoreResult<Option<String>> {
match value {
Some(pr) => Ok(Some(serde_json::to_string(pr)?)),
None => Ok(None),
}
}
pub fn parse_pr(value: Option<String>) -> StoreResult<Option<PullRequest>> {
match value {
Some(raw) if !raw.trim().is_empty() => serde_json::from_str::<PullRequest>(&raw)
.map(Some)
.map_err(StoreError::Serde),
_ => Ok(None),
}
}
pub fn map_wave_row(row: &impl StoreRow) -> StoreResult<Wave> {
let direction = parse_json_vec(&row.text(3)?)?;
let area = parse_json_vec(&row.text(4)?)?;
let paused = row.int(5)? != 0;
let status_value = row.int(6)?;
let iteration = row.int(7)? as u32;
let cycle_start_iteration = row.int(8)? as u32;
let created_at = unix_to_datetime(row.bigint(9)?);
let workers = row.int(10)? as u32;
let mode_str = row.text(11)?;
let mode = mode_str.parse::<WaveMode>().unwrap_or_default();
let primary_flow = row.text(12)?;
let mut status = WaveStatus::from_i32(status_value);
if paused {
status = WaveStatus::Paused;
}
Ok(Wave {
id: LfdId::from_raw(row.text(0)?),
name: row.text(1)?,
repo: row.text(2)?,
mode,
primary_flow,
crons: Vec::new(),
direction,
area,
status,
iteration,
cycle_start_iteration,
created_at: Some(created_at),
workers,
})
}
pub fn map_wave_cron_row(row: &impl StoreRow) -> StoreResult<WaveCron> {
Ok(WaveCron {
id: LfdId::from_raw(row.text(0)?),
wave_id: LfdId::from_raw(row.text(1)?),
flow: row.text(2)?,
schedule: row.text(3)?,
last_triggered_at: row.opt_bigint(4)?,
created_at: row.opt_bigint(5)?.map(unix_to_datetime),
})
}
pub fn map_repo_row(row: &impl StoreRow) -> StoreResult<Repo> {
Ok(Repo {
path: row.text(0)?,
repo_id: RepoId::from_raw(row.text(1)?),
name: row.text(2)?,
added_at: unix_to_datetime(row.bigint(3)?),
})
}
pub fn map_repo_edge_row(row: &impl StoreRow) -> StoreResult<RepoEdge> {
Ok(RepoEdge {
parent_repo_id: RepoId::from_raw(row.text(0)?),
child_repo_id: RepoId::from_raw(row.text(1)?),
})
}
pub fn map_wave_run_row(row: &impl StoreRow) -> StoreResult<WaveRun> {
let started_at = unix_to_datetime(row.bigint(7)?);
let ended_at = row.opt_bigint(8)?;
let snapshot_direction = parse_json_vec(&row.text(12)?)?;
let snapshot_area = parse_json_vec(&row.text(13)?)?;
let snapshot_pr = parse_pr(row.opt_text(14)?)?;
let flow_parents = parse_json_vec(&row.text(15)?)?;
let execution_cursor = row.opt_text(16)?;
let activation_log_id = row.opt_text(17)?.map(LfdId::from_raw);
let parent_run_id = row.opt_text(18)?.map(LfdId::from_raw);
let parent_pr_number = row.opt_bigint(19)?.map(|value| value as u32);
let stack_position = row.int(20)? as u32;
let stack_group_id = row.text(21)?;
let stack_status = WaveRunStackStatus::from_i32(row.int(22)?);
let lineage_inferred = row.int(23)? != 0;
let target_branch = row.text(24)?;
let repair_of = row.opt_text(25)?.map(LfdId::from_raw);
let snapshot = WaveRunSnapshot {
repo: row.text(10)?,
flow: row.text(11)?,
direction: snapshot_direction,
area: snapshot_area,
};
Ok(WaveRun {
id: LfdId::from_raw(row.text(0)?),
wave_id: LfdId::from_raw(row.text(1)?),
snapshot,
iteration: row.int(2)? as u32,
step_index: row.int(3)? as u32,
status: WaveRunStatus::from_i32(row.int(4)?),
worktree: row.text(5)?,
branch: row.text(6)?,
started_at: Some(started_at),
ended_at: ended_at.map(unix_to_datetime),
error: row.opt_text(9)?,
flow_parents,
execution_cursor,
activation_log_id,
parent_run_id,
parent_pr_number,
stack_position,
stack_group_id,
stack_status,
lineage_inferred,
target_branch,
repair_of,
pr: snapshot_pr,
})
}
pub fn map_live_pr_state_row(row: &impl StoreRow) -> StoreResult<LivePullRequestState> {
Ok(LivePullRequestState {
repo_id: row.text(0)?,
pr_number: row.bigint(1)? as u32,
state: LivePrState::from_i32(row.int(2)?),
is_draft: row.int(3)? != 0,
head_ref: row.text(4)?,
head_sha: row.text(5)?,
base_ref: row.text(6)?,
updated_at: unix_to_datetime(row.bigint(7)?),
merged_at: row.opt_bigint(8)?.map(unix_to_datetime),
synced_at: unix_to_datetime(row.bigint(9)?),
})
}
pub fn map_trigger_row(row: &impl StoreRow) -> StoreResult<Trigger> {
let created_at = unix_to_datetime(row.bigint(6)?);
Ok(Trigger {
id: LfdId::from_raw(row.text(0)?),
wave_id: LfdId::from_raw(row.text(1)?),
signal: Signal::from_i32(row.int(2)?).expect("stored signal must be valid"),
flow: row.opt_text(3)?,
last_main_sha: row.opt_text(4)?,
last_triggered_at: row.opt_bigint(5)?,
created_at: Some(created_at),
enabled: row.int(7)? != 0,
source_wave_id: row.opt_text(8)?.map(LfdId::from_raw),
max_iterations: row.opt_int(9)?.map(|v| v as u32),
})
}
pub fn map_pending_activation_row(row: &impl StoreRow) -> StoreResult<PendingActivation> {
Ok(PendingActivation {
id: LfdId::from_raw(row.text(0)?),
wave_id: LfdId::from_raw(row.text(1)?),
trigger_id: row.opt_text(2)?.map(LfdId::from_raw),
reason: row.text(3)?,
from_sha: row.text(4)?,
to_sha: row.text(5)?,
queued_at: row.bigint(6)?,
target_branch: row.text(7)?,
})
}
pub fn map_activation_log_row(row: &impl StoreRow) -> StoreResult<ActivationLog> {
let outcome = ActivationOutcome::parse(&row.text(4)?)
.ok_or_else(|| StoreError::InvalidData("invalid activation log outcome".to_string()))?;
Ok(ActivationLog {
id: LfdId::from_raw(row.text(0)?),
wave_id: LfdId::from_raw(row.text(1)?),
trigger_id: row.opt_text(2)?.map(LfdId::from_raw),
reason: row.text(3)?,
outcome,
created_at: row.bigint(5)?,
})
}
pub fn map_fork_run_row(row: &impl StoreRow) -> StoreResult<ForkRun> {
let status = ForkRunStatus::from_i64(row.int(4)? as i64)
.ok_or_else(|| StoreError::InvalidData("invalid fork run status".to_string()))?;
Ok(ForkRun {
id: LfdId::from_raw(row.text(0)?),
wave_run_id: LfdId::from_raw(row.text(1)?),
step_index: row.int(2)? as u32,
branch_index: row.int(3)? as u32,
status,
worktree: row.text(5)?,
})
}
pub fn map_agent_row(row: &impl StoreRow) -> StoreResult<AgentRun> {
let started_at = unix_to_datetime(row.bigint(6)?);
let ended_at = row.opt_bigint(7)?;
let pid = row.opt_int(8)?;
let container_id = row.opt_text(9)?;
let wave_run_id = row.opt_text(4)?;
Ok(AgentRun {
id: LfdId::from_raw(row.text(0)?),
step: row.text(1)?,
repo: row.text(2)?,
worktree: row.text(3)?,
wave_run_id: wave_run_id.map(LfdId::from_raw),
status: AgentStatus::from_i32(row.int(5)?),
started_at: Some(started_at),
ended_at: ended_at.map(unix_to_datetime),
pid: pid.map(|v| v as u32),
container_id,
agent: row.text(10)?,
run_mode: row.text(11)?,
})
}
pub fn map_summary_row(row: &impl StoreRow) -> StoreResult<Summary> {
Ok(Summary {
id: LfdId::from_raw(row.text(0)?),
wave_id: LfdId::from_raw(row.text(1)?),
content: row.text(2)?,
source_hash: row.text(3)?,
token_budget: row.int(4)? as u32,
agent: row.text(5)?,
created_at: Some(unix_to_datetime(row.bigint(6)?)),
})
}
pub fn map_chat_memory_block_row(row: &impl StoreRow) -> StoreResult<ChatMemoryBlock> {
Ok(ChatMemoryBlock {
wave_id: LfdId::from_raw(row.text(0)?),
name: row.text(1)?,
content: row.text(2)?,
position: row.int(3)? as u32,
updated_at: Some(unix_to_datetime(row.bigint(4)?)),
})
}
pub fn map_chat_message_row(row: &impl StoreRow) -> StoreResult<ChatMessage> {
Ok(ChatMessage {
id: LfdId::from_raw(row.text(0)?),
wave_id: LfdId::from_raw(row.text(1)?),
role: row.text(2)?,
content: row.text(3)?,
created_at: unix_to_datetime(row.bigint(4)?),
})
}