use std::collections::HashMap;
use std::path::Path;
use std::sync::{Arc, Mutex};
use rusqlite::{params, Connection, OptionalExtension, ToSql};
use crate::lfd::attention::{queue_block_attention_item, queue_block_from_attention};
use crate::lfd::id::LfdId;
use crate::lfd::sessions::types::{
PersistedSessionEvent, Session, SessionConfig, SessionEvent, SessionStatus,
};
use crate::lfd::store::catalog::{
list_agent_history_query, list_triggers_query, list_wave_runs_query, list_waves_query, sql,
Query, SqlDialect,
};
use crate::lfd::store::rows::{
map_activation_log_row, map_agent_row, map_chat_memory_block_row, map_chat_message_row,
map_fork_run_row, map_live_pr_state_row, map_pending_activation_row, map_repo_edge_row,
map_repo_row, map_summary_row, map_trigger_row, map_wave_cron_row, map_wave_row,
map_wave_run_row, now_unix, serialize_pr,
};
use crate::lfd::store::token_crypto;
use crate::lfd::store::{ForkRun, ForkRunStatus, SessionFilters, StoreError, StoreResult};
use crate::lfd::types::{
ActivationLog, AgentRun, AgentStatus, AttentionItem, AttentionKind, AttentionStatus,
ChatMemoryBlock, ChatMessage, LivePullRequestState, PendingActivation, QueueBlock,
QueueMergeEvent, Repo, RepoEdge, RepoId, Summary, TerminalSession, TerminalSessionStatus,
Trigger, Wave, WaveCron, WaveRun, WaveRunStatus, WaveStatus,
};
#[derive(Debug, Clone)]
pub struct SqliteStore {
conn: Arc<Mutex<Connection>>,
}
fn migrate_plaintext_provider_tokens(conn: &mut Connection) -> StoreResult<()> {
let mut scan = conn.prepare(
"SELECT provider, access_token, refresh_token
FROM provider_tokens
WHERE encrypted = 0",
)?;
let rows = scan.query_map([], |row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, Option<String>>(2)?,
))
})?;
let mut pending = Vec::new();
for row in rows {
pending.push(row?);
}
drop(scan);
if pending.is_empty() {
return Ok(());
}
let tx = conn.transaction()?;
for (provider, access_token, refresh_token) in pending {
let encrypted_access = token_crypto::encrypt_token(&access_token).map_err(|error| {
StoreError::InvalidData(format!(
"failed to encrypt existing access token for provider '{provider}': {error}"
))
})?;
let encrypted_refresh =
token_crypto::encrypt_optional(refresh_token.as_deref()).map_err(|error| {
StoreError::InvalidData(format!(
"failed to encrypt existing refresh token for provider '{provider}': {error}"
))
})?;
tx.execute(
"UPDATE provider_tokens
SET access_token = ?1,
refresh_token = ?2,
encrypted = 1
WHERE provider = ?3",
params![encrypted_access, encrypted_refresh, provider],
)?;
}
tx.commit()?;
Ok(())
}
type TokenRow = (
String,
String,
Option<String>,
Option<i64>,
Option<String>,
i64,
String,
bool,
);
fn read_token_row(row: &rusqlite::Row) -> rusqlite::Result<TokenRow> {
Ok((
row.get(0)?,
row.get(1)?,
row.get(2)?,
row.get(3)?,
row.get(4)?,
row.get(5)?,
row.get(6)?,
row.get(7)?,
))
}
fn decrypt_token_row(row: TokenRow) -> StoreResult<super::ProviderToken> {
let (provider, access_token, refresh_token, expires_at, login, updated_at, ct, encrypted) = row;
let access_token =
token_crypto::decrypt_if_needed(&access_token, encrypted).map_err(|error| {
StoreError::InvalidData(format!(
"failed to decrypt access token for provider '{provider}': {error}"
))
})?;
let refresh_token = refresh_token
.as_deref()
.map(|token| token_crypto::decrypt_if_needed(token, encrypted))
.transpose()
.map_err(|error| {
StoreError::InvalidData(format!(
"failed to decrypt refresh token for provider '{provider}': {error}"
))
})?;
Ok(super::ProviderToken {
provider,
access_token,
refresh_token,
expires_at,
login,
updated_at,
credential_type: super::CredentialType::from_db(&ct),
})
}
fn read_secrets_config_row(row: &rusqlite::Row) -> StoreResult<super::SecretsProviderConfig> {
Ok(super::SecretsProviderConfig {
provider: row.get(0)?,
project: row.get(1)?,
config: row.get(2)?,
updated_at: row.get(3)?,
})
}
impl SqliteStore {
fn sql(query: Query) -> &'static str {
sql(query, SqlDialect::Sqlite)
}
pub fn new(path: &Path) -> StoreResult<Self> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).map_err(|err| {
StoreError::InvalidData(format!("failed to create db dir: {err}"))
})?;
}
let mut conn = Connection::open(path)?;
conn.execute_batch(
"PRAGMA journal_mode = WAL; PRAGMA busy_timeout = 5000; PRAGMA foreign_keys = ON;",
)?;
super::migrations::apply_sqlite(&conn)?;
migrate_plaintext_provider_tokens(&mut conn)?;
Ok(Self {
conn: Arc::new(Mutex::new(conn)),
})
}
fn read_waves(&self, repo: Option<&str>) -> StoreResult<Vec<Wave>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let query = Self::sql(list_waves_query(repo.is_some()));
let params: Vec<Box<dyn ToSql>> = if let Some(repo) = repo {
vec![Box::new(repo.to_string())]
} else {
vec![]
};
let mut stmt = conn.prepare(query)?;
let params_iter = params.iter().map(|v| v.as_ref() as &dyn ToSql);
let rows = stmt.query_map(rusqlite::params_from_iter(params_iter), |row| {
Ok(map_wave_row(row))
})?;
let mut waves = Vec::new();
for wave in rows {
waves.push(wave??);
}
Ok(waves)
}
fn read_wave_crons<P>(&self, query: Query, params: P) -> StoreResult<Vec<WaveCron>>
where
P: rusqlite::Params,
{
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(query))?;
let rows = stmt.query_map(params, |row| Ok(map_wave_cron_row(row)))?;
let mut crons = Vec::new();
for cron in rows {
crons.push(cron??);
}
Ok(crons)
}
fn upsert_wave(&self, wave: &Wave) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let direction_json = serde_json::to_string(wave.direction())?;
let area_json = serde_json::to_string(wave.area())?;
let created_at = wave
.created_at()
.map(|dt| dt.unix_timestamp())
.unwrap_or_else(now_unix);
conn.execute(
Self::sql(Query::UpsertWave),
params![
wave.id(),
wave.name(),
wave.repo(),
direction_json,
area_json,
if wave.status() == WaveStatus::Paused {
1i64
} else {
0i64
},
wave.status().as_i32() as i64,
wave.iteration() as i64,
wave.cycle_start_iteration() as i64,
created_at,
wave.workers as i64,
wave.mode().as_str(),
wave.primary_flow(),
],
)?;
Ok(())
}
fn map_session_row(row: &rusqlite::Row<'_>) -> Result<Session, rusqlite::Error> {
let config_text: String = row.get(5)?;
let config: SessionConfig = serde_json::from_str(&config_text).map_err(|err| {
rusqlite::Error::FromSqlConversionFailure(5, rusqlite::types::Type::Text, Box::new(err))
})?;
Ok(Session {
id: row.get(0)?,
harness: row.get(1)?,
status: SessionStatus::from_i32(row.get::<_, i64>(2)? as i32),
wave_run_id: row.get(3)?,
provider_session_id: row.get(4)?,
config,
created_at: crate::lfd::store::rows::unix_to_datetime(row.get(6)?),
ended_at: row
.get::<_, Option<i64>>(7)?
.map(crate::lfd::store::rows::unix_to_datetime),
})
}
fn map_terminal_session_row(
row: &rusqlite::Row<'_>,
) -> Result<TerminalSession, rusqlite::Error> {
let argv_text: String = row.get(6)?;
let env_text: String = row.get(7)?;
let argv: Vec<String> = serde_json::from_str(&argv_text).map_err(|err| {
rusqlite::Error::FromSqlConversionFailure(6, rusqlite::types::Type::Text, Box::new(err))
})?;
let env = serde_json::from_str(&env_text).map_err(|err| {
rusqlite::Error::FromSqlConversionFailure(7, rusqlite::types::Type::Text, Box::new(err))
})?;
Ok(TerminalSession {
id: row.get(0)?,
wave_id: row.get(1)?,
wave_run_id: row.get(2)?,
step: row.get(3)?,
agent: row.get(4)?,
cwd: row.get(5)?,
argv,
env,
source: row.get(8)?,
tmux_name: row.get(9)?,
status: TerminalSessionStatus::from_i32(row.get::<_, i64>(10)? as i32),
completion_token: row.get(11)?,
created_at: crate::lfd::store::rows::unix_to_datetime(row.get(12)?),
attached_at: row
.get::<_, Option<i64>>(13)?
.map(crate::lfd::store::rows::unix_to_datetime),
started_at: row
.get::<_, Option<i64>>(14)?
.map(crate::lfd::store::rows::unix_to_datetime),
completed_at: row
.get::<_, Option<i64>>(15)?
.map(crate::lfd::store::rows::unix_to_datetime),
})
}
}
impl SqliteStore {
pub fn health_check(&self) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(Self::sql(Query::HealthCheck), [])?;
Ok(())
}
pub fn schema_version(&self) -> StoreResult<String> {
let conn = self.conn.lock().expect("store mutex poisoned");
super::migrations::latest_version_sqlite(&conn)
}
pub fn get_provider_token(&self, provider: &str) -> StoreResult<Option<super::ProviderToken>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT provider, access_token, refresh_token, expires_at, login, updated_at, credential_type, encrypted
FROM provider_tokens WHERE provider = ?1",
)?;
let row = stmt
.query_row(params![provider], read_token_row)
.optional()?;
row.map(decrypt_token_row).transpose()
}
pub fn upsert_provider_token(&self, token: &super::ProviderToken) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let encrypted_access =
token_crypto::encrypt_token(&token.access_token).map_err(|error| {
StoreError::InvalidData(format!(
"failed to encrypt access token for provider '{}': {error}",
token.provider
))
})?;
let encrypted_refresh = token_crypto::encrypt_optional(token.refresh_token.as_deref())
.map_err(|error| {
StoreError::InvalidData(format!(
"failed to encrypt refresh token for provider '{}': {error}",
token.provider
))
})?;
conn.execute(
"INSERT INTO provider_tokens (provider, access_token, refresh_token, expires_at, login, updated_at, credential_type, encrypted)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, 1)
ON CONFLICT(provider) DO UPDATE SET
access_token = excluded.access_token,
refresh_token = excluded.refresh_token,
expires_at = excluded.expires_at,
login = excluded.login,
updated_at = excluded.updated_at,
credential_type = excluded.credential_type,
encrypted = excluded.encrypted",
params![
token.provider,
encrypted_access,
encrypted_refresh,
token.expires_at,
token.login,
token.updated_at,
token.credential_type.as_str(),
],
)?;
Ok(())
}
pub fn delete_provider_token(&self, provider: &str) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"DELETE FROM provider_tokens WHERE provider = ?1",
params![provider],
)?;
Ok(())
}
pub fn list_provider_tokens(&self) -> StoreResult<Vec<super::ProviderToken>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT provider, access_token, refresh_token, expires_at, login, updated_at, credential_type, encrypted
FROM provider_tokens ORDER BY provider",
)?;
let rows = stmt.query_map([], read_token_row)?;
let mut tokens = Vec::new();
for row in rows {
tokens.push(decrypt_token_row(row?)?);
}
Ok(tokens)
}
pub fn upsert_secrets_provider_config(
&self,
config: &super::SecretsProviderConfig,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"INSERT INTO secrets_provider_config (provider, project, config, updated_at)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(provider) DO UPDATE SET
project = excluded.project,
config = excluded.config,
updated_at = excluded.updated_at",
params![
config.provider,
config.project,
config.config,
config.updated_at,
],
)?;
Ok(())
}
pub fn delete_secrets_provider_config(&self, provider: &str) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"DELETE FROM secrets_provider_config WHERE provider = ?1",
params![provider],
)?;
Ok(())
}
pub fn list_secrets_provider_configs(&self) -> StoreResult<Vec<super::SecretsProviderConfig>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT provider, project, config, updated_at
FROM secrets_provider_config ORDER BY provider",
)?;
let rows = stmt.query_map([], |row| Ok(read_secrets_config_row(row)))?;
let mut configs = Vec::new();
for row in rows {
configs.push(row??);
}
Ok(configs)
}
pub fn list_repos(&self) -> StoreResult<Vec<Repo>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt =
conn.prepare("SELECT path, repo_id, name, added_at FROM repos ORDER BY path ASC")?;
let rows = stmt.query_map([], |row| Ok(map_repo_row(row)))?;
let mut repos = Vec::new();
for row in rows {
repos.push(row??);
}
Ok(repos)
}
pub fn get_repo(&self, path: &str) -> StoreResult<Option<Repo>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn
.prepare("SELECT path, repo_id, name, added_at FROM repos WHERE path = ?1 LIMIT 1")?;
let row = stmt
.query_row(params![path], |row| Ok(map_repo_row(row)))
.optional()?;
row.transpose()
}
pub fn get_repo_by_repo_id(&self, repo_id: &RepoId) -> StoreResult<Option<Repo>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT path, repo_id, name, added_at FROM repos WHERE repo_id = ?1 LIMIT 1",
)?;
let row = stmt
.query_row(params![repo_id.as_str()], |row| Ok(map_repo_row(row)))
.optional()?;
row.transpose()
}
pub fn upsert_repo(&self, repo: &Repo) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"INSERT INTO repos (path, repo_id, name, added_at)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(path) DO UPDATE SET repo_id = excluded.repo_id, name = excluded.name, added_at = excluded.added_at",
params![
repo.path,
repo.repo_id.as_str(),
repo.name,
repo.added_at.unix_timestamp()
],
)?;
Ok(())
}
pub fn delete_repo(&self, path: &str) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute("DELETE FROM repos WHERE path = ?1", params![path])?;
Ok(())
}
pub fn list_edges(&self) -> StoreResult<Vec<RepoEdge>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT parent_repo_id, child_repo_id FROM repo_edges ORDER BY parent_repo_id, child_repo_id",
)?;
let rows = stmt.query_map([], |row| Ok(map_repo_edge_row(row)))?;
let mut edges = Vec::new();
for row in rows {
edges.push(row??);
}
Ok(edges)
}
pub fn add_edge(&self, edge: &RepoEdge) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"INSERT OR IGNORE INTO repo_edges (parent_repo_id, child_repo_id) VALUES (?1, ?2)",
params![edge.parent_repo_id.as_str(), edge.child_repo_id.as_str()],
)?;
Ok(())
}
pub fn remove_edge(&self, parent_id: &RepoId, child_id: &RepoId) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"DELETE FROM repo_edges WHERE parent_repo_id = ?1 AND child_repo_id = ?2",
params![parent_id.as_str(), child_id.as_str()],
)?;
Ok(())
}
pub fn children(&self, repo_id: &RepoId) -> StoreResult<Vec<Repo>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT repos.path, repos.repo_id, repos.name, repos.added_at
FROM repo_edges
INNER JOIN repos ON repos.repo_id = repo_edges.child_repo_id
WHERE repo_edges.parent_repo_id = ?1
ORDER BY repos.path ASC",
)?;
let rows = stmt.query_map(params![repo_id.as_str()], |row| Ok(map_repo_row(row)))?;
let mut repos = Vec::new();
for row in rows {
repos.push(row??);
}
Ok(repos)
}
pub fn parents(&self, repo_id: &RepoId) -> StoreResult<Vec<Repo>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT repos.path, repos.repo_id, repos.name, repos.added_at
FROM repo_edges
INNER JOIN repos ON repos.repo_id = repo_edges.parent_repo_id
WHERE repo_edges.child_repo_id = ?1
ORDER BY repos.path ASC",
)?;
let rows = stmt.query_map(params![repo_id.as_str()], |row| Ok(map_repo_row(row)))?;
let mut repos = Vec::new();
for row in rows {
repos.push(row??);
}
Ok(repos)
}
pub fn create_session(&self, session: &Session) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"INSERT INTO sessions (id, harness, status, wave_run_id, provider_session_id, config, created_at, ended_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
params![
session.id,
session.harness,
session.status.as_i32() as i64,
session.wave_run_id,
session.provider_session_id,
serde_json::to_string(&session.config)?,
session.created_at.unix_timestamp(),
session.ended_at.map(|dt| dt.unix_timestamp()),
],
)?;
Ok(())
}
pub fn get_session(&self, session_id: &LfdId) -> StoreResult<Option<Session>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT id, harness, status, wave_run_id, provider_session_id, config, created_at, ended_at
FROM sessions WHERE id = ?1",
)?;
let row = stmt
.query_row(params![session_id], Self::map_session_row)
.optional()?;
Ok(row)
}
pub fn get_active_session_for_wave_run(
&self,
wave_run_id: &str,
) -> StoreResult<Option<Session>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT id, harness, status, wave_run_id, provider_session_id, config, created_at, ended_at
FROM sessions
WHERE wave_run_id = ?1 AND status IN (?2, ?3, ?4)
ORDER BY created_at DESC
LIMIT 1",
)?;
let row = stmt
.query_row(
params![
wave_run_id,
SessionStatus::Starting.as_i32() as i64,
SessionStatus::Active.as_i32() as i64,
SessionStatus::Ending.as_i32() as i64
],
Self::map_session_row,
)
.optional()?;
Ok(row)
}
pub fn update_provider_session_id(
&self,
session_id: &LfdId,
provider_session_id: &str,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
"UPDATE sessions SET provider_session_id = ?2 WHERE id = ?1",
params![session_id, provider_session_id],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn update_session_status(
&self,
session_id: &LfdId,
status: SessionStatus,
ended_at: Option<i64>,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
"UPDATE sessions
SET status = ?2, ended_at = COALESCE(?3, ended_at)
WHERE id = ?1",
params![session_id, status.as_i32() as i64, ended_at],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn append_session_event(
&self,
session_id: &LfdId,
seq: i64,
event: &SessionEvent,
created_at: i64,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"INSERT INTO session_events (session_id, seq, event_type, data, created_at)
VALUES (?1, ?2, ?3, ?4, ?5)",
params![
session_id,
seq,
event.event_type(),
serde_json::to_string(event)?,
created_at,
],
)?;
Ok(())
}
pub fn list_session_events(
&self,
session_id: &LfdId,
after_seq: Option<i64>,
) -> StoreResult<Vec<PersistedSessionEvent>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = if after_seq.is_some() {
conn.prepare(
"SELECT session_id, seq, data, created_at
FROM session_events
WHERE session_id = ?1 AND seq > ?2
ORDER BY seq ASC",
)?
} else {
conn.prepare(
"SELECT session_id, seq, data, created_at
FROM session_events
WHERE session_id = ?1
ORDER BY seq ASC",
)?
};
let mut rows = if let Some(after_seq) = after_seq {
stmt.query(params![session_id, after_seq])?
} else {
stmt.query(params![session_id])?
};
let mut events = Vec::new();
while let Some(row) = rows.next()? {
let data: String = row.get(2)?;
let event: SessionEvent = serde_json::from_str(&data)?;
events.push(PersistedSessionEvent {
session_id: row.get(0)?,
seq: row.get(1)?,
event,
created_at: crate::lfd::store::rows::unix_to_datetime(row.get(3)?),
});
}
Ok(events)
}
pub fn list_sessions_by_statuses(
&self,
statuses: &[SessionStatus],
) -> StoreResult<Vec<Session>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let placeholders: Vec<String> = statuses
.iter()
.enumerate()
.map(|(i, _)| format!("?{}", i + 1))
.collect();
let sql = format!(
"SELECT id, harness, status, wave_run_id, provider_session_id, config, created_at, ended_at
FROM sessions WHERE status IN ({})
ORDER BY created_at ASC",
placeholders.join(", ")
);
let mut stmt = conn.prepare(&sql)?;
let params: Vec<Box<dyn ToSql>> = statuses
.iter()
.map(|s| Box::new(s.as_i32() as i64) as Box<dyn ToSql>)
.collect();
let param_refs: Vec<&dyn ToSql> = params.iter().map(|p| p.as_ref()).collect();
let mut rows = stmt.query(param_refs.as_slice())?;
let mut sessions = Vec::new();
while let Some(row) = rows.next()? {
sessions.push(Self::map_session_row(row)?);
}
Ok(sessions)
}
pub fn list_events_for_sessions(
&self,
session_ids: &[LfdId],
) -> StoreResult<HashMap<LfdId, Vec<PersistedSessionEvent>>> {
if session_ids.is_empty() {
return Ok(HashMap::new());
}
let conn = self.conn.lock().expect("store mutex poisoned");
let placeholders: Vec<String> = (1..=session_ids.len()).map(|i| format!("?{i}")).collect();
let sql = format!(
"SELECT session_id, seq, data, created_at
FROM session_events
WHERE session_id IN ({})
ORDER BY session_id, seq ASC",
placeholders.join(", ")
);
let mut stmt = conn.prepare(&sql)?;
let params: Vec<&dyn ToSql> = session_ids.iter().map(|id| id as &dyn ToSql).collect();
let mut rows = stmt.query(params.as_slice())?;
let mut result: HashMap<LfdId, Vec<PersistedSessionEvent>> = HashMap::new();
while let Some(row) = rows.next()? {
let session_id: LfdId = row.get(0)?;
let data: String = row.get(2)?;
let event: SessionEvent = serde_json::from_str(&data)?;
result
.entry(session_id.clone())
.or_default()
.push(PersistedSessionEvent {
session_id,
seq: row.get(1)?,
event,
created_at: crate::lfd::store::rows::unix_to_datetime(row.get(3)?),
});
}
Ok(result)
}
pub fn list_sessions_for_wave(&self, wave_id: &str) -> StoreResult<Vec<Session>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT s.id, s.harness, s.status, s.wave_run_id, s.provider_session_id, s.config, s.created_at, s.ended_at
FROM sessions s
JOIN wave_runs wr ON wr.id = s.wave_run_id
WHERE wr.wave_id = ?1
ORDER BY s.created_at ASC",
)?;
let mut rows = stmt.query(params![wave_id])?;
let mut sessions = Vec::new();
while let Some(row) = rows.next()? {
sessions.push(Self::map_session_row(row)?);
}
Ok(sessions)
}
pub fn list_sessions_filtered(&self, filters: &SessionFilters) -> StoreResult<Vec<Session>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut sql = String::from(
"SELECT s.id, s.harness, s.status, s.wave_run_id, s.provider_session_id, s.config, s.created_at, s.ended_at
FROM sessions s",
);
let mut predicates = Vec::new();
let mut params: Vec<Box<dyn ToSql>> = Vec::new();
if filters.wave.is_some() || filters.flow.is_some() {
sql.push_str(" JOIN wave_runs wr ON wr.id = s.wave_run_id");
}
if let Some(wave) = filters.wave.as_ref() {
predicates.push(format!("wr.wave_id = ?{}", params.len() + 1));
params.push(Box::new(wave.clone()));
}
if let Some(flow) = filters.flow.as_ref() {
predicates.push(format!("wr.snapshot_flow = ?{}", params.len() + 1));
params.push(Box::new(flow.clone()));
}
if let Some(step) = filters.step.as_ref() {
predicates.push(format!(
"json_extract(s.config, '$.step') = ?{}",
params.len() + 1
));
params.push(Box::new(step.clone()));
}
if let Some(from) = filters.from {
predicates.push(format!("s.created_at >= ?{}", params.len() + 1));
params.push(Box::new(from));
}
if let Some(to) = filters.to {
predicates.push(format!("s.created_at <= ?{}", params.len() + 1));
params.push(Box::new(to));
}
if !predicates.is_empty() {
sql.push_str(" WHERE ");
sql.push_str(&predicates.join(" AND "));
}
sql.push_str(" ORDER BY s.created_at ASC");
let mut stmt = conn.prepare(&sql)?;
let param_refs: Vec<&dyn ToSql> = params.iter().map(|value| value.as_ref()).collect();
let mut rows = stmt.query(param_refs.as_slice())?;
let mut sessions = Vec::new();
while let Some(row) = rows.next()? {
sessions.push(Self::map_session_row(row)?);
}
Ok(sessions)
}
const TERMINAL_SESSION_COLS: &str =
"id, wave_id, wave_run_id, step, agent, cwd, argv, env, source, tmux_name, status, \
completion_token, created_at, attached_at, started_at, completed_at";
pub fn create_terminal_session(&self, session: &TerminalSession) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
&format!(
"INSERT INTO terminal_sessions ({}) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)",
Self::TERMINAL_SESSION_COLS
),
params![
session.id,
session.wave_id,
session.wave_run_id,
session.step,
session.agent,
session.cwd,
serde_json::to_string(&session.argv)?,
serde_json::to_string(&session.env)?,
session.source,
session.tmux_name,
session.status.as_i32() as i64,
session.completion_token,
session.created_at.unix_timestamp(),
session.attached_at.map(|dt| dt.unix_timestamp()),
session.started_at.map(|dt| dt.unix_timestamp()),
session.completed_at.map(|dt| dt.unix_timestamp()),
],
)?;
Ok(())
}
pub fn get_terminal_session(&self, session_id: &LfdId) -> StoreResult<Option<TerminalSession>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(&format!(
"SELECT {} FROM terminal_sessions WHERE id = ?1",
Self::TERMINAL_SESSION_COLS
))?;
let row = stmt
.query_row(params![session_id], Self::map_terminal_session_row)
.optional()?;
Ok(row)
}
pub fn list_terminal_sessions(
&self,
wave_id: Option<&LfdId>,
statuses: Option<&[TerminalSessionStatus]>,
) -> StoreResult<Vec<TerminalSession>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut sql = format!(
"SELECT {} FROM terminal_sessions",
Self::TERMINAL_SESSION_COLS
);
let mut predicates = Vec::new();
let mut params: Vec<Box<dyn ToSql>> = Vec::new();
if let Some(wave_id) = wave_id {
predicates.push(format!("wave_id = ?{}", params.len() + 1));
params.push(Box::new(wave_id.clone()));
}
if let Some(statuses) = statuses {
let placeholders = statuses
.iter()
.enumerate()
.map(|(index, _)| format!("?{}", params.len() + index + 1))
.collect::<Vec<_>>();
predicates.push(format!("status IN ({})", placeholders.join(", ")));
params.extend(
statuses
.iter()
.map(|status| Box::new(status.as_i32() as i64) as Box<dyn ToSql>),
);
}
if !predicates.is_empty() {
sql.push_str(" WHERE ");
sql.push_str(&predicates.join(" AND "));
}
sql.push_str(" ORDER BY created_at ASC");
let mut stmt = conn.prepare(&sql)?;
let param_refs: Vec<&dyn ToSql> = params.iter().map(|value| value.as_ref()).collect();
let mut rows = stmt.query(param_refs.as_slice())?;
let mut sessions = Vec::new();
while let Some(row) = rows.next()? {
sessions.push(Self::map_terminal_session_row(row)?);
}
Ok(sessions)
}
pub fn get_active_terminal_session_for_wave_run(
&self,
wave_run_id: &LfdId,
) -> StoreResult<Option<TerminalSession>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(&format!(
"SELECT {} FROM terminal_sessions \
WHERE wave_run_id = ?1 AND status IN (?2, ?3, ?4) \
ORDER BY created_at DESC LIMIT 1",
Self::TERMINAL_SESSION_COLS
))?;
let row = stmt
.query_row(
params![
wave_run_id,
TerminalSessionStatus::Pending.as_i32() as i64,
TerminalSessionStatus::Attached.as_i32() as i64,
TerminalSessionStatus::Running.as_i32() as i64,
],
Self::map_terminal_session_row,
)
.optional()?;
Ok(row)
}
pub fn update_terminal_session(&self, session: &TerminalSession) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
"UPDATE terminal_sessions
SET wave_id = ?2,
wave_run_id = ?3,
step = ?4,
agent = ?5,
cwd = ?6,
argv = ?7,
env = ?8,
source = ?9,
tmux_name = ?10,
status = ?11,
completion_token = ?12,
created_at = ?13,
attached_at = ?14,
started_at = ?15,
completed_at = ?16
WHERE id = ?1",
params![
session.id,
session.wave_id,
session.wave_run_id,
session.step,
session.agent,
session.cwd,
serde_json::to_string(&session.argv)?,
serde_json::to_string(&session.env)?,
session.source,
session.tmux_name,
session.status.as_i32() as i64,
session.completion_token,
session.created_at.unix_timestamp(),
session.attached_at.map(|dt| dt.unix_timestamp()),
session.started_at.map(|dt| dt.unix_timestamp()),
session.completed_at.map(|dt| dt.unix_timestamp()),
],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn list_waves(&self, repo: Option<&str>) -> StoreResult<Vec<Wave>> {
self.read_waves(repo)
}
pub fn list_loopable_waves(&self) -> StoreResult<Vec<Wave>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::ListLoopableWaves))?;
let rows = stmt.query_map([], |row| Ok(map_wave_row(row)))?;
let mut waves = Vec::new();
for wave in rows {
waves.push(wave??);
}
Ok(waves)
}
pub fn list_wave_crons(&self, wave_id: &LfdId) -> StoreResult<Vec<WaveCron>> {
self.read_wave_crons(Query::ListWaveCrons, params![wave_id])
}
pub fn list_all_active_crons(&self) -> StoreResult<Vec<WaveCron>> {
self.read_wave_crons(Query::ListAllActiveCrons, [])
}
pub fn get_wave(&self, wave_id: &LfdId) -> StoreResult<Option<Wave>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetWaveById))?;
let wave = stmt
.query_row(params![wave_id], |row| Ok(map_wave_row(row)))
.optional()?;
wave.transpose()
}
pub fn get_wave_by_name(&self, name: &str) -> StoreResult<Option<Wave>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetWaveByName))?;
let wave = stmt
.query_row(params![name], |row| Ok(map_wave_row(row)))
.optional()?;
wave.transpose()
}
pub fn create_wave(&self, wave: &Wave) -> StoreResult<()> {
self.upsert_wave(wave)
}
pub fn update_wave(&self, wave: &Wave) -> StoreResult<()> {
self.upsert_wave(wave)
}
pub fn delete_wave(&self, wave_id: &LfdId) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"DELETE FROM attention_items WHERE wave_id = ?1",
params![wave_id],
)?;
conn.execute(Self::sql(Query::DeleteWaveById), params![wave_id])?;
Ok(())
}
pub fn create_wave_cron(&self, cron: &WaveCron) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
Self::sql(Query::InsertWaveCron),
params![
cron.id.as_str(),
cron.wave_id.as_str(),
cron.flow,
cron.schedule,
cron.last_triggered_at,
cron.created_at
.map(|value| value.unix_timestamp())
.unwrap_or_else(now_unix),
],
)?;
Ok(())
}
pub fn update_wave_cron_last_triggered(
&self,
cron_id: &LfdId,
last_triggered_at: Option<i64>,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
Self::sql(Query::UpdateWaveCronLastTriggered),
params![last_triggered_at, cron_id],
)?;
Ok(())
}
pub fn delete_wave_crons(&self, wave_id: &LfdId) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(Self::sql(Query::DeleteWaveCronsByWave), params![wave_id])?;
Ok(())
}
pub fn list_wave_runs(
&self,
wave_id: Option<&LfdId>,
limit: Option<u32>,
) -> StoreResult<Vec<WaveRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let query = Self::sql(list_wave_runs_query(wave_id.is_some(), limit.is_some()));
let mut params_vec: Vec<Box<dyn ToSql>> = Vec::new();
if let Some(wave_id) = wave_id {
params_vec.push(Box::new(wave_id.clone()));
}
if let Some(limit) = limit {
params_vec.push(Box::new(limit as i64));
}
let mut stmt = conn.prepare(query)?;
let params_iter = params_vec.iter().map(|v| v.as_ref() as &dyn ToSql);
let rows = stmt.query_map(rusqlite::params_from_iter(params_iter), |row| {
Ok(map_wave_run_row(row))
})?;
let mut runs = Vec::new();
for run in rows {
runs.push(run??);
}
Ok(runs)
}
pub fn get_wave_run(&self, wave_run_id: &LfdId) -> StoreResult<Option<WaveRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetWaveRunById))?;
let run = stmt
.query_row(params![wave_run_id], |row| Ok(map_wave_run_row(row)))
.optional()?;
run.transpose()
}
pub fn get_active_wave_run(&self, wave_id: &LfdId) -> StoreResult<Option<WaveRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetActiveWaveRun))?;
let run = stmt
.query_row(
params![
wave_id,
WaveRunStatus::Pending.as_i32() as i64,
WaveRunStatus::Running.as_i32() as i64,
WaveRunStatus::Waiting.as_i32() as i64,
],
|row| Ok(map_wave_run_row(row)),
)
.optional()?;
run.transpose()
}
pub fn count_active_wave_runs(&self, wave_id: &LfdId) -> StoreResult<u32> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::CountActiveWaveRuns))?;
let count = stmt.query_row(
params![
wave_id,
WaveRunStatus::Pending.as_i32() as i64,
WaveRunStatus::Running.as_i32() as i64,
WaveRunStatus::Waiting.as_i32() as i64,
],
|row| row.get::<_, i64>(0),
)?;
Ok(count.max(0) as u32)
}
pub fn get_latest_wave_run(&self, wave_id: &LfdId) -> StoreResult<Option<WaveRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetLatestWaveRun))?;
let run = stmt
.query_row(params![wave_id], |row| Ok(map_wave_run_row(row)))
.optional()?;
run.transpose()
}
pub fn create_wave_run(&self, run: &WaveRun) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let started_at = run
.started_at
.map(|dt| dt.unix_timestamp())
.unwrap_or_else(now_unix);
let flow_parents_json = serde_json::to_string(&run.flow_parents)?;
let execution_cursor = run.execution_cursor.clone();
conn.execute(
Self::sql(Query::InsertWaveRun),
params![
run.id,
run.wave_id,
run.iteration as i64,
run.step_index as i64,
run.status.as_i32() as i64,
run.worktree,
run.branch,
started_at,
run.ended_at.map(|dt| dt.unix_timestamp()),
run.error,
run.snapshot.repo,
run.snapshot.flow,
serde_json::to_string(&run.snapshot.direction)?,
serde_json::to_string(&run.snapshot.area)?,
serialize_pr(&run.pr)?,
flow_parents_json,
execution_cursor,
run.activation_log_id.as_ref(),
run.parent_run_id.as_ref(),
run.parent_pr_number.map(|value| value as i64),
run.stack_position as i64,
run.stack_group_id,
run.stack_status.as_i32() as i64,
if run.lineage_inferred { 1i64 } else { 0i64 },
run.target_branch,
run.repair_of.as_ref(),
],
)?;
Ok(())
}
pub fn update_wave_run(&self, run: &WaveRun) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let flow_parents_json = serde_json::to_string(&run.flow_parents)?;
let execution_cursor = run.execution_cursor.clone();
let updated = conn.execute(
Self::sql(Query::UpdateWaveRun),
params![
run.iteration as i64,
run.step_index as i64,
run.status.as_i32() as i64,
run.worktree,
run.branch,
run.started_at.map(|dt| dt.unix_timestamp()),
run.ended_at.map(|dt| dt.unix_timestamp()),
run.error,
run.snapshot.repo,
run.snapshot.flow,
serde_json::to_string(&run.snapshot.direction)?,
serde_json::to_string(&run.snapshot.area)?,
serialize_pr(&run.pr)?,
flow_parents_json,
execution_cursor,
run.activation_log_id.as_ref(),
run.parent_run_id.as_ref(),
run.parent_pr_number.map(|value| value as i64),
run.stack_position as i64,
run.stack_group_id,
run.stack_status.as_i32() as i64,
if run.lineage_inferred { 1i64 } else { 0i64 },
run.target_branch,
run.repair_of.as_ref(),
run.id,
],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn list_stack_runs(&self, wave_id: &LfdId) -> StoreResult<Vec<WaveRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::ListStackRuns))?;
let rows = stmt.query_map(params![wave_id], |row| Ok(map_wave_run_row(row)))?;
let mut runs = Vec::new();
for run in rows {
runs.push(run??);
}
Ok(runs)
}
pub fn get_live_pr_state(
&self,
repo_id: &str,
pr_number: u32,
) -> StoreResult<Option<LivePullRequestState>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetLivePrState))?;
let state = stmt
.query_row(params![repo_id, pr_number as i64], |row| {
Ok(map_live_pr_state_row(row))
})
.optional()?;
state.transpose()
}
pub fn upsert_live_pr_state(&self, state: &LivePullRequestState) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
Self::sql(Query::UpsertLivePrState),
params![
state.repo_id,
state.pr_number as i64,
state.state.as_i32() as i64,
if state.is_draft { 1i64 } else { 0i64 },
state.head_ref,
state.head_sha,
state.base_ref,
state.updated_at.unix_timestamp(),
state.merged_at.map(|value| value.unix_timestamp()),
state.synced_at.unix_timestamp(),
],
)?;
Ok(())
}
fn map_attention_item_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<AttentionItem> {
let kind_raw: String = row.get(3)?;
let status_raw: String = row.get(4)?;
let context_raw: String = row.get(7)?;
let kind = kind_raw.parse::<AttentionKind>().map_err(|err| {
rusqlite::Error::FromSqlConversionFailure(
3,
rusqlite::types::Type::Text,
Box::new(StoreError::InvalidData(err)),
)
})?;
let status = status_raw.parse::<AttentionStatus>().map_err(|err| {
rusqlite::Error::FromSqlConversionFailure(
4,
rusqlite::types::Type::Text,
Box::new(StoreError::InvalidData(err)),
)
})?;
Ok(AttentionItem {
id: LfdId::from_raw(row.get::<_, String>(0)?),
wave_id: LfdId::from_raw(row.get::<_, String>(1)?),
run_id: row.get::<_, Option<String>>(2)?.map(LfdId::from_raw),
kind,
status,
title: row.get(5)?,
summary: row.get(6)?,
context: serde_json::from_str(&context_raw)
.unwrap_or(serde_json::Value::Object(Default::default())),
surfaced_at: crate::lfd::store::rows::unix_to_datetime(row.get(8)?),
viewed_at: row
.get::<_, Option<i64>>(9)?
.map(crate::lfd::store::rows::unix_to_datetime),
resolved_at: row
.get::<_, Option<i64>>(10)?
.map(crate::lfd::store::rows::unix_to_datetime),
})
}
pub fn list_attention_items(
&self,
status: Option<AttentionStatus>,
kind: Option<AttentionKind>,
) -> StoreResult<Vec<AttentionItem>> {
let mut sql = String::from(
"SELECT id, wave_id, run_id, kind, status, title, summary, context, surfaced_at, viewed_at, resolved_at\n FROM attention_items",
);
let mut params: Vec<String> = Vec::new();
let mut clauses: Vec<String> = Vec::new();
if let Some(status) = status {
clauses.push("status = ?".to_string());
params.push(status.as_str().to_string());
}
if let Some(kind) = kind {
clauses.push("kind = ?".to_string());
params.push(kind.as_str().to_string());
}
if !clauses.is_empty() {
sql.push_str(" WHERE ");
sql.push_str(&clauses.join(" AND "));
}
sql.push_str(" ORDER BY surfaced_at DESC");
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(&sql)?;
let rows = stmt.query_map(
rusqlite::params_from_iter(params.iter()),
Self::map_attention_item_row,
)?;
rows.collect::<rusqlite::Result<Vec<_>>>()
.map_err(StoreError::from)
}
pub fn get_attention_item(&self, attention_id: &LfdId) -> StoreResult<Option<AttentionItem>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT id, wave_id, run_id, kind, status, title, summary, context, surfaced_at, viewed_at, resolved_at\n FROM attention_items WHERE id = ?1",
)?;
stmt.query_row(
rusqlite::params![attention_id],
Self::map_attention_item_row,
)
.optional()
.map_err(StoreError::from)
}
pub fn find_attention_item_for_run(
&self,
run_id: &LfdId,
kind: AttentionKind,
) -> StoreResult<Option<AttentionItem>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT id, wave_id, run_id, kind, status, title, summary, context, surfaced_at, viewed_at, resolved_at
FROM attention_items
WHERE run_id = ?1 AND kind = ?2 AND status != ?3
ORDER BY surfaced_at DESC
LIMIT 1",
)?;
stmt.query_row(
rusqlite::params![run_id, kind.as_str(), AttentionStatus::Resolved.as_str()],
Self::map_attention_item_row,
)
.optional()
.map_err(StoreError::from)
}
pub fn upsert_attention_item(&self, item: &AttentionItem) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"INSERT INTO attention_items (id, wave_id, run_id, kind, status, title, summary, context, surfaced_at, viewed_at, resolved_at)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)
ON CONFLICT(id) DO UPDATE SET
wave_id = excluded.wave_id,
run_id = excluded.run_id,
kind = excluded.kind,
status = excluded.status,
title = excluded.title,
summary = excluded.summary,
context = excluded.context,
surfaced_at = excluded.surfaced_at,
viewed_at = excluded.viewed_at,
resolved_at = excluded.resolved_at",
rusqlite::params![
&item.id,
&item.wave_id,
&item.run_id,
&item.kind.as_str(),
&item.status.as_str(),
&item.title,
&item.summary,
&serde_json::to_string(&item.context)?,
&item.surfaced_at.unix_timestamp(),
&item.viewed_at.map(|value: time::OffsetDateTime| value.unix_timestamp()),
&item.resolved_at.map(|value: time::OffsetDateTime| value.unix_timestamp()),
],
)?;
Ok(())
}
pub fn delete_attention_item(&self, attention_id: &LfdId) -> StoreResult<u32> {
let conn = self.conn.lock().expect("store mutex poisoned");
let deleted = conn.execute(
"DELETE FROM attention_items WHERE id = ?1",
rusqlite::params![attention_id],
)?;
Ok(deleted as u32)
}
pub fn list_queue_blocks(&self, wave_id: &LfdId) -> StoreResult<Vec<QueueBlock>> {
let blocks = self
.list_attention_items(
Some(AttentionStatus::Surfaced),
Some(AttentionKind::Algedonic),
)?
.into_iter()
.filter(|item| &item.wave_id == wave_id)
.filter_map(|item| queue_block_from_attention(&item).ok().flatten())
.collect();
Ok(blocks)
}
pub fn upsert_queue_block(&self, block: &QueueBlock) -> StoreResult<()> {
self.upsert_attention_item(&queue_block_attention_item(block))
}
pub fn delete_queue_block(&self, _wave_id: &LfdId, run_id: &LfdId) -> StoreResult<u32> {
let id = crate::lfd::attention::attention_id_for_queue_block(run_id);
let Some(mut item) = self.get_attention_item(&id)? else {
return Ok(0);
};
item.status = AttentionStatus::Resolved;
item.resolved_at = Some(time::OffsetDateTime::now_utc());
self.upsert_attention_item(&item)?;
Ok(1)
}
pub fn record_merge_event(&self, event: &QueueMergeEvent) -> StoreResult<bool> {
let conn = self.conn.lock().expect("store mutex poisoned");
let inserted = conn.execute(
"INSERT INTO wave_pr_merge_events (wave_id, pr_number, merged_at, processed_at)
VALUES (?1, ?2, ?3, ?4)
ON CONFLICT(wave_id, pr_number, merged_at) DO NOTHING",
params![
event.wave_id,
event.pr_number as i64,
event.merged_at.unix_timestamp(),
event.processed_at.unix_timestamp(),
],
)?;
Ok(inserted > 0)
}
pub fn fail_orphaned_runs(&self) -> StoreResult<u32> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
Self::sql(Query::FailOrphanedRuns),
params![
WaveRunStatus::Failed.as_i32() as i64,
"orphaned: lfd restarted",
now_unix(),
WaveRunStatus::Pending.as_i32() as i64,
WaveRunStatus::Running.as_i32() as i64,
WaveRunStatus::Waiting.as_i32() as i64,
],
)?;
conn.execute(
Self::sql(Query::ResetStaleActiveWaves),
params![
WaveStatus::Idle.as_i32() as i64,
WaveStatus::Running.as_i32() as i64,
WaveStatus::Waiting.as_i32() as i64,
],
)?;
Ok(updated as u32)
}
pub fn list_triggers(&self, wave_id: Option<&LfdId>) -> StoreResult<Vec<Trigger>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let query = Self::sql(list_triggers_query(wave_id.is_some()));
let params: Vec<Box<dyn ToSql>> = if let Some(wave_id) = wave_id {
vec![Box::new(wave_id.clone())]
} else {
vec![]
};
let mut stmt = conn.prepare(query)?;
let params_iter = params.iter().map(|v| v.as_ref() as &dyn ToSql);
let rows = stmt.query_map(rusqlite::params_from_iter(params_iter), |row| {
Ok(map_trigger_row(row))
})?;
let mut triggers = Vec::new();
for trigger in rows {
triggers.push(trigger??);
}
Ok(triggers)
}
pub fn list_triggers_by_signal(&self, signal: i32) -> StoreResult<Vec<Trigger>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::ListTriggersBySignal))?;
let rows = stmt.query_map(params![signal as i64], |row| Ok(map_trigger_row(row)))?;
let mut triggers = Vec::new();
for trigger in rows {
triggers.push(trigger??);
}
Ok(triggers)
}
pub fn get_trigger(&self, trigger_id: &LfdId) -> StoreResult<Option<Trigger>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetTriggerById))?;
let trigger = stmt
.query_row(params![trigger_id], |row| Ok(map_trigger_row(row)))
.optional()?;
trigger.transpose()
}
pub fn create_trigger(&self, trigger: &Trigger) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let created_at = trigger
.created_at
.map(|dt| dt.unix_timestamp())
.unwrap_or_else(now_unix);
conn.execute(
Self::sql(Query::InsertTrigger),
params![
trigger.id,
trigger.wave_id,
trigger.signal.as_i32() as i64,
trigger.flow,
trigger.last_main_sha,
trigger.last_triggered_at,
created_at,
trigger.enabled as i64,
trigger.source_wave_id,
trigger.max_iterations.map(|v| v as i64),
],
)?;
Ok(())
}
pub fn update_trigger(&self, trigger: &Trigger) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
Self::sql(Query::UpdateTrigger),
params![
trigger.signal.as_i32() as i64,
trigger.flow,
trigger.last_main_sha,
trigger.last_triggered_at,
trigger.enabled as i64,
trigger.source_wave_id,
trigger.max_iterations.map(|v| v as i64),
trigger.id,
],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn delete_trigger(&self, trigger_id: &LfdId) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(Self::sql(Query::DeleteTriggerById), params![trigger_id])?;
Ok(())
}
pub fn list_pending_activations(&self, wave_id: &LfdId) -> StoreResult<Vec<PendingActivation>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::ListPendingActivationsByWave))?;
let rows = stmt.query_map(params![wave_id], |row| Ok(map_pending_activation_row(row)))?;
let mut activations = Vec::new();
for activation in rows {
activations.push(activation??);
}
Ok(activations)
}
pub fn create_pending_activation(&self, activation: &PendingActivation) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
Self::sql(Query::InsertPendingActivation),
params![
activation.id,
activation.wave_id,
activation.trigger_id,
activation.reason,
activation.from_sha,
activation.to_sha,
activation.queued_at,
activation.target_branch,
],
)?;
Ok(())
}
pub fn update_pending_activation(&self, activation: &PendingActivation) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
Self::sql(Query::UpdatePendingActivation),
params![
activation.reason,
activation.from_sha,
activation.to_sha,
activation.target_branch,
activation.id
],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn get_pending_for_trigger(
&self,
wave_id: &LfdId,
trigger_id: Option<&LfdId>,
) -> StoreResult<Option<PendingActivation>> {
let conn = self.conn.lock().expect("store mutex poisoned");
match trigger_id {
Some(tid) => {
let mut stmt = conn.prepare(Self::sql(Query::GetPendingActivationForTrigger))?;
let activation = stmt
.query_row(params![wave_id, tid], |row| {
Ok(map_pending_activation_row(row))
})
.optional()?;
activation.transpose()
}
None => {
let mut stmt = conn.prepare(Self::sql(Query::GetPendingActivationForWave))?;
let activation = stmt
.query_row(params![wave_id], |row| Ok(map_pending_activation_row(row)))
.optional()?;
activation.transpose()
}
}
}
pub fn delete_pending_activation_by_id(&self, activation_id: &LfdId) -> StoreResult<u32> {
let conn = self.conn.lock().expect("store mutex poisoned");
let deleted = conn.execute(
Self::sql(Query::DeletePendingActivationById),
params![activation_id],
)?;
Ok(deleted as u32)
}
pub fn create_activation_log(&self, log: &ActivationLog) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
Self::sql(Query::InsertActivationLog),
params![
log.id,
log.wave_id,
log.trigger_id,
log.reason,
log.outcome.as_str(),
log.created_at,
],
)?;
Ok(())
}
pub fn list_activation_log(
&self,
wave_id: &LfdId,
limit: u32,
) -> StoreResult<Vec<ActivationLog>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::ListActivationLogByWave))?;
let rows = stmt.query_map(params![wave_id, limit as i64], |row| {
Ok(map_activation_log_row(row))
})?;
let mut entries = Vec::new();
for row in rows {
entries.push(row??);
}
Ok(entries)
}
pub fn get_activation_log(
&self,
activation_log_id: &LfdId,
) -> StoreResult<Option<ActivationLog>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetActivationLogById))?;
let entry = stmt
.query_row(params![activation_log_id], |row| {
Ok(map_activation_log_row(row))
})
.optional()?;
entry.transpose()
}
pub fn list_fork_runs(
&self,
wave_run_id: &LfdId,
step_index: u32,
) -> StoreResult<Vec<ForkRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::ListForkRuns))?;
let rows = stmt.query_map(params![wave_run_id, step_index as i64], |row| {
Ok(map_fork_run_row(row))
})?;
let mut runs = Vec::new();
for run in rows {
runs.push(run??);
}
Ok(runs)
}
pub fn list_orphaned_fork_runs(&self) -> StoreResult<Vec<ForkRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT fr.id, fr.wave_run_id, fr.step_index, fr.branch_index, fr.status, fr.worktree
FROM fork_runs fr
LEFT JOIN wave_runs wr ON wr.id = fr.wave_run_id
WHERE fr.status IN (?1, ?2)
AND (
wr.id IS NULL
OR wr.status NOT IN (?3, ?4, ?5)
OR fr.step_index != wr.step_index
)
ORDER BY fr.wave_run_id ASC, fr.step_index ASC, fr.branch_index ASC",
)?;
let rows = stmt.query_map(
params![
ForkRunStatus::Pending as i32 as i64,
ForkRunStatus::Running as i32 as i64,
WaveRunStatus::Pending.as_i32() as i64,
WaveRunStatus::Running.as_i32() as i64,
WaveRunStatus::Waiting.as_i32() as i64
],
|row| Ok(map_fork_run_row(row)),
)?;
let mut runs = Vec::new();
for run in rows {
runs.push(run??);
}
Ok(runs)
}
pub fn upsert_fork_run(&self, fork_run: &ForkRun) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
Self::sql(Query::UpsertForkRun),
params![
fork_run.id,
fork_run.wave_run_id,
fork_run.step_index as i64,
fork_run.branch_index as i64,
fork_run.status as i32 as i64,
fork_run.worktree,
],
)?;
Ok(())
}
pub fn delete_fork_runs(&self, wave_run_id: &LfdId, step_index: u32) -> StoreResult<u32> {
let conn = self.conn.lock().expect("store mutex poisoned");
let deleted = conn.execute(
Self::sql(Query::DeleteForkRuns),
params![wave_run_id, step_index as i64],
)?;
Ok(deleted as u32)
}
pub fn list_agents(&self) -> StoreResult<Vec<AgentRun>> {
self.list_agent_history(None, None, None)
}
pub fn list_agent_history(
&self,
worktree: Option<&str>,
repo: Option<&str>,
limit: Option<u32>,
) -> StoreResult<Vec<AgentRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let query = Self::sql(list_agent_history_query(
worktree.is_some(),
repo.is_some(),
limit.is_some(),
));
let mut params_vec: Vec<Box<dyn ToSql>> = Vec::new();
if let Some(worktree) = worktree {
params_vec.push(Box::new(worktree.to_string()));
}
if let Some(repo) = repo {
params_vec.push(Box::new(repo.to_string()));
}
if let Some(limit) = limit {
params_vec.push(Box::new(limit as i64));
}
let mut stmt = conn.prepare(query)?;
let params_iter = params_vec.iter().map(|v| v.as_ref() as &dyn ToSql);
let rows = stmt.query_map(rusqlite::params_from_iter(params_iter), |row| {
Ok(map_agent_row(row))
})?;
let mut runs = Vec::new();
for run in rows {
runs.push(run??);
}
Ok(runs)
}
pub fn get_agent(&self, agent_id: &LfdId) -> StoreResult<Option<AgentRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetAgentById))?;
let run = stmt
.query_row(params![agent_id], |row| Ok(map_agent_row(row)))
.optional()?;
run.transpose()
}
pub fn get_waiting_agent_for_wave(&self, wave_id: &LfdId) -> StoreResult<Option<AgentRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetWaitingAgentForWave))?;
let run = stmt
.query_row(
params![wave_id, AgentStatus::Waiting.as_i32() as i64],
|row| Ok(map_agent_row(row)),
)
.optional()?;
run.transpose()
}
pub fn start_agent(&self, agent_run: &AgentRun) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let started_at = agent_run
.started_at
.map(|dt| dt.unix_timestamp())
.unwrap_or_else(now_unix);
conn.execute(
Self::sql(Query::InsertAgent),
params![
agent_run.id,
agent_run.step,
agent_run.repo,
agent_run.worktree,
agent_run.wave_run_id,
agent_run.status.as_i32() as i64,
started_at,
agent_run.ended_at.map(|dt| dt.unix_timestamp()),
agent_run.pid.map(|v| v as i64),
agent_run.container_id.as_deref(),
agent_run.agent,
agent_run.run_mode,
],
)?;
Ok(())
}
pub fn update_agent_status(
&self,
agent_id: &LfdId,
status: i32,
pid: Option<u32>,
container_id: Option<&str>,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
Self::sql(Query::UpdateAgentStatus),
params![status as i64, pid.map(|v| v as i64), container_id, agent_id],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn end_agent(&self, agent_id: &LfdId, status: i32, ended_at: i64) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
Self::sql(Query::EndAgent),
params![status as i64, ended_at, agent_id],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn get_active_agents_for_wave(&self, wave_id: &LfdId) -> StoreResult<Vec<AgentRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetActiveAgentsForWave))?;
let rows = stmt.query_map(params![wave_id], |row| Ok(map_agent_row(row)))?;
let mut agents = Vec::new();
for row in rows {
agents.push(row??);
}
Ok(agents)
}
pub fn end_active_agent_for_wave(
&self,
wave_id: &LfdId,
status: i32,
ended_at: i64,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated = conn.execute(
Self::sql(Query::EndActiveAgentsForWave),
params![status as i64, ended_at, wave_id.as_str()],
)?;
if updated == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn get_stuck_agents(&self, older_than_secs: u64) -> StoreResult<Vec<AgentRun>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let cutoff = now_unix() - older_than_secs as i64;
let mut stmt = conn.prepare(Self::sql(Query::GetStuckAgents))?;
let rows = stmt.query_map(params![cutoff], |row| Ok(map_agent_row(row)))?;
let mut runs = Vec::new();
for run in rows {
runs.push(run??);
}
Ok(runs)
}
pub fn get_summary(&self, wave_id: &LfdId) -> StoreResult<Option<Summary>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::GetSummaryByWave))?;
let summary = stmt
.query_row(params![wave_id], |row| Ok(map_summary_row(row)))
.optional()?;
summary.transpose()
}
pub fn upsert_summary(&self, summary: &Summary) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let created_at = summary
.created_at
.map(|dt| dt.unix_timestamp())
.unwrap_or_else(now_unix);
conn.execute(
Self::sql(Query::UpsertSummary),
params![
summary.id,
summary.wave_id,
summary.content,
summary.source_hash,
summary.token_budget as i64,
summary.agent,
created_at,
],
)?;
Ok(())
}
pub fn list_chat_memory_blocks(&self, wave_id: &LfdId) -> StoreResult<Vec<ChatMemoryBlock>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(Self::sql(Query::ListChatMemoryBlocks))?;
let rows = stmt.query_map(params![wave_id], |row| Ok(map_chat_memory_block_row(row)))?;
let mut blocks = Vec::new();
for row in rows {
blocks.push(row??);
}
Ok(blocks)
}
pub fn upsert_chat_memory_block(&self, block: &ChatMemoryBlock) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let updated_at = block
.updated_at
.map(|dt| dt.unix_timestamp())
.unwrap_or_else(now_unix);
conn.execute(
Self::sql(Query::UpsertChatMemoryBlock),
params![
block.wave_id,
block.name,
block.content,
block.position as i64,
updated_at,
],
)?;
Ok(())
}
pub fn delete_chat_memory_block(&self, wave_id: &LfdId, name: &str) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
Self::sql(Query::DeleteChatMemoryBlock),
params![wave_id, name],
)?;
Ok(())
}
pub fn list_chat_messages(&self, wave_id: &LfdId) -> StoreResult<Vec<ChatMessage>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut stmt = conn.prepare(
"SELECT id, wave_id, role, content, created_at
FROM chat_messages
WHERE wave_id = ?1
ORDER BY created_at ASC",
)?;
let rows = stmt.query_map(params![wave_id], |row| Ok(map_chat_message_row(row)))?;
let mut messages = Vec::new();
for row in rows {
messages.push(row??);
}
Ok(messages)
}
pub fn create_chat_message(&self, message: &ChatMessage) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"INSERT INTO chat_messages (id, wave_id, role, content, created_at)
VALUES (?1, ?2, ?3, ?4, ?5)",
params![
message.id,
message.wave_id,
message.role,
message.content,
message.created_at.unix_timestamp(),
],
)?;
Ok(())
}
}