use std::path::PathBuf;
use rusqlite::{params, Connection, OptionalExtension, ToSql, TransactionBehavior};
use time::OffsetDateTime;
use crate::child::{AbandonIntent, ChildRef, ObservationRecipient};
use crate::durable::{Author, FlowPosition, TaskFlowBlocker, TaskWorkerClaim};
use crate::id::WaveId;
use crate::planning::{LinearIssueId, LinearProjectId, ProjectPlan, TaskPlan};
use crate::store::rows::now_unix;
use crate::store::{StoreError, StoreResult};
use crate::work::project::{
ChildEventPayload, ObservationOutboxRow, Project, ProjectEvent, ProjectEventKind, ProjectId,
};
use crate::work::task::{
CiObservation, GithubObservation, GithubPr, LinearObservationApply, LinearObservationOutcome,
PmWritebackState, PrMergeRequest, PrPhase, PrPresentation, PrPublication, Task, TaskEvent,
TaskEventKind, TaskId, TaskLinearObservation, TaskPr, TaskPrId, TaskPrRepairKind,
};
use super::durable::{create_project_work, create_task_work};
use super::SqliteStore;
impl SqliteStore {
pub fn insert_task(&self, task: &Task, pr: &TaskPr) -> StoreResult<()> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
insert_initial_task(&transaction, task, pr)?;
transaction.commit()?;
Ok(())
}
pub fn insert_task_with_worktree(&self, task: &Task, pr: &TaskPr) -> StoreResult<()> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
insert_initial_task(&transaction, task, pr)?;
insert_task_event_in(
&transaction,
task,
&TaskEventKind::WorktreeInitializing {
pr_id: pr.id.clone(),
sequence: pr.sequence,
branch: pr.branch.clone(),
path: task.worktree.display().to_string(),
base_commit: pr.base_commit.clone(),
},
)?;
transaction.commit()?;
Ok(())
}
pub fn update_task_plan(&self, task_id: &TaskId, plan: &TaskPlan) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let changed = conn.execute(
"UPDATE tasks SET issue_identifier=?2, issue_title=?3, issue_description=?4,
pm_snapshot_synced_at=?5 WHERE id=?1 AND external_issue_id=?6",
params![
task_id.as_str(),
plan.identifier,
plan.title,
plan.description,
plan.pm_snapshot_synced_at,
plan.id.as_str()
],
)?;
if changed == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn update_task_pm_writeback(
&self,
task_id: &TaskId,
state: &PmWritebackState,
updated_at: OffsetDateTime,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
update_task_pm_writeback_in(&conn, task_id, state, updated_at)
}
pub fn update_task(&self, task: &Task) -> StoreResult<()> {
validate_task(task)?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
validate_task_project(&transaction, task)?;
let parameters = task_params(task);
let changed = transaction.execute(
TASK_UPDATE,
rusqlite::params_from_iter(parameters.iter().map(|value| value.as_ref())),
)?;
if changed == 0 {
return Err(StoreError::NotFound);
}
transaction.commit()?;
Ok(())
}
pub fn set_task_agent(&self, task_id: &TaskId, agent: &str) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let changed = conn.execute(
"UPDATE tasks SET agent=?2, updated_at=?3 WHERE id=?1",
params![task_id.as_str(), agent, now_unix()],
)?;
if changed == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn settle_task_worker(
&self,
task: &Task,
expected: &TaskWorkerClaim,
next: &FlowPosition,
progress: Option<&str>,
) -> StoreResult<FlowPosition> {
validate_task(task)?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
validate_task_project(&transaction, task)?;
update_task_timestamp_in(&transaction, task)?;
let position =
super::durable::settle_task_worker_in(&transaction, &task.id, expected, next)?;
if let Some(summary) = progress.filter(|summary| !summary.is_empty()) {
insert_task_event_in(
&transaction,
task,
&TaskEventKind::Progress {
summary: summary.to_string(),
},
)?;
}
transaction.commit()?;
Ok(position)
}
pub fn finish_task_flow(
&self,
task: &Task,
expected: &TaskWorkerClaim,
progress: Option<&str>,
) -> StoreResult<()> {
validate_task(task)?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
validate_task_project(&transaction, task)?;
update_task_timestamp_in(&transaction, task)?;
let position = super::durable::flow_position_in(&transaction, &task.id)?
.ok_or(StoreError::NotFound)?;
super::durable::finish_task_flow_in(&transaction, &task.id, expected)?;
insert_task_event_in(
&transaction,
task,
&TaskEventKind::FlowFinished {
invocation_id: position.invocation.id,
flow: position.invocation.flow,
summary: progress.unwrap_or_default().to_string(),
},
)?;
transaction.commit()?;
Ok(())
}
pub fn block_task_flow(
&self,
task_id: &TaskId,
expected: &TaskWorkerClaim,
failure: &TaskFlowBlocker,
) -> StoreResult<FlowPosition> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let current = matching_claim_in(&transaction, task_id, expected)?;
let position =
super::durable::block_task_flow_in(&transaction, task_id, ¤t, failure)?;
let task = task_on(&transaction, task_id)?.ok_or(StoreError::NotFound)?;
insert_task_event_in(
&transaction,
&task,
&TaskEventKind::Failed {
error: failure.reason.clone(),
resumable: !failure.restart_required,
},
)?;
transaction.commit()?;
Ok(position)
}
pub fn release_task_worker(
&self,
task_id: &TaskId,
expected: &TaskWorkerClaim,
) -> StoreResult<FlowPosition> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let current = matching_claim_in(&transaction, task_id, expected)?;
let position = super::durable::release_task_worker_in(&transaction, task_id, ¤t)?;
transaction.commit()?;
Ok(position)
}
pub fn complete_human_task_boundary(
&self,
task: &Task,
expected: &FlowPosition,
next: &FlowPosition,
summary: &str,
) -> StoreResult<FlowPosition> {
validate_task(task)?;
if expected.task_id != task.id
|| !expected.is_human()
|| expected.claim.is_some()
|| expected.failure.is_some()
{
return Err(StoreError::InvalidAuthority(
"human Task settlement requires its exact unclaimed position".to_string(),
));
}
if next.version != expected.version {
return Err(StoreError::InvalidAuthority(
"human Task settlement has a stale position version".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let current = super::durable::flow_position_in(&transaction, &task.id)?
.ok_or(StoreError::NotFound)?;
if current != *expected {
return Err(StoreError::InvalidAuthority(
"human Task position changed before settlement".to_string(),
));
}
validate_task_project(&transaction, task)?;
update_task_timestamp_in(&transaction, task)?;
let position = super::durable::set_flow_position_in(&transaction, &task.id, next)?;
insert_task_event_in(
&transaction,
task,
&TaskEventKind::Progress {
summary: summary.to_string(),
},
)?;
transaction.commit()?;
Ok(position)
}
pub fn finish_human_task_boundary(
&self,
task: &Task,
expected: &FlowPosition,
summary: &str,
) -> StoreResult<()> {
validate_task(task)?;
if expected.task_id != task.id
|| !expected.is_human()
|| expected.claim.is_some()
|| expected.failure.is_some()
{
return Err(StoreError::InvalidAuthority(
"Task review completion requires its exact unclaimed position".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let current = super::durable::flow_position_in(&transaction, &task.id)?
.ok_or(StoreError::NotFound)?;
if current != *expected {
return Err(StoreError::InvalidAuthority(
"Task review position changed before completion".to_string(),
));
}
validate_task_project(&transaction, task)?;
if transaction.execute(
"DELETE FROM task_flow_positions WHERE task_id=?1 AND position_version=?2",
params![
task.id.as_str(),
i64::try_from(expected.version)
.map_err(|error| StoreError::InvalidData(error.to_string()))?
],
)? != 1
{
return Err(StoreError::InvalidAuthority(
"Task review position changed before completion".to_string(),
));
}
insert_task_event_in(
&transaction,
task,
&TaskEventKind::FlowFinished {
invocation_id: expected.invocation.id.clone(),
flow: expected.invocation.flow.clone(),
summary: summary.to_string(),
},
)?;
transaction.commit()?;
Ok(())
}
pub fn retry_task_flow(
&self,
task_id: &TaskId,
expected: &FlowPosition,
feedback: Option<&str>,
) -> StoreResult<FlowPosition> {
if expected.task_id != *task_id || expected.claim.is_some() || expected.failure.is_none() {
return Err(StoreError::InvalidAuthority(
"Task retry requires its exact failed Flow position".to_string(),
));
}
if expected
.failure
.as_ref()
.is_some_and(|failure| failure.restart_required)
{
return Err(StoreError::InvalidAuthority(
"Task failure requires an explicit Flow restart".to_string(),
));
}
let mut next = expected.clone();
next.failure = None;
if let Some(feedback) = feedback {
next.cursor.leaf_mut().progress.direction = Some(feedback.to_string());
}
next.cursor.leaf_mut().progress.verdict = None;
next.cursor.leaf_mut().route = None;
next.updated_at = OffsetDateTime::now_utc();
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let current =
super::durable::flow_position_in(&transaction, task_id)?.ok_or(StoreError::NotFound)?;
if current != *expected {
return Err(StoreError::InvalidAuthority(
"Task failure changed before retry".to_string(),
));
}
let position = super::durable::set_flow_position_in(&transaction, task_id, &next)?;
transaction.commit()?;
Ok(position)
}
pub(crate) fn restart_task_flow(
&self,
task: &Task,
expected: Option<&FlowPosition>,
checkpoint_head: &str,
) -> StoreResult<()> {
validate_task(task)?;
if checkpoint_head.trim().is_empty() {
return Err(StoreError::InvalidData(
"Task restart requires a checkpoint head".to_string(),
));
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
super::durable::require_ready_work(
&transaction,
&crate::durable::WorkRef::Task(task.id.clone()),
)?;
let current = super::durable::flow_position_in(&transaction, &task.id)?;
if current.as_ref() != expected
|| current
.as_ref()
.is_some_and(|position| position.claim.is_some())
{
return Err(StoreError::InvalidAuthority(
"Task execution changed before restart; stop the current worker before retrying"
.into(),
));
}
validate_task_project(&transaction, task)?;
let task_work = task_on(&transaction, &task.id)?.ok_or(StoreError::NotFound)?;
transaction.execute(
"DELETE FROM task_flow_positions WHERE task_id=?1",
[task.id.as_str()],
)?;
let parameters = task_params(task);
transaction.execute(
TASK_UPDATE,
rusqlite::params_from_iter(parameters.iter().map(|value| value.as_ref())),
)?;
insert_task_event_in(
&transaction,
&task_work,
&TaskEventKind::Progress {
summary: format!("Task restarted from checkpoint {checkpoint_head}"),
},
)?;
transaction.commit()?;
Ok(())
}
pub fn rebind_task_issue_identifier(
&self,
issue_id: &str,
old_identifier: &str,
new_identifier: &str,
) -> StoreResult<bool> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let Some((task_id, current_identifier)) = tx
.query_row(
"SELECT id, issue_identifier FROM tasks WHERE external_issue_id=?1",
[issue_id],
|row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
)
.optional()?
else {
return Ok(false);
};
if current_identifier == new_identifier {
return Ok(false);
}
if current_identifier != old_identifier {
return Err(StoreError::InvalidData(format!(
"Task {task_id} identifies issue {issue_id} as {current_identifier}, not {old_identifier}"
)));
}
let changed = tx.execute(
"UPDATE tasks SET issue_identifier=?3 WHERE id=?1 AND issue_identifier=?2",
params![task_id, old_identifier, new_identifier],
)?;
if changed == 0 {
return Err(StoreError::InvalidData(format!(
"Task {task_id} changed during its team migration"
)));
}
tx.execute(
"UPDATE tasks SET updated_at=?2 WHERE id=?1",
params![task_id, now_unix()],
)?;
tx.commit()?;
Ok(true)
}
pub fn complete_task(&self, task: &Task, skipped_pr: Option<&TaskPr>) -> StoreResult<()> {
validate_task(task)?;
if let Some(pr) = skipped_pr {
validate_task_pr(pr)?;
if pr.task_id != task.id || pr.phase() != PrPhase::Working {
return Err(StoreError::InvalidData(
"empty completion requires an unpublished Working Task PR".to_string(),
));
}
}
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
validate_task_project(&transaction, task)?;
if let Some(pr) = skipped_pr {
if transaction.execute(
"DELETE FROM task_prs
WHERE id=?1 AND task_id=?2
AND publication_requested_at IS NULL
AND merge_commit IS NULL AND abandoned_at IS NULL",
params![pr.id.as_str(), pr.task_id.as_str()],
)? == 0
{
return Err(StoreError::NotFound);
}
}
update_task_pm_writeback_in(&transaction, &task.id, &task.pm_writeback, task.updated_at)?;
complete_task_work_in(&transaction, task)?;
transaction.commit()?;
Ok(())
}
pub fn task(&self, task_id: &TaskId) -> StoreResult<Option<Task>> {
let conn = self.conn.lock().expect("store mutex poisoned");
task_on(&conn, task_id)
}
pub fn task_by_issue(&self, issue: &str) -> StoreResult<Option<Task>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let query = format!(
"{TASK_COLUMNS} WHERE t.id=?1 OR t.external_issue_id=?1 OR t.issue_identifier=?1"
);
let mut statement = conn.prepare(&query)?;
let rows = statement.query_map(params![issue], map_task_row)?;
let mut tasks = Vec::new();
for row in rows {
tasks.push(row?);
}
resolve_current_task(issue, tasks)
}
pub fn task_by_branch(&self, branch: &str) -> StoreResult<Option<Task>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let query = format!(
"{TASK_COLUMNS} WHERE EXISTS (
SELECT 1 FROM task_prs matched
WHERE matched.task_id=t.id AND matched.branch=?1
AND matched.abandoned_at IS NULL
AND (
matched.merge_commit IS NULL
OR NOT EXISTS (
SELECT 1 FROM task_prs active
WHERE active.task_id=t.id
AND active.merge_commit IS NULL
AND active.abandoned_at IS NULL
)
)
)"
);
let mut statement = conn.prepare(&query)?;
let rows = statement.query_map(params![branch], map_task_row)?;
let mut tasks = Vec::new();
for row in rows {
tasks.push(row?);
}
resolve_current_task(branch, tasks)
}
pub fn list_tasks(&self, wave_id: Option<&WaveId>) -> StoreResult<Vec<Task>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let (query, parameter): (String, Option<&dyn ToSql>) = match wave_id {
Some(wave_id) => (
format!("{TASK_COLUMNS} WHERE p.wave_id=?1 AND {TASK_VISIBLE} ORDER BY t.updated_at DESC"),
Some(wave_id as &dyn ToSql),
),
None => (format!("{TASK_COLUMNS} WHERE {TASK_VISIBLE} ORDER BY t.updated_at DESC"), None),
};
let mut statement = conn.prepare(&query)?;
let mut tasks = Vec::new();
if let Some(parameter) = parameter {
let rows = statement.query_map([parameter], map_task_row)?;
for row in rows {
tasks.push(row?);
}
} else {
let rows = statement.query_map([], map_task_row)?;
for row in rows {
tasks.push(row?);
}
}
Ok(tasks)
}
pub fn update_task_pr(&self, pr: &TaskPr) -> StoreResult<()> {
validate_task_pr(pr)?;
let conn = self.conn.lock().expect("store mutex poisoned");
let changed = update_task_pr(&conn, pr)?;
if changed == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub(crate) fn record_task_pr_repair_incident(
&self,
pr_id: &TaskPrId,
kind: TaskPrRepairKind,
occurred_at: OffsetDateTime,
) -> StoreResult<bool> {
let conn = self.conn.lock().expect("store mutex poisoned");
record_task_pr_repair_incident_on(&conn, pr_id, kind, occurred_at)
}
pub fn heal_task_pr_base(&self, pr: &TaskPr) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
if heal_task_pr_base(&conn, pr)? == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub fn task_prs(&self, task_id: &TaskId) -> StoreResult<Vec<TaskPr>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut statement = conn.prepare(&format!(
"{TASK_PR_COLUMNS} WHERE task_id=?1 ORDER BY sequence"
))?;
let rows = statement.query_map(params![task_id.as_str()], map_task_pr_row)?;
Ok(rows.collect::<Result<Vec<_>, _>>()?)
}
pub fn task_pr(&self, pr_id: &TaskPrId) -> StoreResult<Option<TaskPr>> {
let conn = self.conn.lock().expect("store mutex poisoned");
task_pr_on(&conn, pr_id)
}
pub fn active_task_pr(&self, task_id: &TaskId) -> StoreResult<Option<TaskPr>> {
let conn = self.conn.lock().expect("store mutex poisoned");
active_task_pr_on(&conn, task_id)
}
pub fn settle_task_pr(&self, settled: &TaskPr, next: Option<&TaskPr>) -> StoreResult<()> {
validate_task_pr_settlement(settled, next)?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
settle_task_pr_in(&transaction, settled, next)?;
transaction.commit()?;
Ok(())
}
pub fn rebase_task_pr(
&self,
pr_id: &TaskPrId,
new_base: &str,
clear_parent: bool,
updated_at: OffsetDateTime,
) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
let changed = if clear_parent {
conn.execute(
"UPDATE task_prs SET base_commit=?2, parent_pr_id=NULL, updated_at=?3 WHERE id=?1",
params![pr_id.as_str(), new_base, updated_at.unix_timestamp()],
)?
} else {
conn.execute(
"UPDATE task_prs SET base_commit=?2, updated_at=?3 WHERE id=?1 AND parent_pr_id IS NOT NULL",
params![pr_id.as_str(), new_base, updated_at.unix_timestamp()],
)?
};
if changed == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
pub(crate) fn settle_task_pr_merged(
&self,
settled: &TaskPr,
merged_at: Option<OffsetDateTime>,
) -> StoreResult<crate::store::TaskPrMergeEvidenceOutcome> {
validate_task_pr_settlement(settled, None)?;
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let outcome = settle_task_pr_merged_in(&transaction, settled, merged_at)?;
transaction.commit()?;
Ok(outcome)
}
pub fn task_linear_observation(
&self,
task_id: &TaskId,
) -> StoreResult<Option<TaskLinearObservation>> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.query_row(
"SELECT task_id, last_revision, last_title, last_description,
last_success_at, degraded_reason, updated_at
FROM task_linear_observations WHERE task_id=?1",
params![task_id.as_str()],
map_task_linear_observation_row,
)
.optional()
.map_err(StoreError::from)
}
pub fn apply_linear_observation(
&self,
apply: &LinearObservationApply,
) -> StoreResult<LinearObservationOutcome> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let existing = transaction
.query_row(
"SELECT last_revision, last_title, last_description
FROM task_linear_observations WHERE task_id=?1",
params![apply.task_id.as_str()],
|row| {
Ok((
row.get::<_, String>(0)?,
row.get::<_, String>(1)?,
row.get::<_, String>(2)?,
))
},
)
.optional()?;
let observed_at = apply.observed_at.unix_timestamp();
let mut follow_ups_created = Vec::new();
for follow_up in &apply.follow_ups {
if let Some(id) = ingest_linear_comment(
&transaction,
apply.task_id.as_str(),
&follow_up.comment_id,
&follow_up.text,
observed_at,
)? {
follow_ups_created.push(id);
}
}
let Some((last_revision, last_title, last_description)) = existing else {
transaction.execute(
"INSERT INTO task_linear_observations (
task_id, last_revision, last_title, last_description,
last_success_at, degraded_reason, updated_at
) VALUES (?1, ?2, ?3, ?4, ?5, NULL, ?5)",
params![
apply.task_id.as_str(),
apply.revision,
apply.title,
apply.description,
observed_at,
],
)?;
transaction.commit()?;
return Ok(LinearObservationOutcome {
baselined: true,
content_steer_applied: false,
follow_ups_created,
});
};
if apply.revision.as_str() < last_revision.as_str() {
transaction.commit()?;
return Ok(LinearObservationOutcome {
baselined: false,
content_steer_applied: false,
follow_ups_created,
});
}
let mut content_steer_applied = false;
if let Some(text) = &apply.content_steer {
if last_title != apply.title || last_description != apply.description {
Self::append_task_steer_in(&transaction, &apply.task_id, &Author::User, text)?;
content_steer_applied = true;
}
}
transaction.execute(
"UPDATE task_linear_observations
SET last_revision=?2, last_title=?3, last_description=?4,
last_success_at=?5, degraded_reason=NULL, updated_at=?5
WHERE task_id=?1",
params![
apply.task_id.as_str(),
apply.revision,
apply.title,
apply.description,
observed_at,
],
)?;
transaction.commit()?;
Ok(LinearObservationOutcome {
baselined: false,
content_steer_applied,
follow_ups_created,
})
}
pub fn apply_linear_comment(
&self,
task_id: &TaskId,
comment_id: &str,
text: &str,
observed_at: OffsetDateTime,
) -> StoreResult<Option<i64>> {
let observed_at = observed_at.unix_timestamp();
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let created = ingest_linear_comment(
&transaction,
task_id.as_str(),
comment_id,
text,
observed_at,
)?;
if created.is_some() {
transaction.execute(
"UPDATE task_linear_observations
SET last_success_at=?2, degraded_reason=NULL, updated_at=?2
WHERE task_id=?1",
params![task_id.as_str(), observed_at],
)?;
}
transaction.commit()?;
Ok(created)
}
pub fn mark_task_linear_degraded(&self, task_id: &TaskId, reason: &str) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"UPDATE task_linear_observations SET degraded_reason=?2, updated_at=?3
WHERE task_id=?1",
params![task_id.as_str(), reason, now_unix()],
)?;
Ok(())
}
pub fn append_task_event(
&self,
task_id: &TaskId,
kind: &TaskEventKind,
) -> StoreResult<TaskEvent> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let task = task_on(&transaction, task_id)?.ok_or(StoreError::NotFound)?;
let event = insert_task_event_in(&transaction, &task, kind)?;
transaction.commit()?;
Ok(event)
}
pub fn task_events_after(&self, task_id: &TaskId, cursor: i64) -> StoreResult<Vec<TaskEvent>> {
let conn = self.conn.lock().expect("store mutex poisoned");
task_events_after_in(&conn, task_id, cursor)
}
pub fn task_event(&self, task_id: &TaskId, event_id: i64) -> StoreResult<Option<TaskEvent>> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.query_row(
"SELECT id, task_id, kind_json, created_at
FROM task_events WHERE task_id = ?1 AND id = ?2",
params![task_id.as_str(), event_id],
map_task_event_row,
)
.optional()
.map_err(StoreError::from)
}
pub fn latest_task_event_at(&self, task_id: &TaskId) -> StoreResult<Option<OffsetDateTime>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let seconds: Option<i64> = conn.query_row(
"SELECT MAX(created_at) FROM task_events WHERE task_id = ?1",
params![task_id.as_str()],
|row| row.get(0),
)?;
Ok(seconds.map(crate::store::rows::unix_to_datetime))
}
pub fn recent_task_events(&self, task_id: &TaskId, limit: u32) -> StoreResult<Vec<TaskEvent>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut statement = conn.prepare(
"SELECT id, task_id, kind_json, created_at
FROM task_events WHERE task_id = ?1 ORDER BY id DESC LIMIT ?2",
)?;
let rows = statement.query_map(params![task_id.as_str(), limit], map_task_event_row)?;
let mut events = Vec::new();
for row in rows {
events.push(row?);
}
Ok(events)
}
pub fn latest_task_event(&self, task_id: &TaskId) -> StoreResult<Option<TaskEvent>> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.query_row(
"SELECT id, task_id, kind_json, created_at
FROM task_events WHERE task_id = ?1 ORDER BY id DESC LIMIT 1",
params![task_id.as_str()],
map_task_event_row,
)
.optional()
.map_err(StoreError::from)
}
pub fn insert_project(&self, project: &Project) -> StoreResult<()> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
transaction.execute(
PROJECT_INSERT,
rusqlite::params_from_iter(project_params(project).iter().map(|value| value.as_ref())),
)?;
create_project_work(&transaction, project)?;
transaction.commit()?;
Ok(())
}
pub fn update_project(&self, project: &Project) -> StoreResult<()> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let parameters = project_fact_params(project);
let changed = transaction.execute(
PROJECT_FACT_UPDATE,
rusqlite::params_from_iter(parameters.iter().map(|value| value.as_ref())),
)?;
if changed == 0 {
return Err(StoreError::NotFound);
}
transaction.commit()?;
Ok(())
}
pub fn project(&self, project_id: &ProjectId) -> StoreResult<Option<Project>> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.query_row(
PROJECT_SELECT,
params![project_id.as_str()],
map_project_row,
)
.optional()
.map_err(StoreError::from)
}
pub fn project_by_project(&self, project: &str) -> StoreResult<Option<Project>> {
if let Ok(project_id) = ProjectId::parse(project) {
return self.project(&project_id);
}
let conn = self.conn.lock().expect("store mutex poisoned");
let query = format!(
"{PROJECT_COLUMNS}
WHERE external_project_id=?1 OR project_slug=?1
ORDER BY created_at DESC, id DESC
LIMIT 1"
);
conn.query_row(&query, params![project], map_project_row)
.optional()
.map_err(StoreError::from)
}
pub fn list_projects(&self, wave_id: Option<&WaveId>) -> StoreResult<Vec<Project>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let query = match wave_id {
Some(_) => {
format!("{PROJECT_COLUMNS} WHERE wave_id=?1 ORDER BY updated_at DESC")
}
None => format!("{PROJECT_COLUMNS} ORDER BY updated_at DESC"),
};
let mut statement = conn.prepare(&query)?;
let mut projects = Vec::new();
if let Some(wave_id) = wave_id {
let rows = statement.query_map(params![wave_id], map_project_row)?;
for row in rows {
projects.push(row?);
}
} else {
let rows = statement.query_map([], map_project_row)?;
for row in rows {
projects.push(row?);
}
}
Ok(projects)
}
pub fn append_project_event(
&self,
project_id: &ProjectId,
kind: &ProjectEventKind,
) -> StoreResult<ProjectEvent> {
let mut conn = self.conn.lock().expect("store mutex poisoned");
let transaction = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let project = transaction.query_row(
PROJECT_SELECT,
params![project_id.as_str()],
map_project_row,
)?;
let event = insert_project_event_in(&transaction, &project, kind)?;
transaction.commit()?;
Ok(event)
}
pub fn project_events_after(
&self,
project_id: &ProjectId,
cursor: i64,
) -> StoreResult<Vec<ProjectEvent>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut statement = conn.prepare(
"SELECT id, project_id, kind_json, created_at
FROM project_events WHERE project_id=?1 AND id>?2 ORDER BY id",
)?;
let rows =
statement.query_map(params![project_id.as_str(), cursor], map_project_event_row)?;
let mut events = Vec::new();
for row in rows {
events.push(row?);
}
Ok(events)
}
pub fn latest_project_event_at(
&self,
project_id: &ProjectId,
) -> StoreResult<Option<OffsetDateTime>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let seconds: Option<i64> = conn.query_row(
"SELECT MAX(created_at) FROM project_events WHERE project_id = ?1",
params![project_id.as_str()],
|row| row.get(0),
)?;
Ok(seconds.map(crate::store::rows::unix_to_datetime))
}
pub fn latest_project_event(
&self,
project_id: &ProjectId,
) -> StoreResult<Option<ProjectEvent>> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.query_row(
"SELECT id, project_id, kind_json, created_at
FROM project_events WHERE project_id=?1 ORDER BY id DESC LIMIT 1",
params![project_id.as_str()],
map_project_event_row,
)
.optional()
.map_err(StoreError::from)
}
pub fn latest_project_failure(
&self,
project_id: &ProjectId,
) -> StoreResult<Option<crate::work::project::HistoricalFailure>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let event = conn
.query_row(
"SELECT id, project_id, kind_json, created_at
FROM project_events
WHERE project_id=?1 AND json_extract(kind_json, '$.kind')='failed'
ORDER BY id DESC LIMIT 1",
params![project_id.as_str()],
map_project_event_row,
)
.optional()?;
Ok(event
.as_ref()
.and_then(crate::work::project::HistoricalFailure::from_event))
}
pub fn pending_observations(
&self,
recipient: &ObservationRecipient,
) -> StoreResult<Vec<ObservationOutboxRow>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let (kind, id) = recipient_columns(recipient);
let mut statement = conn.prepare(
"SELECT id, recipient_kind, recipient_id, source_kind, source_id,
event_id, payload_json, delivered_at
FROM observation_outbox
WHERE recipient_kind=?1 AND recipient_id=?2 AND delivered_at IS NULL
ORDER BY id",
)?;
let rows = statement.query_map(params![kind, id], map_observation_row)?;
let mut observations = Vec::new();
for row in rows {
observations.push(row?);
}
Ok(observations)
}
pub fn pending_project_observations(
&self,
project_id: &ProjectId,
) -> StoreResult<Vec<ObservationOutboxRow>> {
let conn = self.conn.lock().expect("store mutex poisoned");
let mut statement = conn.prepare(
"SELECT id, recipient_kind, recipient_id, source_kind, source_id,
event_id, payload_json, delivered_at
FROM observation_outbox
WHERE recipient_kind='project'
AND recipient_id=?1
AND delivered_at IS NULL
ORDER BY id",
)?;
let rows = statement.query_map(params![project_id.as_str()], map_observation_row)?;
let mut observations = Vec::new();
for row in rows {
observations.push(row?);
}
Ok(observations)
}
pub fn mark_observation_delivered(&self, id: i64) -> StoreResult<()> {
let conn = self.conn.lock().expect("store mutex poisoned");
conn.execute(
"UPDATE observation_outbox SET delivered_at=?1
WHERE id=?2 AND delivered_at IS NULL",
params![now_unix(), id],
)?;
Ok(())
}
}
fn validate_task(task: &Task) -> StoreResult<()> {
task.validate()
.map_err(|error| StoreError::InvalidData(error.to_string()))
}
fn update_task_timestamp_in(conn: &Connection, task: &Task) -> StoreResult<()> {
if conn.execute(
"UPDATE tasks SET updated_at=?2 WHERE id=?1",
params![task.id.as_str(), task.updated_at.unix_timestamp()],
)? != 1
{
return Err(StoreError::NotFound);
}
Ok(())
}
fn matching_claim_in(
conn: &Connection,
task_id: &TaskId,
expected: &TaskWorkerClaim,
) -> StoreResult<TaskWorkerClaim> {
let current = super::durable::flow_position_in(conn, task_id)?
.and_then(|position| position.claim)
.ok_or_else(|| {
StoreError::InvalidAuthority(format!("Task {task_id} has no active worker"))
})?;
if current.invocation_id != expected.invocation_id
|| current.position_version != expected.position_version
|| current.generation != expected.generation
{
return Err(StoreError::InvalidAuthority(format!(
"Task {task_id} worker changed"
)));
}
Ok(current)
}
fn update_task_pm_writeback_in(
conn: &Connection,
task_id: &TaskId,
state: &PmWritebackState,
updated_at: OffsetDateTime,
) -> StoreResult<()> {
let changed = conn.execute(
"UPDATE tasks SET pm_writeback_json=?2, updated_at=?3 WHERE id=?1",
params![
task_id.as_str(),
serde_json::to_string(state).expect("Task PM writeback state must serialize"),
updated_at.unix_timestamp(),
],
)?;
if changed == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
fn complete_task_work_in(conn: &Connection, task: &Task) -> StoreResult<()> {
let state: String = conn.query_row(
"SELECT work_state FROM tasks WHERE id=?1",
[task.id.as_str()],
|row| row.get(0),
)?;
match state.as_str() {
"done" => return Ok(()),
"ready" => {}
"abandoned" => {
return Err(StoreError::InvalidData(format!(
"Task {} Work is abandoned and cannot be completed",
task.id
)))
}
other => {
return Err(StoreError::InvalidData(format!(
"Task {} Work has invalid state {other:?}",
task.id
)))
}
}
if conn.execute(
"UPDATE tasks SET work_state='done', work_terminal_at=?2
WHERE id=?1 AND work_state='ready'",
params![task.id.as_str(), now_unix()],
)? != 1
{
return Err(StoreError::InvalidData(format!(
"Task {} Work changed while completion was being recorded",
task.id
)));
}
insert_task_event_in(
conn,
task,
&TaskEventKind::Completed {
summary: "Task completed".to_string(),
},
)?;
Ok(())
}
fn resolve_current_task(key: &str, mut tasks: Vec<Task>) -> StoreResult<Option<Task>> {
if tasks.len() > 1 {
return Err(StoreError::InvalidData(format!(
"multiple stable Tasks resolve to {key:?}"
)));
}
Ok(tasks.pop())
}
fn validate_task_pr(pr: &TaskPr) -> StoreResult<()> {
pr.validate()
.map_err(|error| StoreError::InvalidData(error.to_string()))
}
fn validate_initial_task_pr(task: &Task, pr: &TaskPr) -> StoreResult<()> {
validate_task_pr(pr)?;
if pr.task_id != task.id || pr.sequence != 1 || pr.phase() != PrPhase::Working {
return Err(StoreError::InvalidData(
"Task requires its sequence-1 Working PR".to_string(),
));
}
Ok(())
}
fn require_task_not_deleted(conn: &Connection, task: &Task) -> StoreResult<()> {
let deleted: bool = conn.query_row(
"SELECT EXISTS(SELECT 1 FROM task_deletions WHERE wave_id=?1 AND issue_id=?2)",
params![task.wave_id.as_str(), task.plan.id.as_str()],
|row| row.get(0),
)?;
if deleted {
return Err(StoreError::InvalidAuthority(format!(
"Task {} was deleted; create a new Task",
task.plan.identifier
)));
}
Ok(())
}
fn insert_initial_task(
conn: &rusqlite::Transaction<'_>,
task: &Task,
pr: &TaskPr,
) -> StoreResult<()> {
validate_task(task)?;
validate_initial_task_pr(task, pr)?;
validate_task_project(conn, task)?;
require_task_not_deleted(conn, task)?;
let expired: bool = conn.query_row(
"SELECT EXISTS(SELECT 1 FROM projects p JOIN wave_chapters c ON c.wave_id=p.wave_id AND c.current=1
WHERE p.id=?1 AND p.external_project_id != c.project_id)",
[task.project_id.as_str()], |row| row.get(0),
)?;
if expired {
return Err(StoreError::InvalidData(
"cannot prepare a new Task in chapter history; refresh the Wave".into(),
));
}
let mut parameters = task_params(task);
parameters.push(Box::new(task.agent.clone()));
conn.execute(
TASK_INSERT,
rusqlite::params_from_iter(parameters.iter().map(|value| value.as_ref())),
)?;
create_task_work(conn, task)?;
insert_task_pr(conn, pr)?;
seed_task_linear_observation(conn, task)
}
fn seed_task_linear_observation(conn: &Connection, task: &Task) -> StoreResult<()> {
conn.execute(
"INSERT OR IGNORE INTO task_linear_observations (
task_id, last_revision, last_title, last_description,
last_success_at, degraded_reason, updated_at
) VALUES (?1, '', ?2, ?3, ?4, NULL, ?4)",
params![
task.id.as_str(),
task.plan.title,
task.plan.description,
now_unix(),
],
)?;
Ok(())
}
fn validate_task_project(conn: &Connection, task: &Task) -> StoreResult<()> {
let owner = conn
.query_row(
"SELECT wave_id FROM projects WHERE id=?1",
params![task.project_id.as_str()],
|row| row.get::<_, String>(0),
)
.optional()?;
let Some(wave_id) = owner else {
return Err(StoreError::InvalidData(format!(
"Task {} requires Project {}",
task.id, task.project_id
)));
};
if wave_id != task.wave_id.as_str() {
return Err(StoreError::InvalidData(format!(
"Project {} does not belong to Task {}'s Wave {}",
task.project_id, task.id, task.wave_id
)));
}
Ok(())
}
const TASK_INSERT: &str = "INSERT INTO tasks (
id, project_id, external_issue_id, issue_identifier, issue_title,
issue_description, pm_snapshot_synced_at, pm_writeback_json,
worktree, workspace_slug,
abandon_requested_at, abandon_reason, created_at, updated_at, agent
) VALUES (
?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15
)";
const TASK_VISIBLE: &str = "NOT EXISTS (
SELECT 1 FROM task_deletions d WHERE d.wave_id=p.wave_id AND d.issue_id=t.external_issue_id
)";
const TASK_COLUMNS: &str = "SELECT
t.id, t.external_issue_id, t.issue_identifier, t.issue_title, t.issue_description,
p.wave_id, t.worktree, t.workspace_slug,
t.created_at, t.updated_at, t.pm_snapshot_synced_at, t.pm_writeback_json,
t.project_id, t.abandon_requested_at, t.abandon_reason, t.agent
FROM tasks t JOIN projects p ON p.id=t.project_id";
const TASK_UPDATE: &str = "UPDATE tasks SET
external_issue_id=?3, issue_identifier=?4,
issue_title=?5, issue_description=?6, pm_snapshot_synced_at=?7,
pm_writeback_json=?8, worktree=?9, workspace_slug=?10,
abandon_requested_at=?11, abandon_reason=?12, created_at=?13, updated_at=?14
WHERE id=?1";
const TASK_PR_COLUMNS: &str = "SELECT
id, task_id, sequence, slug, branch, base_commit,
publication_requested_at, after_merge, next_slug, github_number, github_url,
merge_commit, abandoned_at, created_at, updated_at,
github_head_sha, ci_observation, parent_pr_id, github_observation,
linear_attachment_id, linear_comment_id, linear_link_error,
merge_mode, merge_requested_at, merge_head_sha,
pr_title, pr_body, pr_copy_head_sha
FROM task_prs";
const TASK_PR_SELECT: &str = "SELECT
id, task_id, sequence, slug, branch, base_commit,
publication_requested_at, after_merge, next_slug, github_number, github_url,
merge_commit, abandoned_at, created_at, updated_at,
github_head_sha, ci_observation, parent_pr_id, github_observation,
linear_attachment_id, linear_comment_id, linear_link_error,
merge_mode, merge_requested_at, merge_head_sha,
pr_title, pr_body, pr_copy_head_sha
FROM task_prs WHERE id=?1";
fn ingest_linear_comment(
conn: &rusqlite::Transaction<'_>,
task_id: &str,
comment_id: &str,
text: &str,
observed_at: i64,
) -> StoreResult<Option<i64>> {
if let Some((id, _)) = comment_id.split_once('@') {
let prefix = format!("{id}@");
let latest: Option<String> = conn.query_row(
"SELECT MAX(comment_id) FROM task_linear_ingested_comments WHERE task_id=?1 AND substr(comment_id, 1, length(?2))=?2",
params![task_id, prefix], |row| row.get(0),
)?;
if latest.as_deref().is_some_and(|latest| latest > comment_id) {
return Ok(None);
}
}
let inserted = conn.execute(
"INSERT OR IGNORE INTO task_linear_ingested_comments
(task_id, comment_id, ingested_at) VALUES (?1, ?2, ?3)",
params![task_id, comment_id, observed_at],
)?;
if inserted == 1 {
let steer = SqliteStore::append_task_steer_in(
conn,
&TaskId::from_raw(task_id),
&Author::User,
text,
)?;
Ok(Some(steer.id))
} else {
Ok(None)
}
}
fn task_params(task: &Task) -> Vec<Box<dyn ToSql>> {
vec![
Box::new(task.id.as_str().to_string()),
Box::new(task.project_id.as_str().to_string()),
Box::new(task.plan.id.as_str().to_string()),
Box::new(task.plan.identifier.clone()),
Box::new(task.plan.title.clone()),
Box::new(task.plan.description.clone()),
Box::new(task.plan.pm_snapshot_synced_at),
Box::new(
serde_json::to_string(&task.pm_writeback)
.expect("Task PM writeback state must serialize"),
),
Box::new(task.worktree.display().to_string()),
Box::new(task.workspace_slug.clone()),
Box::new(
task.abandon_intent
.as_ref()
.map(|intent| intent.requested_at.unix_timestamp()),
),
Box::new(
task.abandon_intent
.as_ref()
.map(|intent| intent.reason.clone()),
),
Box::new(task.created_at.unix_timestamp()),
Box::new(task.updated_at.unix_timestamp()),
]
}
fn insert_task_pr(conn: &Connection, pr: &TaskPr) -> StoreResult<()> {
validate_task_pr(pr)?;
let publication = pr.publication.as_ref();
let presentation = publication.and_then(|publication| publication.presentation.as_ref());
let github = publication.and_then(|publication| publication.github.as_ref());
let merge = publication.and_then(|publication| publication.merge.as_ref());
conn.execute(
"INSERT INTO task_prs (
id, task_id, sequence, slug, branch, base_commit,
publication_requested_at, after_merge, next_slug,
github_number, github_url, merge_commit, abandoned_at,
created_at, updated_at, github_head_sha, ci_observation, parent_pr_id,
github_observation,
linear_attachment_id, linear_comment_id, linear_link_error,
merge_mode, merge_requested_at, merge_head_sha,
pr_title, pr_body, pr_copy_head_sha
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20, ?21, ?22, ?23, ?24, ?25, ?26, ?27, ?28)",
params![
pr.id.as_str(),
pr.task_id.as_str(),
i64::from(pr.sequence),
pr.slug,
pr.branch,
pr.base_commit,
publication.map(|publication| publication.requested_at.unix_timestamp()),
merge.map(|request| request.after_merge.as_str()),
merge.and_then(|request| request.next_slug.as_deref()),
github.map(|github| i64::from(github.number)),
github.map(|github| github.url.as_str()),
pr.merge_commit,
pr.abandoned_at.map(OffsetDateTime::unix_timestamp),
pr.created_at.unix_timestamp(),
pr.updated_at.unix_timestamp(),
github.and_then(|github| github.head_sha.as_deref()),
task_pr_ci_json(pr)?,
pr.parent_pr_id.as_ref().map(TaskPrId::as_str),
task_pr_github_observation_json(pr)?,
pr.linear_attachment_id.as_deref(),
pr.linear_comment_id.as_deref(),
pr.linear_link_error.as_deref(),
merge.map(|request| request.mode.as_str()),
merge.map(|request| request.requested_at.unix_timestamp()),
merge.map(|request| request.head_sha.as_str()),
presentation.map(|copy| copy.title.as_str()),
presentation.map(|copy| copy.body.as_str()),
presentation.map(|copy| copy.head_sha.as_str()),
],
)?;
Ok(())
}
fn record_task_pr_repair_incident_on(
conn: &Connection,
pr_id: &TaskPrId,
kind: TaskPrRepairKind,
occurred_at: OffsetDateTime,
) -> StoreResult<bool> {
Ok(conn.execute(
"INSERT INTO task_pr_repair_incidents (task_pr_id, kind, occurred_at)
VALUES (?1, ?2, ?3)
ON CONFLICT(task_pr_id, kind) DO NOTHING",
params![pr_id.as_str(), kind.as_str(), occurred_at.unix_timestamp()],
)? == 1)
}
fn update_task_pr(conn: &Connection, pr: &TaskPr) -> StoreResult<usize> {
validate_task_pr(pr)?;
let publication = pr.publication.as_ref();
let presentation = publication.and_then(|publication| publication.presentation.as_ref());
let github = publication.and_then(|publication| publication.github.as_ref());
let merge = publication.and_then(|publication| publication.merge.as_ref());
conn.execute(
"UPDATE task_prs SET
publication_requested_at=?7, after_merge=?8, next_slug=?9,
github_number=?10, github_url=?11, merge_commit=?12,
abandoned_at=?13, updated_at=?15, github_head_sha=?16,
ci_observation=?17, parent_pr_id=?18, github_observation=?19,
linear_attachment_id=?20, linear_comment_id=?21, linear_link_error=?22,
merge_mode=?23, merge_requested_at=?24, merge_head_sha=?25,
pr_title=?26, pr_body=?27, pr_copy_head_sha=?28
WHERE id=?1 AND task_id=?2 AND sequence=?3 AND slug=?4
AND branch=?5 AND base_commit=?6 AND created_at=?14",
params![
pr.id.as_str(),
pr.task_id.as_str(),
i64::from(pr.sequence),
pr.slug,
pr.branch,
pr.base_commit,
publication.map(|publication| publication.requested_at.unix_timestamp()),
merge.map(|request| request.after_merge.as_str()),
merge.and_then(|request| request.next_slug.as_deref()),
github.map(|github| i64::from(github.number)),
github.map(|github| github.url.as_str()),
pr.merge_commit,
pr.abandoned_at.map(OffsetDateTime::unix_timestamp),
pr.created_at.unix_timestamp(),
pr.updated_at.unix_timestamp(),
github.and_then(|github| github.head_sha.as_deref()),
task_pr_ci_json(pr)?,
pr.parent_pr_id.as_ref().map(TaskPrId::as_str),
task_pr_github_observation_json(pr)?,
pr.linear_attachment_id.as_deref(),
pr.linear_comment_id.as_deref(),
pr.linear_link_error.as_deref(),
merge.map(|request| request.mode.as_str()),
merge.map(|request| request.requested_at.unix_timestamp()),
merge.map(|request| request.head_sha.as_str()),
presentation.map(|copy| copy.title.as_str()),
presentation.map(|copy| copy.body.as_str()),
presentation.map(|copy| copy.head_sha.as_str()),
],
)
.map_err(StoreError::from)
}
fn heal_task_pr_base(conn: &Connection, pr: &TaskPr) -> StoreResult<usize> {
validate_task_pr(pr)?;
conn.execute(
"UPDATE task_prs SET base_commit=?4, updated_at=?5
WHERE id=?1 AND task_id=?2 AND sequence=?3",
params![
pr.id.as_str(),
pr.task_id.as_str(),
i64::from(pr.sequence),
pr.base_commit,
pr.updated_at.unix_timestamp(),
],
)
.map_err(StoreError::from)
}
fn task_pr_ci_json(pr: &TaskPr) -> StoreResult<Option<String>> {
pr.ci_observation
.as_ref()
.map(serde_json::to_string)
.transpose()
.map_err(StoreError::from)
}
fn task_pr_github_observation_json(pr: &TaskPr) -> StoreResult<Option<String>> {
pr.github_observation
.as_ref()
.map(serde_json::to_string)
.transpose()
.map_err(StoreError::from)
}
pub(super) fn task_on(conn: &Connection, task_id: &TaskId) -> StoreResult<Option<Task>> {
conn.query_row(
&format!("{TASK_COLUMNS} WHERE t.id=?1"),
[task_id.as_str()],
map_task_row,
)
.optional()
.map_err(StoreError::from)
}
fn task_pr_on(conn: &Connection, pr_id: &TaskPrId) -> StoreResult<Option<TaskPr>> {
conn.query_row(TASK_PR_SELECT, params![pr_id.as_str()], map_task_pr_row)
.optional()
.map_err(StoreError::from)
}
fn active_task_pr_on(conn: &Connection, task_id: &TaskId) -> StoreResult<Option<TaskPr>> {
let query = format!(
"{TASK_PR_COLUMNS}
WHERE task_id=?1 AND merge_commit IS NULL AND abandoned_at IS NULL"
);
conn.query_row(&query, [task_id.as_str()], map_task_pr_row)
.optional()
.map_err(StoreError::from)
}
fn settle_task_pr_on(conn: &Connection, settled: &TaskPr) -> StoreResult<()> {
let current = task_pr_on(conn, &settled.id)?.ok_or(StoreError::NotFound)?;
if current.is_settled() {
if !same_task_pr(¤t, settled) {
return Err(StoreError::InvalidData(format!(
"Task PR {} is already settled differently",
settled.id
)));
}
} else if update_task_pr(conn, settled)? == 0 {
return Err(StoreError::NotFound);
}
Ok(())
}
fn validate_task_pr_settlement(settled: &TaskPr, next: Option<&TaskPr>) -> StoreResult<()> {
validate_task_pr(settled)?;
if !settled.is_settled() {
return Err(StoreError::InvalidData(
"Task PR transition requires a settled PR".to_string(),
));
}
if let Some(next) = next {
validate_task_pr(next)?;
if next.task_id != settled.task_id
|| next.sequence != settled.sequence + 1
|| next.phase() != PrPhase::Working
{
return Err(StoreError::InvalidData(
"next Task PR must be the following Working PR for the same Task".to_string(),
));
}
}
Ok(())
}
fn settle_task_pr_in(
conn: &Connection,
settled: &TaskPr,
next: Option<&TaskPr>,
) -> StoreResult<()> {
settle_task_pr_on(conn, settled)?;
let Some(next) = next else {
return Ok(());
};
let query = format!("{TASK_PR_COLUMNS} WHERE task_id=?1 AND sequence=?2");
let existing = conn
.query_row(
&query,
params![next.task_id.as_str(), i64::from(next.sequence)],
map_task_pr_row,
)
.optional()?;
match existing {
Some(existing) if same_task_pr(&existing, next) => Ok(()),
Some(existing) => Err(StoreError::InvalidData(format!(
"Task PR sequence {} already belongs to {}",
next.sequence, existing.id
))),
None => insert_task_pr(conn, next),
}
}
fn settle_task_pr_merged_in(
conn: &Connection,
settled: &TaskPr,
merged_at: Option<OffsetDateTime>,
) -> StoreResult<crate::store::TaskPrMergeEvidenceOutcome> {
let has_authority_column: bool = conn.query_row(
"SELECT EXISTS(
SELECT 1 FROM pragma_table_info('task_prs') WHERE name='merged_at'
)",
[],
|row| row.get(0),
)?;
if !has_authority_column {
settle_task_pr_in(conn, settled, None)?;
return Ok(crate::store::TaskPrMergeEvidenceOutcome::SchemaUnavailable);
}
let accepted_at = conn.query_row(
"SELECT merged_at FROM task_prs WHERE id=?1",
[settled.id.as_str()],
|row| row.get::<_, Option<i64>>(0),
)?;
let observed_at = merged_at.map(OffsetDateTime::unix_timestamp);
let outcome = match (accepted_at, observed_at) {
(None, Some(_)) => crate::store::TaskPrMergeEvidenceOutcome::Accepted,
(Some(accepted), Some(observed)) if accepted == observed => {
crate::store::TaskPrMergeEvidenceOutcome::Repeated
}
(Some(accepted), Some(_)) => crate::store::TaskPrMergeEvidenceOutcome::Conflict {
accepted_at: accepted,
},
(_, None) => crate::store::TaskPrMergeEvidenceOutcome::Missing,
};
settle_task_pr_in(conn, settled, None)?;
if let crate::store::TaskPrMergeEvidenceOutcome::Conflict { accepted_at } = outcome {
let checked_at = settled
.github_observation
.as_ref()
.map(|observation| observation.checked_at)
.unwrap_or_else(OffsetDateTime::now_utc);
let observation = GithubObservation {
checked_at,
result: crate::work::task::GithubObservationResult::Partial {
reason: format!(
"GitHub merged_at conflicts with first accepted value {accepted_at}"
),
},
};
conn.execute(
"UPDATE task_prs SET github_observation=?2 WHERE id=?1",
params![settled.id.as_str(), serde_json::to_string(&observation)?],
)?;
}
if let Some(observed_at) = observed_at {
conn.execute(
"UPDATE task_prs SET merged_at=COALESCE(merged_at, ?2) WHERE id=?1",
params![settled.id.as_str(), observed_at],
)?;
}
Ok(outcome)
}
fn same_task_pr(left: &TaskPr, right: &TaskPr) -> bool {
left.id == right.id
&& left.task_id == right.task_id
&& left.sequence == right.sequence
&& left.slug == right.slug
&& left.branch == right.branch
&& left.base_commit == right.base_commit
&& left.publication == right.publication
&& left.merge_commit == right.merge_commit
&& same_settle_instant(left.abandoned_at, right.abandoned_at)
}
fn same_settle_instant(
left: Option<time::OffsetDateTime>,
right: Option<time::OffsetDateTime>,
) -> bool {
left.map(time::OffsetDateTime::unix_timestamp)
== right.map(time::OffsetDateTime::unix_timestamp)
}
fn invalid_column(
index: usize,
error: impl std::error::Error + Send + Sync + 'static,
) -> rusqlite::Error {
rusqlite::Error::FromSqlConversionFailure(index, rusqlite::types::Type::Text, Box::new(error))
}
fn map_task_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<Task> {
let abandon_intent = match (
row.get::<_, Option<i64>>(13)?,
row.get::<_, Option<String>>(14)?,
) {
(Some(requested_at), Some(reason)) => Some(AbandonIntent {
requested_at: crate::store::rows::unix_to_datetime(requested_at),
reason,
}),
_ => None,
};
Ok(Task {
id: TaskId::from_raw(row.get::<_, String>(0)?),
plan: TaskPlan {
id: LinearIssueId::from_raw(row.get::<_, String>(1)?),
identifier: row.get(2)?,
title: row.get(3)?,
description: row.get(4)?,
pm_snapshot_synced_at: row.get(10)?,
},
pm_writeback: serde_json::from_str(&row.get::<_, String>(11)?)
.map_err(|error| invalid_column(11, error))?,
wave_id: row.get(5)?,
project_id: ProjectId::from_raw(row.get::<_, String>(12)?),
worktree: PathBuf::from(row.get::<_, String>(6)?),
workspace_slug: row.get(7)?,
agent: row.get(15)?,
abandon_intent,
created_at: crate::store::rows::unix_to_datetime(row.get(8)?),
updated_at: crate::store::rows::unix_to_datetime(row.get(9)?),
observation: crate::work::task::Observation::NotRequired,
})
}
fn map_task_pr_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<TaskPr> {
let publication_requested_at = row.get::<_, Option<i64>>(6)?;
let after_merge = row
.get::<_, Option<String>>(7)?
.map(|value| value.parse())
.transpose()
.map_err(|error| invalid_column(7, error))?;
let github_number = row.get::<_, Option<i64>>(9)?.map(|number| number as u32);
let github_url = row.get::<_, Option<String>>(10)?;
let github_head_sha = row.get::<_, Option<String>>(15)?;
let ci_observation = row
.get::<_, Option<String>>(16)?
.map(|json| serde_json::from_str::<CiObservation>(&json))
.transpose()
.map_err(|error| invalid_column(16, error))?;
let github_observation = row
.get::<_, Option<String>>(18)?
.map(|json| serde_json::from_str::<GithubObservation>(&json))
.transpose()
.map_err(|error| invalid_column(18, error))?;
let merge = match (
row.get::<_, Option<String>>(22)?,
row.get::<_, Option<i64>>(23)?,
row.get::<_, Option<String>>(24)?,
after_merge,
) {
(Some(mode), Some(requested_at), Some(head_sha), Some(after_merge)) => {
Some(PrMergeRequest {
mode: mode.parse().map_err(|error| invalid_column(22, error))?,
requested_at: crate::store::rows::unix_to_datetime(requested_at),
head_sha,
after_merge,
next_slug: row.get(8)?,
})
}
(None, None, None, None) => None,
_ => {
return Err(invalid_column(
22,
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"PR merge request fields must all be present or absent",
),
))
}
};
let presentation = match (
row.get::<_, Option<String>>(25)?,
row.get::<_, Option<String>>(26)?,
row.get::<_, Option<String>>(27)?,
) {
(Some(title), Some(body), Some(head_sha)) => Some(PrPresentation {
title,
body,
head_sha,
}),
(None, None, None) => None,
_ => {
return Err(invalid_column(
25,
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"PR presentation fields must all be present or absent",
),
))
}
};
let publication = match publication_requested_at {
Some(requested_at) => Some(PrPublication {
requested_at: crate::store::rows::unix_to_datetime(requested_at),
presentation,
github: match (github_number, github_url) {
(Some(number), Some(url)) => Some(GithubPr {
number,
url,
head_sha: github_head_sha,
}),
(None, None) => None,
_ => {
return Err(invalid_column(
9,
std::io::Error::new(
std::io::ErrorKind::InvalidData,
"GitHub PR number and URL must both be present or absent",
),
))
}
},
merge,
}),
None => None,
};
let pr = TaskPr {
id: TaskPrId::from_raw(row.get::<_, String>(0)?),
task_id: TaskId::from_raw(row.get::<_, String>(1)?),
sequence: row.get::<_, i64>(2)? as u32,
slug: row.get(3)?,
branch: row.get(4)?,
base_commit: row.get(5)?,
parent_pr_id: row.get::<_, Option<String>>(17)?.map(TaskPrId::from_raw),
publication,
merge_commit: row.get(11)?,
abandoned_at: row
.get::<_, Option<i64>>(12)?
.map(crate::store::rows::unix_to_datetime),
ci_observation,
github_observation,
linear_attachment_id: row.get::<_, Option<String>>(19)?,
linear_comment_id: row.get::<_, Option<String>>(20)?,
linear_link_error: row.get::<_, Option<String>>(21)?,
created_at: crate::store::rows::unix_to_datetime(row.get(13)?),
updated_at: crate::store::rows::unix_to_datetime(row.get(14)?),
};
pr.validate_persisted()
.map_err(|error| invalid_column(6, error))?;
Ok(pr)
}
fn map_task_linear_observation_row(
row: &rusqlite::Row<'_>,
) -> rusqlite::Result<TaskLinearObservation> {
let last_success_at = OffsetDateTime::from_unix_timestamp(row.get::<_, i64>(4)?)
.map_err(|error| invalid_column(4, error))?;
let updated_at = OffsetDateTime::from_unix_timestamp(row.get::<_, i64>(6)?)
.map_err(|error| invalid_column(6, error))?;
Ok(TaskLinearObservation {
task_id: TaskId::from_raw(row.get::<_, String>(0)?),
last_revision: row.get(1)?,
last_title: row.get(2)?,
last_description: row.get(3)?,
last_success_at,
degraded_reason: row.get(5)?,
updated_at,
})
}
fn map_task_event_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<TaskEvent> {
let kind_json: String = row.get(2)?;
let kind: TaskEventKind =
serde_json::from_str(&kind_json).map_err(|error| invalid_column(2, error))?;
Ok(TaskEvent {
id: row.get(0)?,
task_id: TaskId::from_raw(row.get::<_, String>(1)?),
kind,
created_at: crate::store::rows::unix_to_datetime(row.get(3)?),
})
}
pub(super) fn task_events_after_in(
conn: &Connection,
task_id: &TaskId,
cursor: i64,
) -> StoreResult<Vec<TaskEvent>> {
let mut statement = conn.prepare(
"SELECT id, task_id, kind_json, created_at
FROM task_events WHERE task_id=?1 AND id>?2 ORDER BY id",
)?;
let rows = statement.query_map(params![task_id.as_str(), cursor], map_task_event_row)?;
rows.collect::<Result<Vec<_>, _>>()
.map_err(StoreError::from)
}
const PROJECT_INSERT: &str = "INSERT INTO projects (
id, wave_id, external_project_id, project_slug, project_name,
project_prompt_context, pm_snapshot_synced_at,
abandon_requested_at, abandon_reason,
created_at, updated_at, iteration
) VALUES (
?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12
)";
const PROJECT_COLUMNS: &str = "SELECT
id, external_project_id, project_slug, project_name, project_prompt_context,
wave_id, pm_snapshot_synced_at, abandon_requested_at, abandon_reason,
created_at, updated_at, iteration
FROM projects";
pub(super) const PROJECT_SELECT: &str = "SELECT
id, external_project_id, project_slug, project_name, project_prompt_context,
wave_id, pm_snapshot_synced_at, abandon_requested_at, abandon_reason,
created_at, updated_at, iteration
FROM projects WHERE id=?1";
const PROJECT_FACT_UPDATE: &str = "UPDATE projects SET
wave_id=?2, external_project_id=?3, project_slug=?4, project_name=?5,
project_prompt_context=?6, pm_snapshot_synced_at=?7,
abandon_requested_at=?8, abandon_reason=?9,
created_at=?10, updated_at=?11
WHERE id=?1";
fn project_params(project: &Project) -> Vec<Box<dyn ToSql>> {
vec![
Box::new(project.id.as_str().to_string()),
Box::new(project.wave_id.clone()),
Box::new(project.plan.id.as_str().to_string()),
Box::new(project.plan.slug.clone()),
Box::new(project.plan.name.clone()),
Box::new(project.plan.prompt_context.clone()),
Box::new(project.plan.pm_snapshot_synced_at),
Box::new(
project
.abandon_intent
.as_ref()
.map(|intent| intent.requested_at.unix_timestamp()),
),
Box::new(
project
.abandon_intent
.as_ref()
.map(|intent| intent.reason.clone()),
),
Box::new(project.created_at.unix_timestamp()),
Box::new(project.updated_at.unix_timestamp()),
Box::new(project.iteration),
]
}
fn project_fact_params(project: &Project) -> Vec<Box<dyn ToSql>> {
project_params(project).into_iter().take(11).collect()
}
pub(super) fn map_project_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<Project> {
let abandon_intent = match (
row.get::<_, Option<i64>>(7)?,
row.get::<_, Option<String>>(8)?,
) {
(Some(requested_at), Some(reason)) => Some(AbandonIntent {
requested_at: crate::store::rows::unix_to_datetime(requested_at),
reason,
}),
_ => None,
};
Ok(Project {
id: ProjectId::from_raw(row.get::<_, String>(0)?),
plan: ProjectPlan {
id: LinearProjectId::from_raw(row.get::<_, String>(1)?),
slug: row.get(2)?,
name: row.get(3)?,
prompt_context: row.get(4)?,
pm_snapshot_synced_at: row.get(6)?,
},
wave_id: row.get(5)?,
iteration: row.get::<_, i64>(11)? as u32,
abandon_intent,
created_at: crate::store::rows::unix_to_datetime(row.get(9)?),
updated_at: crate::store::rows::unix_to_datetime(row.get(10)?),
})
}
fn map_project_event_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<ProjectEvent> {
let kind: ProjectEventKind = serde_json::from_str(&row.get::<_, String>(2)?)
.map_err(|error| invalid_column(2, error))?;
Ok(ProjectEvent {
id: row.get(0)?,
project_id: ProjectId::from_raw(row.get::<_, String>(1)?),
kind,
created_at: crate::store::rows::unix_to_datetime(row.get(3)?),
})
}
fn recipient_columns(recipient: &ObservationRecipient) -> (&'static str, String) {
match recipient {
ObservationRecipient::Wave { wave_id } => ("wave", wave_id.as_str().to_string()),
ObservationRecipient::Project { project_id } => {
("project", project_id.as_str().to_string())
}
}
}
fn child_columns(source: &ChildRef) -> (&'static str, String) {
match source {
ChildRef::Project(project_id) => ("project", project_id.as_str().to_string()),
ChildRef::Task(task_id) => ("task", task_id.as_str().to_string()),
}
}
pub(super) fn insert_task_event_in(
conn: &Connection,
task: &Task,
kind: &TaskEventKind,
) -> StoreResult<TaskEvent> {
let created_at = now_unix();
conn.execute(
"INSERT INTO task_events (task_id, kind_json, created_at) VALUES (?1, ?2, ?3)",
params![task.id.as_str(), serde_json::to_string(kind)?, created_at],
)?;
let event_id = conn.last_insert_rowid();
if kind.is_wave_observable() {
insert_observation(
conn,
&ObservationRecipient::Wave {
wave_id: task.wave_id.clone(),
},
&ChildRef::Task(task.id.clone()),
event_id,
&ChildEventPayload::Task {
event: kind.clone(),
},
created_at,
)?;
}
Ok(TaskEvent {
id: event_id,
task_id: task.id.clone(),
kind: kind.clone(),
created_at: crate::store::rows::unix_to_datetime(created_at),
})
}
pub(super) fn insert_project_event_in(
conn: &Connection,
project: &Project,
kind: &ProjectEventKind,
) -> StoreResult<ProjectEvent> {
let created_at = now_unix();
conn.execute(
"INSERT INTO project_events (project_id, kind_json, created_at)
VALUES (?1, ?2, ?3)",
params![
project.id.as_str(),
serde_json::to_string(kind)?,
created_at
],
)?;
let event_id = conn.last_insert_rowid();
if kind.is_wave_observable() {
insert_observation(
conn,
&ObservationRecipient::Wave {
wave_id: project.wave_id.clone(),
},
&ChildRef::Project(project.id.clone()),
event_id,
&ChildEventPayload::Project {
event: kind.clone(),
},
created_at,
)?;
}
Ok(ProjectEvent {
id: event_id,
project_id: project.id.clone(),
kind: kind.clone(),
created_at: crate::store::rows::unix_to_datetime(created_at),
})
}
fn insert_observation(
conn: &Connection,
recipient: &ObservationRecipient,
source: &ChildRef,
event_id: i64,
payload: &ChildEventPayload,
created_at: i64,
) -> StoreResult<()> {
let (recipient_kind, recipient_id) = recipient_columns(recipient);
let (source_kind, source_id) = child_columns(source);
conn.execute(
"INSERT OR IGNORE INTO observation_outbox (
recipient_kind, recipient_id, source_kind, source_id,
event_id, payload_json, created_at, delivered_at
) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, NULL)",
params![
recipient_kind,
recipient_id,
source_kind,
source_id,
event_id,
serde_json::to_string(payload)?,
created_at,
],
)?;
Ok(())
}
fn map_observation_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<ObservationOutboxRow> {
let recipient_kind: String = row.get(1)?;
let recipient_id: String = row.get(2)?;
let recipient = match recipient_kind.as_str() {
"wave" => ObservationRecipient::Wave {
wave_id: WaveId::parse(&recipient_id).map_err(|error| {
rusqlite::Error::FromSqlConversionFailure(
2,
rusqlite::types::Type::Text,
Box::new(error),
)
})?,
},
"project" => ObservationRecipient::Project {
project_id: ProjectId::from_raw(recipient_id),
},
value => {
return Err(invalid_column(
1,
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("unknown observation recipient {value:?}"),
),
))
}
};
let source_kind: String = row.get(3)?;
let source_id: String = row.get(4)?;
let source = match source_kind.as_str() {
"project" => ChildRef::Project(ProjectId::from_raw(source_id)),
"task" => ChildRef::Task(TaskId::from_raw(source_id)),
value => {
return Err(invalid_column(
3,
std::io::Error::new(
std::io::ErrorKind::InvalidData,
format!("unknown observation source {value:?}"),
),
))
}
};
let payload: ChildEventPayload = serde_json::from_str(&row.get::<_, String>(6)?)
.map_err(|error| invalid_column(6, error))?;
Ok(ObservationOutboxRow {
id: row.get(0)?,
recipient,
source,
event_id: row.get(5)?,
payload,
delivered_at: row
.get::<_, Option<i64>>(7)?
.map(crate::store::rows::unix_to_datetime),
})
}