Menu
AkurAI-Build
publicLatest change ed12ffe7602938c977a45314c87c888094c335ea - AKURAI-BUILD-5 stabilize paginated run history by Ólafur Búi Ólafsson
use std::{
env,
fs::{self, OpenOptions},
io::{Read, Write},
path::{Path, PathBuf},
sync::{Arc, Mutex, MutexGuard},
time::Duration,
};
use anyhow::{Context, Result, anyhow, bail, ensure};
use rusqlite::{
Connection, OpenFlags, TransactionBehavior, config::DbConfig, limits::Limit, params,
params_from_iter, types::Value,
};
use serde::{Deserialize, Serialize};
use zeroize::Zeroizing;
const APPLICATION_ID: i64 = 0x414B_5552;
const MIGRATIONS: &[(i64, &str, &str)] = &[
(
1,
"initial CI schema",
include_str!("../migrations/001_init.sql"),
),
(
2,
"repository visibility",
include_str!("../migrations/002_repository_visibility.sql"),
),
(
3,
"build workers",
include_str!("../migrations/003_workers.sql"),
),
];
const LEDGER: &str = "CREATE TABLE IF NOT EXISTS _migrations (
version INTEGER PRIMARY KEY,
name TEXT NOT NULL UNIQUE,
sql TEXT NOT NULL,
applied_at INTEGER NOT NULL DEFAULT (unixepoch())
) STRICT;";
#[derive(Clone)]
pub struct Database(Arc<Mutex<Connection>>);
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Repository {
pub id: i64,
pub name: String,
pub url: String,
pub default_branch: String,
pub visibility: String,
pub created_at: i64,
}
#[derive(Clone, Debug, Serialize)]
pub struct PublicRepository {
pub id: i64,
pub name: String,
pub url: String,
pub default_branch: String,
pub created_at: i64,
pub run_count: i64,
pub successful_runs: i64,
pub last_activity_at: i64,
}
#[derive(Clone, Debug)]
pub struct RunQuery {
pub repository: Option<String>,
pub statuses: Vec<String>,
pub git_ref: Option<String>,
pub triggers: Vec<String>,
pub search: Option<String>,
pub before_id: Option<i64>,
pub limit: usize,
pub offset: usize,
}
impl Default for RunQuery {
fn default() -> Self {
Self {
repository: None,
statuses: Vec::new(),
git_ref: None,
triggers: Vec::new(),
search: None,
before_id: None,
limit: 20,
offset: 0,
}
}
}
struct ValidatedRunQuery {
predicate: String,
values: Vec<Value>,
}
fn validate_run_query(query: &RunQuery) -> Result<ValidatedRunQuery> {
ensure!(
(1..=200).contains(&query.limit),
"run limit must be 1..=200"
);
ensure!(query.offset <= 10_000, "run offset must be <= 10000");
if let Some(before_id) = query.before_id {
ensure!(before_id > 0, "run cursor must be positive");
}
let repositories = query
.repository
.as_deref()
.map(|value| {
value
.split(',')
.map(str::trim)
.filter(|item| !item.is_empty())
.map(str::to_owned)
.collect::<Vec<_>>()
})
.unwrap_or_default();
ensure!(repositories.len() <= 32, "too many repository filters");
for repository in &repositories {
ensure!(
repository.len() <= 64
&& repository
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || b"._-".contains(&byte)),
"invalid repository filter: {repository}"
);
}
if let Some(git_ref) = &query.git_ref {
crate::config::validate_ref(git_ref)?;
}
for status in &query.statuses {
ensure!(
matches!(
status.as_str(),
"queued"
| "running"
| "waiting"
| "succeeded"
| "failed"
| "canceled"
| "interrupted"
),
"invalid run status: {status}"
);
}
for trigger in &query.triggers {
ensure!(
matches!(trigger.as_str(), "manual" | "webhook" | "retry"),
"invalid run trigger: {trigger}"
);
}
let mut clauses = Vec::new();
let mut values = Vec::new();
if !repositories.is_empty() {
let placeholders = std::iter::repeat_n("?", repositories.len())
.collect::<Vec<_>>()
.join(", ");
clauses.push(format!("p.name IN ({placeholders})"));
values.extend(repositories.into_iter().map(Value::Text));
}
if let Some(before_id) = query.before_id {
clauses.push("r.id <= ?".to_owned());
values.push(Value::Integer(before_id));
}
if let Some(git_ref) = &query.git_ref {
clauses.push("r.git_ref = ?".to_owned());
values.push(Value::Text(git_ref.clone()));
}
if !query.statuses.is_empty() {
let placeholders = std::iter::repeat_n("?", query.statuses.len())
.collect::<Vec<_>>()
.join(", ");
clauses.push(format!("r.status IN ({placeholders})"));
values.extend(query.statuses.iter().cloned().map(Value::Text));
}
if !query.triggers.is_empty() {
let placeholders = std::iter::repeat_n("?", query.triggers.len())
.collect::<Vec<_>>()
.join(", ");
clauses.push(format!("r.trigger IN ({placeholders})"));
values.extend(query.triggers.iter().cloned().map(Value::Text));
}
if let Some(search) = &query.search {
ensure!(
!search.is_empty() && search.len() <= 200,
"search must be 1..=200 bytes"
);
let needle = Value::Text(format!("%{search}%"));
clauses.push(
"(p.name LIKE ? OR r.git_ref LIKE ? OR COALESCE(r.commit_sha, '') LIKE ?
OR COALESCE(r.error, '') LIKE ?)"
.to_owned(),
);
values.extend(std::iter::repeat_n(needle, 4));
}
let predicate = if clauses.is_empty() {
" WHERE 1=1".to_owned()
} else {
format!(" WHERE {}", clauses.join(" AND "))
};
Ok(ValidatedRunQuery { predicate, values })
}
#[derive(Clone, Debug)]
pub struct RepositoryQuery {
pub search: Option<String>,
pub visibility: Option<String>,
pub limit: usize,
pub offset: usize,
}
impl Default for RepositoryQuery {
fn default() -> Self {
Self {
search: None,
visibility: None,
limit: 50,
offset: 0,
}
}
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Worker {
pub id: String,
pub host: String,
pub capabilities: String,
pub status: String,
pub capacity: i64,
pub current_run_id: Option<i64>,
pub started_at: i64,
pub heartbeat_at: i64,
pub completed_runs: i64,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Run {
pub id: i64,
pub repository_id: i64,
pub repository: String,
pub git_ref: String,
pub commit_sha: Option<String>,
pub trigger: String,
pub status: String,
pub error: Option<String>,
pub created_at: i64,
pub started_at: Option<i64>,
pub finished_at: Option<i64>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Job {
pub id: i64,
pub run_id: i64,
pub base_name: String,
pub name: String,
pub executor: String,
pub image: Option<String>,
pub platform: Option<String>,
pub status: String,
pub needs: Vec<String>,
pub spec_json: String,
pub logs: String,
pub exit_code: Option<i64>,
pub environment: Option<String>,
pub approval_required: bool,
pub approved_at: Option<i64>,
pub started_at: Option<i64>,
pub finished_at: Option<i64>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Artifact {
pub id: i64,
pub run_id: i64,
pub job_id: i64,
pub name: String,
pub path: String,
pub sha256: String,
pub bytes: i64,
pub created_at: i64,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct Deployment {
pub id: i64,
pub run_id: i64,
pub job_id: i64,
pub environment: String,
pub status: String,
pub artifacts_json: String,
pub created_at: i64,
pub finished_at: Option<i64>,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct RunDetail {
#[serde(flatten)]
pub run: Run,
pub jobs: Vec<Job>,
pub artifacts: Vec<Artifact>,
pub deployments: Vec<Deployment>,
}
impl Database {
pub fn open(path: &Path, key: &str, create: bool) -> Result<Self> {
validate_secret(key)?;
prepare_database(path, create)?;
let flags = if create {
OpenFlags::SQLITE_OPEN_READ_WRITE | OpenFlags::SQLITE_OPEN_CREATE
} else {
OpenFlags::SQLITE_OPEN_READ_WRITE
} | OpenFlags::SQLITE_OPEN_NO_MUTEX;
let connection = Connection::open_with_flags(path, flags)
.with_context(|| format!("open database {}", path.display()))?;
secure_database_file(path)?;
configure(&connection, key)?;
Ok(Self(Arc::new(Mutex::new(connection))))
}
#[cfg(test)]
pub fn memory(key: &str) -> Result<Self> {
validate_secret(key)?;
let connection = Connection::open_in_memory()?;
configure(&connection, key)?;
let database = Self(Arc::new(Mutex::new(connection)));
database.migrate()?;
Ok(database)
}
pub fn migrate(&self) -> Result<usize> {
let mut connection = self.connection()?;
let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
transaction.execute_batch(LEDGER)?;
let mut applied = 0;
for &(version, name, sql) in MIGRATIONS {
let existing: Option<(String, String)> = transaction
.query_row(
"SELECT name, sql FROM _migrations WHERE version = ?1",
[version],
|row| Ok((row.get(0)?, row.get(1)?)),
)
.optional()?;
if let Some((existing_name, existing_sql)) = existing {
ensure!(
existing_name == name && existing_sql == sql,
"migration history differs from this binary"
);
continue;
}
transaction.execute_batch(sql)?;
transaction.execute(
"INSERT INTO _migrations(version, name, sql) VALUES (?1, ?2, ?3)",
params![version, name, sql],
)?;
applied += 1;
}
transaction.commit()?;
drop(connection);
self.validate_schema()?;
Ok(applied)
}
pub fn validate_schema(&self) -> Result<()> {
let connection = self.connection()?;
let application_id: i64 =
connection.pragma_query_value(None, "application_id", |row| row.get(0))?;
let version: i64 = connection.pragma_query_value(None, "user_version", |row| row.get(0))?;
ensure!(
application_id == APPLICATION_ID && version == 3,
"database schema is not AkurAI Build v3"
);
for table in [
"repositories",
"runs",
"jobs",
"artifacts",
"deployments",
"workers",
"_migrations",
] {
let exists: bool = connection.query_row(
"SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE type='table' AND name=?1)",
[table],
|row| row.get(0),
)?;
ensure!(exists, "database is missing table {table}");
}
Ok(())
}
pub fn check_ready(&self) -> Result<()> {
self.connection()?.query_row("SELECT 1", [], |_| Ok(()))?;
Ok(())
}
pub fn add_repository(&self, name: &str, url: &str, branch: &str) -> Result<Repository> {
let connection = self.connection()?;
connection.execute(
"INSERT INTO repositories(name, url, default_branch) VALUES (?1, ?2, ?3)
ON CONFLICT(name) DO UPDATE SET url=excluded.url, default_branch=excluded.default_branch",
params![name, url, branch],
)?;
drop(connection);
self.repository(name)
}
pub fn set_repository_visibility(&self, name: &str, visibility: &str) -> Result<Repository> {
ensure!(
matches!(visibility, "private" | "public"),
"repository visibility must be private or public"
);
let connection = self.connection()?;
ensure!(
connection.execute(
"UPDATE repositories SET visibility=?1 WHERE name=?2",
params![visibility, name],
)? == 1,
"unknown repository {name}"
);
drop(connection);
self.repository(name)
}
pub fn rename_repository(&self, old: &str, new: &str) -> Result<Repository> {
crate::config::validate_repo_name(new)?;
let connection = self.connection()?;
ensure!(
connection.execute(
"UPDATE repositories SET name=?1 WHERE name=?2",
params![new, old],
)? == 1,
"unknown repository {old}"
);
drop(connection);
self.repository(new)
}
/// Unregister a repository, cascading to its runs, jobs, logs, and artifacts.
/// Returns the removed repository record. Errors if the name is unknown.
pub fn remove_repository(&self, name: &str) -> Result<Repository> {
let repository = self.repository(name)?;
let connection = self.connection()?;
ensure!(
connection.execute("DELETE FROM repositories WHERE name=?1", [name])? == 1,
"unknown repository {name}"
);
Ok(repository)
}
pub fn repository(&self, name: &str) -> Result<Repository> {
self.connection()?
.query_row(
"SELECT id, name, url, default_branch, visibility, created_at FROM repositories WHERE name=?1",
[name],
map_repository,
)
.with_context(|| format!("unknown repository {name}"))
}
pub fn repositories(&self) -> Result<Vec<Repository>> {
let connection = self.connection()?;
let mut statement = connection.prepare(
"SELECT id, name, url, default_branch, visibility, created_at FROM repositories ORDER BY name",
)?;
Ok(statement
.query_map([], map_repository)?
.collect::<rusqlite::Result<_>>()?)
}
pub fn query_repositories(&self, query: &RepositoryQuery) -> Result<Vec<Repository>> {
ensure!(
(1..=500).contains(&query.limit),
"repository limit must be 1..=500"
);
ensure!(query.offset <= 10_000, "repository offset must be <= 10000");
if let Some(visibility) = &query.visibility {
ensure!(
matches!(visibility.as_str(), "private" | "public"),
"repository visibility must be private or public"
);
}
if let Some(search) = &query.search {
ensure!(
!search.is_empty() && search.len() <= 200,
"search must be 1..=200 bytes"
);
}
let search = query.search.as_ref().map(|value| format!("%{value}%"));
let connection = self.connection()?;
let mut statement = connection.prepare(
"SELECT id, name, url, default_branch, visibility, created_at
FROM repositories
WHERE (?1 IS NULL OR visibility=?1)
AND (?2 IS NULL OR name LIKE ?2 OR url LIKE ?2 OR default_branch LIKE ?2)
ORDER BY name LIMIT ?3 OFFSET ?4",
)?;
Ok(statement
.query_map(
params![
query.visibility.as_deref(),
search.as_deref(),
i64::try_from(query.limit)?,
i64::try_from(query.offset)?
],
map_repository,
)?
.collect::<rusqlite::Result<_>>()?)
}
pub fn public_repositories(&self) -> Result<Vec<PublicRepository>> {
let connection = self.connection()?;
let mut statement = connection.prepare(
"SELECT p.id, p.name, p.url, p.default_branch, p.created_at,
COUNT(r.id),
COALESCE(SUM(CASE WHEN r.status='succeeded' THEN 1 ELSE 0 END), 0),
COALESCE(MAX(COALESCE(r.finished_at, r.started_at, r.created_at)), p.created_at)
FROM repositories p
LEFT JOIN runs r ON r.repository_id=p.id
WHERE p.visibility='public'
GROUP BY p.id
ORDER BY COALESCE(MAX(COALESCE(r.finished_at, r.started_at, r.created_at)), p.created_at) DESC,
COUNT(r.id) DESC,
p.name
LIMIT 12",
)?;
Ok(statement
.query_map([], map_public_repository)?
.collect::<rusqlite::Result<_>>()?)
}
pub fn create_run(
&self,
repository_id: i64,
git_ref: &str,
commit: Option<&str>,
trigger: &str,
) -> Result<i64> {
let connection = self.connection()?;
connection.execute(
"INSERT INTO runs(repository_id, git_ref, commit_sha, trigger, status) VALUES (?1, ?2, ?3, ?4, 'queued')",
params![repository_id, git_ref, commit, trigger],
)?;
Ok(connection.last_insert_rowid())
}
pub fn run(&self, id: i64) -> Result<Run> {
self.connection()?
.query_row(
"SELECT r.id, r.repository_id, p.name, r.git_ref, r.commit_sha, r.trigger, r.status,
r.error, r.created_at, r.started_at, r.finished_at
FROM runs r JOIN repositories p ON p.id=r.repository_id WHERE r.id=?1",
[id],
map_run,
)
.with_context(|| format!("unknown run {id}"))
}
pub fn runs(&self, repository: Option<&str>, limit: usize) -> Result<Vec<Run>> {
self.query_runs(&RunQuery {
repository: repository.map(str::to_owned),
limit,
..RunQuery::default()
})
}
pub fn query_runs(&self, query: &RunQuery) -> Result<Vec<Run>> {
let ValidatedRunQuery {
predicate,
mut values,
} = validate_run_query(query)?;
let mut sql = String::from(
"SELECT r.id, r.repository_id, p.name, r.git_ref, r.commit_sha, r.trigger, r.status,
r.error, r.created_at, r.started_at, r.finished_at
FROM runs r JOIN repositories p ON p.id=r.repository_id",
);
sql.push_str(&predicate);
sql.push_str(" ORDER BY r.id DESC LIMIT ? OFFSET ?");
values.push(Value::Integer(i64::try_from(query.limit)?));
values.push(Value::Integer(i64::try_from(query.offset)?));
let connection = self.connection()?;
let mut statement = connection.prepare(&sql)?;
Ok(statement
.query_map(params_from_iter(values), map_run)?
.collect::<rusqlite::Result<_>>()?)
}
pub fn count_runs(&self, query: &RunQuery) -> Result<i64> {
let ValidatedRunQuery { predicate, values } = validate_run_query(query)?;
let sql = format!(
"SELECT COUNT(*) FROM runs r JOIN repositories p ON p.id=r.repository_id{predicate}"
);
Ok(self
.connection()?
.query_row(&sql, params_from_iter(values), |row| row.get(0))?)
}
pub fn max_run_id(&self, query: &RunQuery) -> Result<Option<i64>> {
let ValidatedRunQuery { predicate, values } = validate_run_query(query)?;
let sql = format!(
"SELECT MAX(r.id) FROM runs r JOIN repositories p ON p.id=r.repository_id{predicate}"
);
Ok(self
.connection()?
.query_row(&sql, params_from_iter(values), |row| row.get(0))?)
}
pub fn detail(&self, id: i64) -> Result<RunDetail> {
Ok(RunDetail {
run: self.run(id)?,
jobs: self.jobs(id)?,
artifacts: self.artifacts(id)?,
deployments: self.deployments(id)?,
})
}
pub fn claim_next_queued_run(&self) -> Result<Option<i64>> {
let mut connection = self.connection()?;
let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
let id = transaction
.query_row(
"SELECT id FROM runs WHERE status='queued' ORDER BY id LIMIT 1",
[],
|row| row.get(0),
)
.optional()?;
if let Some(id) = id {
ensure!(
transaction.execute(
"UPDATE runs SET status='running',
started_at=COALESCE(started_at, unixepoch()), error=NULL
WHERE id=?1 AND status='queued'",
[id],
)? == 1,
"queued run changed during claim"
);
}
transaction.commit()?;
Ok(id)
}
pub fn claim_run(&self, id: i64) -> Result<bool> {
Ok(self.connection()?.execute(
"UPDATE runs SET status='running', started_at=COALESCE(started_at, unixepoch()), error=NULL WHERE id=?1 AND status='queued'",
[id],
)? == 1)
}
pub fn set_run_commit(&self, id: i64, commit: &str) -> Result<()> {
self.connection()?.execute(
"UPDATE runs SET commit_sha=?2 WHERE id=?1",
params![id, commit],
)?;
Ok(())
}
pub fn finish_run(&self, id: i64, status: &str, error: Option<&str>) -> Result<()> {
self.connection()?.execute(
"UPDATE runs SET status=?2, error=?3, finished_at=CASE WHEN ?2 IN ('succeeded','failed','canceled','interrupted') THEN unixepoch() ELSE NULL END WHERE id=?1",
params![id, status, error],
)?;
Ok(())
}
pub fn insert_jobs(&self, run_id: i64, specs: &[crate::config::JobSpec]) -> Result<()> {
let mut connection = self.connection()?;
let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
for spec in specs {
let spec_json = serde_json::to_string(spec)?;
let needs = serde_json::to_string(&spec.needs)?;
let platform = spec.matrix.get("platform");
transaction.execute(
"INSERT INTO jobs(run_id, base_name, name, executor, image, platform, status, needs_json, spec_json, environment, approval_required)
VALUES (?1, ?2, ?3, ?4, ?5, ?6, 'queued', ?7, ?8, ?9, ?10)",
params![run_id, spec.base_name, spec.name, spec.executor, spec.image, platform, needs, spec_json, spec.environment, spec.approval],
)?;
}
transaction.commit()?;
Ok(())
}
pub fn jobs(&self, run_id: i64) -> Result<Vec<Job>> {
let connection = self.connection()?;
let mut statement = connection.prepare(
"SELECT id, run_id, base_name, name, executor, image, platform, status, needs_json,
spec_json, logs, exit_code, environment, approval_required, approved_at, started_at, finished_at
FROM jobs WHERE run_id=?1 ORDER BY id",
)?;
Ok(statement
.query_map([run_id], map_job)?
.collect::<rusqlite::Result<_>>()?)
}
pub fn set_job_status(
&self,
id: i64,
status: &str,
logs: Option<&str>,
exit_code: Option<i32>,
) -> Result<()> {
self.connection()?.execute(
"UPDATE jobs SET status=?2, logs=COALESCE(?3, logs), exit_code=?4,
started_at=CASE WHEN ?2='running' THEN COALESCE(started_at, unixepoch()) ELSE started_at END,
finished_at=CASE WHEN ?2 IN ('succeeded','failed','skipped','canceled','interrupted') THEN unixepoch() ELSE finished_at END
WHERE id=?1",
params![id, status, logs, exit_code],
)?;
Ok(())
}
pub fn approve_environment(&self, run_id: i64, environment: &str) -> Result<usize> {
let connection = self.connection()?;
let changed = connection.execute(
"UPDATE jobs SET approved_at=unixepoch(), status=CASE WHEN status='waiting' THEN 'queued' ELSE status END
WHERE run_id=?1 AND environment=?2 AND approval_required=1 AND status IN ('waiting','queued')",
params![run_id, environment],
)?;
ensure!(
changed > 0,
"no promotable {environment} job in run {run_id}"
);
connection.execute(
"UPDATE runs SET status='queued', finished_at=NULL WHERE id=?1 AND status='waiting'",
[run_id],
)?;
Ok(changed)
}
/// Cancel a run that has not started executing (status queued or
/// waiting), cascading to its not-yet-final jobs. Running and terminal
/// runs are refused: the worker owns a running run's lifecycle.
pub fn cancel_run(&self, id: i64) -> Result<()> {
let mut connection = self.connection()?;
let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
let changed = transaction.execute(
"UPDATE runs SET status='canceled', error=NULL, finished_at=unixepoch()
WHERE id=?1 AND status IN ('queued','waiting')",
[id],
)?;
if changed != 1 {
let status: Option<String> = transaction
.query_row("SELECT status FROM runs WHERE id=?1", [id], |row| {
row.get(0)
})
.optional()?;
transaction.commit()?;
match status {
Some(status) => {
bail!("run {id} is {status}; only queued or waiting runs can be canceled")
}
None => bail!("unknown run {id}"),
}
}
transaction.execute(
"UPDATE jobs SET status='canceled', finished_at=unixepoch()
WHERE run_id=?1 AND status IN ('queued','waiting','running')",
[id],
)?;
transaction.commit()?;
Ok(())
}
pub fn retry(&self, run_id: i64) -> Result<i64> {
let run = self.run(run_id)?;
self.create_run(
run.repository_id,
&run.git_ref,
run.commit_sha.as_deref(),
"retry",
)
}
pub fn add_artifact(
&self,
run_id: i64,
job_id: i64,
name: &str,
path: &str,
sha256: &str,
bytes: u64,
) -> Result<i64> {
let connection = self.connection()?;
connection.execute(
"INSERT INTO artifacts(run_id, job_id, name, path, sha256, bytes) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
params![run_id, job_id, name, path, sha256, i64::try_from(bytes)?],
)?;
Ok(connection.last_insert_rowid())
}
pub fn artifact(&self, id: i64) -> Result<Artifact> {
self.connection()?.query_row(
"SELECT id, run_id, job_id, name, path, sha256, bytes, created_at FROM artifacts WHERE id=?1",
[id],
map_artifact,
).with_context(|| format!("unknown artifact {id}"))
}
pub fn artifacts(&self, run_id: i64) -> Result<Vec<Artifact>> {
let connection = self.connection()?;
let mut statement = connection.prepare(
"SELECT id, run_id, job_id, name, path, sha256, bytes, created_at FROM artifacts WHERE run_id=?1 ORDER BY id",
)?;
Ok(statement
.query_map([run_id], map_artifact)?
.collect::<rusqlite::Result<_>>()?)
}
pub fn begin_deployment(
&self,
run_id: i64,
job_id: i64,
environment: &str,
artifacts_json: &str,
) -> Result<i64> {
let connection = self.connection()?;
connection.execute(
"INSERT INTO deployments(run_id, job_id, environment, status, artifacts_json) VALUES (?1, ?2, ?3, 'running', ?4)",
params![run_id, job_id, environment, artifacts_json],
)?;
Ok(connection.last_insert_rowid())
}
pub fn finish_deployment(&self, id: i64, status: &str) -> Result<()> {
self.connection()?.execute(
"UPDATE deployments SET status=?2, finished_at=unixepoch() WHERE id=?1",
params![id, status],
)?;
Ok(())
}
pub fn deployments(&self, run_id: i64) -> Result<Vec<Deployment>> {
let connection = self.connection()?;
let mut statement = connection.prepare(
"SELECT id, run_id, job_id, environment, status, artifacts_json, created_at, finished_at
FROM deployments WHERE run_id=?1 ORDER BY id",
)?;
Ok(statement
.query_map([run_id], map_deployment)?
.collect::<rusqlite::Result<_>>()?)
}
pub fn queued_runs(&self) -> Result<i64> {
Ok(self.connection()?.query_row(
"SELECT COUNT(*) FROM runs WHERE status='queued'",
[],
|row| row.get(0),
)?)
}
pub fn reset_workers(&self) -> Result<()> {
self.connection()?.execute(
"UPDATE workers SET status='offline', current_run_id=NULL, heartbeat_at=unixepoch()",
[],
)?;
Ok(())
}
pub fn register_worker(
&self,
id: &str,
host: &str,
capabilities: &str,
capacity: usize,
) -> Result<()> {
ensure!(!id.is_empty() && id.len() <= 64, "invalid worker id");
ensure!(!host.is_empty() && host.len() <= 128, "invalid worker host");
ensure!(
!capabilities.is_empty() && capabilities.len() <= 256,
"invalid capabilities"
);
ensure!(
(1..=16).contains(&capacity),
"worker capacity must be 1..=16"
);
self.connection()?.execute(
"INSERT INTO workers(id, host, capabilities, status, capacity, started_at, heartbeat_at)
VALUES (?1, ?2, ?3, 'starting', ?4, unixepoch(), unixepoch())
ON CONFLICT(id) DO UPDATE SET host=excluded.host, capabilities=excluded.capabilities,
status='starting', capacity=excluded.capacity, current_run_id=NULL,
started_at=excluded.started_at, heartbeat_at=excluded.heartbeat_at, completed_runs=0",
params![id, host, capabilities, i64::try_from(capacity)?],
)?;
Ok(())
}
pub fn heartbeat_worker(
&self,
id: &str,
status: &str,
current_run_id: Option<i64>,
completed: bool,
) -> Result<()> {
ensure!(
matches!(status, "starting" | "idle" | "running" | "offline"),
"invalid worker status"
);
ensure!(
self.connection()?.execute(
"UPDATE workers SET status=?2, current_run_id=?3, heartbeat_at=unixepoch(),
completed_runs=completed_runs + ?4 WHERE id=?1",
params![id, status, current_run_id, i64::from(completed)],
)? == 1,
"unknown worker {id}"
);
Ok(())
}
pub fn workers(&self) -> Result<Vec<Worker>> {
let connection = self.connection()?;
let mut statement = connection.prepare(
"SELECT id, host, capabilities, status, capacity, current_run_id, started_at,
heartbeat_at, completed_runs
FROM workers ORDER BY id",
)?;
Ok(statement
.query_map([], map_worker)?
.collect::<rusqlite::Result<_>>()?)
}
pub fn recover_interrupted(&self) -> Result<usize> {
let connection = self.connection()?;
let deployments = connection.execute(
"UPDATE deployments SET status='failed', finished_at=unixepoch() WHERE status='running'",
[],
)?;
let jobs = connection.execute(
"UPDATE jobs SET status='interrupted', finished_at=unixepoch() WHERE status='running'",
[],
)?;
let runs = connection.execute(
"UPDATE runs SET status='interrupted', error='controller restarted during execution', finished_at=unixepoch() WHERE status='running'",
[],
)?;
Ok(deployments + jobs + runs)
}
fn connection(&self) -> Result<MutexGuard<'_, Connection>> {
self.0
.lock()
.map_err(|_| anyhow!("database mutex poisoned"))
}
}
fn map_repository(row: &rusqlite::Row<'_>) -> rusqlite::Result<Repository> {
Ok(Repository {
id: row.get(0)?,
name: row.get(1)?,
url: row.get(2)?,
default_branch: row.get(3)?,
visibility: row.get(4)?,
created_at: row.get(5)?,
})
}
fn map_public_repository(row: &rusqlite::Row<'_>) -> rusqlite::Result<PublicRepository> {
Ok(PublicRepository {
id: row.get(0)?,
name: row.get(1)?,
url: row.get(2)?,
default_branch: row.get(3)?,
created_at: row.get(4)?,
run_count: row.get(5)?,
successful_runs: row.get(6)?,
last_activity_at: row.get(7)?,
})
}
fn map_run(row: &rusqlite::Row<'_>) -> rusqlite::Result<Run> {
Ok(Run {
id: row.get(0)?,
repository_id: row.get(1)?,
repository: row.get(2)?,
git_ref: row.get(3)?,
commit_sha: row.get(4)?,
trigger: row.get(5)?,
status: row.get(6)?,
error: row.get(7)?,
created_at: row.get(8)?,
started_at: row.get(9)?,
finished_at: row.get(10)?,
})
}
fn map_job(row: &rusqlite::Row<'_>) -> rusqlite::Result<Job> {
let needs_json: String = row.get(8)?;
let needs = serde_json::from_str(&needs_json).map_err(|error| {
rusqlite::Error::FromSqlConversionFailure(8, rusqlite::types::Type::Text, Box::new(error))
})?;
Ok(Job {
id: row.get(0)?,
run_id: row.get(1)?,
base_name: row.get(2)?,
name: row.get(3)?,
executor: row.get(4)?,
image: row.get(5)?,
platform: row.get(6)?,
status: row.get(7)?,
needs,
spec_json: row.get(9)?,
logs: row.get(10)?,
exit_code: row.get(11)?,
environment: row.get(12)?,
approval_required: row.get(13)?,
approved_at: row.get(14)?,
started_at: row.get(15)?,
finished_at: row.get(16)?,
})
}
fn map_artifact(row: &rusqlite::Row<'_>) -> rusqlite::Result<Artifact> {
Ok(Artifact {
id: row.get(0)?,
run_id: row.get(1)?,
job_id: row.get(2)?,
name: row.get(3)?,
path: row.get(4)?,
sha256: row.get(5)?,
bytes: row.get(6)?,
created_at: row.get(7)?,
})
}
fn map_deployment(row: &rusqlite::Row<'_>) -> rusqlite::Result<Deployment> {
Ok(Deployment {
id: row.get(0)?,
run_id: row.get(1)?,
job_id: row.get(2)?,
environment: row.get(3)?,
status: row.get(4)?,
artifacts_json: row.get(5)?,
created_at: row.get(6)?,
finished_at: row.get(7)?,
})
}
fn map_worker(row: &rusqlite::Row<'_>) -> rusqlite::Result<Worker> {
Ok(Worker {
id: row.get(0)?,
host: row.get(1)?,
capabilities: row.get(2)?,
status: row.get(3)?,
capacity: row.get(4)?,
current_run_id: row.get(5)?,
started_at: row.get(6)?,
heartbeat_at: row.get(7)?,
completed_runs: row.get(8)?,
})
}
fn prepare_database(path: &Path, create: bool) -> Result<()> {
let parent = path.parent().context("database path has no parent")?;
if create {
fs::create_dir_all(parent)?;
}
ensure!(
parent.is_dir(),
"database parent does not exist: {}",
parent.display()
);
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
fs::set_permissions(parent, fs::Permissions::from_mode(0o700))?;
}
match fs::symlink_metadata(path) {
Ok(metadata) => ensure!(
metadata.is_file() && !metadata.file_type().is_symlink(),
"database must be a regular non-symlink file"
),
Err(error) if error.kind() == std::io::ErrorKind::NotFound && create => {}
Err(error) => return Err(error.into()),
}
Ok(())
}
fn secure_database_file(path: &Path) -> Result<()> {
let metadata = fs::symlink_metadata(path)?;
ensure!(
metadata.is_file() && !metadata.file_type().is_symlink(),
"database must be a regular non-symlink file"
);
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
fs::set_permissions(path, fs::Permissions::from_mode(0o600))?;
}
Ok(())
}
pub fn load_secret(path: &Path, environment: &str) -> Result<Zeroizing<String>> {
let mut value = Zeroizing::new(String::new());
match fs::symlink_metadata(path) {
Ok(metadata) => {
ensure!(
metadata.is_file() && !metadata.file_type().is_symlink(),
"secret must be a regular non-symlink file"
);
#[cfg(unix)]
{
use std::os::unix::fs::MetadataExt;
ensure!(
metadata.mode() & 0o077 == 0,
"secret file permissions must be 0600"
);
}
ensure!(metadata.len() <= 1024, "secret file is too large");
let mut file = crate::fsguard::open_nofollow(path, false)?;
file.read_to_string(&mut value)?;
crate::fsguard::ensure_path_matches_file(path, &file, "secret")?;
let trimmed = value.trim_end().len();
value.truncate(trimmed);
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
value.push_str(
&env::var(environment)
.with_context(|| format!("missing {environment} or {}", path.display()))?,
);
}
Err(error) => return Err(error.into()),
}
validate_secret(&value)?;
Ok(value)
}
pub fn generate_secret(path: &Path) -> Result<()> {
ensure!(!path.exists(), "refusing to overwrite {}", path.display());
let parent = path.parent().context("secret path has no parent")?;
fs::create_dir_all(parent)?;
#[cfg(unix)]
{
use std::os::unix::fs::{OpenOptionsExt, PermissionsExt};
fs::set_permissions(parent, fs::Permissions::from_mode(0o700))?;
let mut file = OpenOptions::new()
.write(true)
.create_new(true)
.mode(0o600)
.open(path)?;
write_secret(&mut file)?;
}
#[cfg(not(unix))]
{
let mut file = OpenOptions::new().write(true).create_new(true).open(path)?;
write_secret(&mut file)?;
}
Ok(())
}
fn write_secret(file: &mut fs::File) -> Result<()> {
let mut bytes = Zeroizing::new([0_u8; 32]);
getrandom::fill(bytes.as_mut()).context("read operating-system randomness")?;
let mut encoded = Zeroizing::new(String::with_capacity(65));
for byte in bytes.iter() {
use std::fmt::Write as _;
write!(encoded, "{byte:02x}")?;
}
encoded.push('\n');
file.write_all(encoded.as_bytes())?;
file.sync_all()?;
Ok(())
}
fn validate_secret(value: &str) -> Result<()> {
ensure!(
value.len() == 64 && value.bytes().all(|byte| byte.is_ascii_hexdigit()),
"secret must be 32 bytes encoded as hexadecimal"
);
Ok(())
}
fn decode_secret(value: &str) -> Result<Zeroizing<Vec<u8>>> {
validate_secret(value)?;
let mut bytes = Zeroizing::new(Vec::with_capacity(32));
for pair in value.as_bytes().chunks_exact(2) {
let text = std::str::from_utf8(pair)?;
bytes.push(u8::from_str_radix(text, 16)?);
}
Ok(bytes)
}
fn configure(connection: &Connection, key: &str) -> Result<()> {
apply_key(connection, key)?;
let cipher: String = connection.pragma_query_value(None, "cipher_version", |row| row.get(0))?;
ensure!(!cipher.is_empty(), "SQLCipher support is unavailable");
connection
.query_row("SELECT count(*) FROM sqlite_schema", [], |_| Ok(()))
.context("database key is incorrect or database is corrupt")?;
connection.busy_timeout(Duration::from_secs(5))?;
connection.pragma_update(None, "foreign_keys", "ON")?;
connection.pragma_update(None, "journal_mode", "WAL")?;
connection.pragma_update(None, "synchronous", "FULL")?;
connection.pragma_update(None, "secure_delete", "ON")?;
connection.set_db_config(DbConfig::SQLITE_DBCONFIG_DEFENSIVE, true)?;
connection.set_db_config(DbConfig::SQLITE_DBCONFIG_TRUSTED_SCHEMA, false)?;
connection.set_limit(Limit::SQLITE_LIMIT_LENGTH, 16 * 1024 * 1024)?;
connection.set_limit(Limit::SQLITE_LIMIT_SQL_LENGTH, 256 * 1024)?;
connection.set_limit(Limit::SQLITE_LIMIT_VARIABLE_NUMBER, 128)?;
Ok(())
}
#[allow(unsafe_code)]
fn apply_key(connection: &Connection, encoded: &str) -> Result<()> {
let key = decode_secret(encoded)?;
let length = i32::try_from(key.len())?;
// SAFETY: SQLCipher copies these bytes during the call; rusqlite owns the live handle.
let result =
unsafe { rusqlite::ffi::sqlite3_key(connection.handle(), key.as_ptr().cast(), length) };
ensure!(
result == rusqlite::ffi::SQLITE_OK,
"SQLCipher rejected the database key"
);
Ok(())
}
trait OptionalRow<T> {
fn optional(self) -> rusqlite::Result<Option<T>>;
}
impl<T> OptionalRow<T> for rusqlite::Result<T> {
fn optional(self) -> rusqlite::Result<Option<T>> {
match self {
Ok(value) => Ok(Some(value)),
Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
Err(error) => Err(error),
}
}
}
pub fn default_paths(root: &Path) -> (PathBuf, PathBuf, PathBuf) {
(
root.join("akurai.db"),
root.join("workspaces"),
root.join("artifacts"),
)
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::BTreeMap;
const KEY: &str = "000102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
const KEY2: &str = "100102030405060708090a0b0c0d0e0f101112131415161718191a1b1c1d1e1f";
fn open_db(dir: &tempfile::TempDir) -> Result<Database> {
let path = dir.path().join("akurai.db");
let db = Database::open(&path, KEY, true)?;
db.migrate()?;
Ok(db)
}
// ── migrate idempotent ──
#[test]
fn migrate_is_idempotent() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let count = db.migrate()?;
assert_eq!(count, 0, "second migrate should apply zero migrations");
db.validate_schema()?;
// schema still usable
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
assert_eq!(repo.name, "app");
Ok(())
}
// ── wrong key ──
#[test]
fn open_with_wrong_key_fails() -> Result<()> {
let dir = tempfile::tempdir()?;
let path = dir.path().join("akurai.db");
// create with correct key
let _db = Database::open(&path, KEY, true)?;
// reopen with wrong key
match Database::open(&path, KEY2, false) {
Ok(_) => panic!("opening with wrong key should have failed"),
Err(err) => {
let msg = format!("{err:#}");
assert!(
msg.contains("database key is incorrect") || msg.contains("corrupt"),
"expected key-error message, got: {msg}"
);
}
}
Ok(())
}
// ── secret generate → load round-trip ──
#[test]
fn secret_generate_then_load_round_trips() -> Result<()> {
let dir = tempfile::tempdir()?;
let path = dir.path().join("secret.key");
generate_secret(&path)?;
let loaded = load_secret(&path, "TEST_SECRET_ENV")?;
// load_secret returns the raw hex string; validate it
validate_secret(&loaded)?;
assert_eq!(loaded.len(), 64);
Ok(())
}
#[test]
fn load_secret_missing_file_and_env_errors() -> Result<()> {
let dir = tempfile::tempdir()?;
let path = dir.path().join("nonexistent.key");
let err = load_secret(&path, "NONEXISTENT_ENV_VAR_AKBUILD_TEST")
.expect_err("loading nonexistent secret should fail");
let msg = format!("{err:#}");
assert!(
msg.contains("missing NONEXISTENT_ENV_VAR_AKBUILD_TEST")
|| msg.contains(&path.display().to_string()),
"expected missing-env-or-file message, got: {msg}"
);
Ok(())
}
#[test]
fn generate_secret_refuses_overwrite() -> Result<()> {
let dir = tempfile::tempdir()?;
let path = dir.path().join("secret.key");
generate_secret(&path)?;
let err = generate_secret(&path).expect_err("generating over existing secret should fail");
let msg = format!("{err:#}");
assert!(msg.contains("refusing to overwrite"), "got: {msg}");
Ok(())
}
// ── repositories ──
#[test]
fn repository_add_then_list() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.add_repository("app", "https://example.com/app.git", "main")?;
db.add_repository("lib", "https://example.com/lib.git", "develop")?;
let repos = db.repositories()?;
assert_eq!(repos.len(), 2);
// sorted by name
assert_eq!(repos[0].name, "app");
assert_eq!(repos[1].name, "lib");
Ok(())
}
#[test]
fn repository_add_upserts_on_duplicate_name() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.add_repository("app", "https://example.com/app.git", "main")?;
// same name, different url — upserts
let updated = db.add_repository("app", "https://example.com/app2.git", "develop")?;
assert_eq!(updated.url, "https://example.com/app2.git");
assert_eq!(updated.default_branch, "develop");
let repos = db.repositories()?;
assert_eq!(repos.len(), 1, "upsert should not increase count");
Ok(())
}
#[test]
fn repository_rename_works() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.add_repository("old", "https://example.com/old.git", "main")?;
let renamed = db.rename_repository("old", "new")?;
assert_eq!(renamed.name, "new");
// old name gone
let err = db
.repository("old")
.expect_err("old name should not exist after rename");
let msg = format!("{err:#}");
assert!(msg.contains("unknown repository"), "got: {msg}");
// new name works
let found = db.repository("new")?;
assert_eq!(found.url, "https://example.com/old.git");
Ok(())
}
#[test]
fn repository_rename_to_existing_name_fails() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.add_repository("alpha", "https://example.com/alpha.git", "main")?;
db.add_repository("beta", "https://example.com/beta.git", "main")?;
let err = db
.rename_repository("alpha", "beta")
.expect_err("renaming to existing name should fail");
let msg = format!("{err:#}");
assert!(
msg.contains("UNIQUE constraint") || msg.contains("constraint"),
"expected constraint violation, got: {msg}"
);
Ok(())
}
#[test]
fn repository_rename_invalid_name_fails() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.add_repository("valid", "https://example.com/valid.git", "main")?;
// validate_repo_name rejects empty and names with bad chars
let err = db
.rename_repository("valid", "")
.expect_err("empty name should be rejected");
let msg = format!("{err:#}");
assert!(msg.contains("invalid repository name"), "got: {msg}");
Ok(())
}
#[test]
fn repository_unknown_name_errors() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let err = db
.repository("nonexistent")
.expect_err("unknown repository should error");
let msg = format!("{err:#}");
assert!(msg.contains("unknown repository"), "got: {msg}");
Ok(())
}
#[test]
fn remove_repository_deletes_and_returns_it() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.add_repository("gone", "https://example.com/gone.git", "main")?;
let removed = db.remove_repository("gone")?;
assert_eq!(removed.name, "gone");
let err = db
.repository("gone")
.expect_err("removed repository should be gone");
assert!(format!("{err:#}").contains("unknown repository"));
Ok(())
}
#[test]
fn remove_repository_unknown_errors() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let err = db
.remove_repository("ghost")
.expect_err("removing unknown repository should error");
assert!(format!("{err:#}").contains("unknown repository"));
Ok(())
}
#[test]
fn remove_repository_cascades_runs() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("casc", "https://example.com/casc.git", "main")?;
let run_id = db.create_run(
repo.id,
"main",
Some("deadbeefdeadbeefdeadbeefdeadbeefdeadbeef"),
"manual",
)?;
assert!(db.run(run_id).is_ok(), "run should exist before removal");
db.remove_repository("casc")?;
assert!(
db.run(run_id).is_err(),
"run row must cascade-delete when its repository is removed"
);
Ok(())
}
#[test]
fn repository_query_by_visibility_and_search() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.add_repository("alpha", "https://gh.com/alpha.git", "main")?;
db.add_repository("beta", "https://gh.com/beta.git", "main")?;
db.set_repository_visibility("beta", "public")?;
let public = db.query_repositories(&RepositoryQuery {
visibility: Some("public".into()),
..RepositoryQuery::default()
})?;
assert_eq!(public.len(), 1);
assert_eq!(public[0].name, "beta");
let by_name = db.query_repositories(&RepositoryQuery {
search: Some("alpha".into()),
..RepositoryQuery::default()
})?;
assert_eq!(by_name.len(), 1);
assert_eq!(by_name[0].name, "alpha");
Ok(())
}
// ── runs ──
#[test]
fn run_insert_then_query_by_repo() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let alpha = db.add_repository("alpha", "https://example.com/alpha.git", "main")?;
let beta = db.add_repository("beta", "https://example.com/beta.git", "main")?;
let r1 = db.create_run(alpha.id, "main", None, "manual")?;
let _r2 = db.create_run(beta.id, "release", None, "webhook")?;
let r3 = db.create_run(alpha.id, "feature/x", None, "manual")?;
// by repo filter
let alpha_runs = db.query_runs(&RunQuery {
repository: Some("alpha".into()),
limit: 20,
..RunQuery::default()
})?;
assert_eq!(alpha_runs.len(), 2);
// ordered by id DESC
assert_eq!(alpha_runs[0].id, r3);
assert_eq!(alpha_runs[1].id, r1);
// by status filter
db.finish_run(r1, "succeeded", None)?;
let succeeded = db.query_runs(&RunQuery {
statuses: vec!["succeeded".into()],
limit: 20,
..RunQuery::default()
})?;
assert_eq!(succeeded.len(), 1);
assert_eq!(succeeded[0].id, r1);
Ok(())
}
#[test]
fn run_query_limit_respected() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
for i in 0..10 {
db.create_run(repo.id, &format!("ref{i}"), None, "manual")?;
}
let limited = db.query_runs(&RunQuery {
limit: 3,
..RunQuery::default()
})?;
assert_eq!(limited.len(), 3);
// first 3 by DESC id
let ids: Vec<i64> = limited.iter().map(|r| r.id).collect();
assert!(ids[0] > ids[1] && ids[1] > ids[2], "not DESC ordered");
let oldest = db.query_runs(&RunQuery {
limit: 3,
offset: 9,
..RunQuery::default()
})?;
assert_eq!(oldest.len(), 1);
assert_eq!(db.count_runs(&RunQuery::default())?, 10);
let cursor = db.max_run_id(&RunQuery::default())?.expect("run cursor");
db.create_run(repo.id, "new-after-snapshot", None, "manual")?;
let snapshot = RunQuery {
before_id: Some(cursor),
..RunQuery::default()
};
assert_eq!(db.count_runs(&snapshot)?, 10);
assert!(db.query_runs(&snapshot)?.iter().all(|run| run.id <= cursor));
assert_eq!(db.count_runs(&RunQuery::default())?, 11);
Ok(())
}
#[test]
fn run_claim_optimistic_concurrency() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
// first claim succeeds
let claimed = db.claim_run(run_id)?;
assert!(claimed, "first claim should succeed");
assert_eq!(db.run(run_id)?.status, "running");
// second claim fails (already claimed)
let again = db.claim_run(run_id)?;
assert!(!again, "second claim of same run should fail");
Ok(())
}
#[test]
fn run_finish_transitions() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
db.finish_run(run_id, "failed", Some("build error"))?;
let run = db.run(run_id)?;
assert_eq!(run.status, "failed");
assert_eq!(run.error.as_deref(), Some("build error"));
assert!(run.finished_at.is_some(), "finished_at should be set");
Ok(())
}
#[test]
fn cancel_run_cancels_queued_run_and_jobs() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
let spec: crate::config::JobSpec = serde_json::from_value(serde_json::json!({
"base_name": "build",
"name": "build",
"needs": [],
"executor": "native",
"image": null,
"shell": null,
"command": "true",
"matrix": {},
"artifacts": [],
"cache": [],
"network": false,
"secrets": [],
"environment": null,
"approval": false,
"branches": [],
"timeout": 60
}))?;
db.insert_jobs(run_id, std::slice::from_ref(&spec))?;
db.cancel_run(run_id)?;
let detail = db.detail(run_id)?;
assert_eq!(detail.run.status, "canceled");
assert!(detail.run.finished_at.is_some());
assert!(detail.jobs.iter().all(|job| job.status == "canceled"));
Ok(())
}
#[test]
fn cancel_run_refuses_running_and_terminal_runs() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let running = db.create_run(repo.id, "main", None, "manual")?;
assert!(db.claim_run(running)?);
let err = db.cancel_run(running).expect_err("running must be refused");
assert!(format!("{err:#}").contains("running"), "got: {err:#}");
let finished = db.create_run(repo.id, "main", None, "manual")?;
db.finish_run(finished, "succeeded", None)?;
let err = db
.cancel_run(finished)
.expect_err("terminal must be refused");
assert!(format!("{err:#}").contains("succeeded"), "got: {err:#}");
let err = db.cancel_run(99999).expect_err("unknown must error");
assert!(format!("{err:#}").contains("unknown run"), "got: {err:#}");
Ok(())
}
#[test]
fn run_unknown_id_errors() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let err = db.run(99999).expect_err("unknown run should error");
let msg = format!("{err:#}");
assert!(msg.contains("unknown run"), "got: {msg}");
Ok(())
}
// ── retry creates a new run ──
#[test]
fn run_retry_creates_new_run_with_retry_trigger() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
db.finish_run(run_id, "failed", Some("oops"))?;
let retry_id = db.retry(run_id)?;
assert_ne!(retry_id, run_id, "retry creates a distinct run");
let retry_run = db.run(retry_id)?;
assert_eq!(retry_run.trigger, "retry");
assert_eq!(retry_run.status, "queued");
Ok(())
}
// ── jobs ──
fn dummy_job_spec(name: &str, base: &str) -> crate::config::JobSpec {
crate::config::JobSpec {
base_name: base.into(),
name: name.into(),
needs: vec![],
executor: "docker".into(),
image: Some("alpine".into()),
shell: None,
command: "echo ok".into(),
matrix: BTreeMap::new(),
artifacts: vec![],
cache: vec![],
network: false,
secrets: vec![],
environment: None,
approval: false,
branches: vec![],
timeout: 60,
}
}
#[test]
fn job_insert_and_status_transitions() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
let specs = vec![dummy_job_spec("test", "test")];
db.insert_jobs(run_id, &specs)?;
let jobs = db.jobs(run_id)?;
assert_eq!(jobs.len(), 1);
let job_id = jobs[0].id;
assert_eq!(jobs[0].status, "queued");
assert_eq!(jobs[0].name, "test");
// transition to running
db.set_job_status(job_id, "running", None, None)?;
let jobs = db.jobs(run_id)?;
assert_eq!(jobs[0].status, "running");
assert!(jobs[0].started_at.is_some());
// transition to succeeded with logs and exit code
db.set_job_status(job_id, "succeeded", Some("all good"), Some(0))?;
let jobs = db.jobs(run_id)?;
assert_eq!(jobs[0].status, "succeeded");
assert_eq!(jobs[0].logs.as_str(), "all good");
assert_eq!(jobs[0].exit_code, Some(0));
assert!(jobs[0].finished_at.is_some());
Ok(())
}
// ── artifacts ──
#[test]
fn artifact_store_and_retrieve_by_id() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
let specs = vec![dummy_job_spec("build", "build")];
db.insert_jobs(run_id, &specs)?;
let job_id = db.jobs(run_id)?[0].id;
let artifact_id = db.add_artifact(
run_id,
job_id,
"binary",
"dist/binary",
"abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789",
1024,
)?;
let art = db.artifact(artifact_id)?;
assert_eq!(art.name, "binary");
assert_eq!(
art.sha256,
"abcdef0123456789abcdef0123456789abcdef0123456789abcdef0123456789"
);
assert_eq!(art.bytes, 1024);
assert_eq!(art.path, "dist/binary");
Ok(())
}
#[test]
fn artifact_unknown_id_errors() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let err = db
.artifact(99999)
.expect_err("unknown artifact should error");
let msg = format!("{err:#}");
assert!(msg.contains("unknown artifact"), "got: {msg}");
Ok(())
}
#[test]
fn artifacts_by_run_empty_when_no_artifacts() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
let arts = db.artifacts(run_id)?;
assert!(arts.is_empty());
Ok(())
}
// ── deployments ──
#[test]
fn deployment_begin_and_finish() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
let specs = vec![dummy_job_spec("deploy", "deploy")];
db.insert_jobs(run_id, &specs)?;
let job_id = db.jobs(run_id)?[0].id;
let dep_id = db.begin_deployment(run_id, job_id, "production", "[]")?;
let deps = db.deployments(run_id)?;
assert_eq!(deps.len(), 1);
assert_eq!(deps[0].status, "running");
assert_eq!(deps[0].environment, "production");
db.finish_deployment(dep_id, "succeeded")?;
let deps = db.deployments(run_id)?;
assert_eq!(deps[0].status, "succeeded");
assert!(deps[0].finished_at.is_some());
Ok(())
}
// ── workers ──
#[test]
fn worker_register_heartbeat_and_list() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.register_worker("w1", "host1", "docker,native", 2)?;
let workers = db.workers()?;
assert_eq!(workers.len(), 1);
assert_eq!(workers[0].id, "w1");
assert_eq!(workers[0].capacity, 2);
assert_eq!(workers[0].status, "starting");
assert_eq!(workers[0].completed_runs, 0);
db.heartbeat_worker("w1", "running", Some(42), false)?;
let workers = db.workers()?;
assert_eq!(workers[0].status, "running");
assert_eq!(workers[0].current_run_id, Some(42));
db.heartbeat_worker("w1", "idle", None, true)?;
let workers = db.workers()?;
assert_eq!(workers[0].completed_runs, 1);
Ok(())
}
#[test]
fn worker_heartbeat_unknown_id_errors() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let err = db
.heartbeat_worker("ghost", "idle", None, false)
.expect_err("heartbeat for unknown worker should fail");
let msg = format!("{err:#}");
assert!(msg.contains("unknown worker"), "got: {msg}");
Ok(())
}
#[test]
fn worker_reset_sets_all_offline() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.register_worker("w1", "h1", "docker", 1)?;
db.register_worker("w2", "h2", "native", 1)?;
db.heartbeat_worker("w1", "running", Some(1), false)?;
db.reset_workers()?;
for w in db.workers()? {
assert_eq!(w.status, "offline");
}
Ok(())
}
// ── recover interrupted ──
#[test]
fn recover_interrupted_resets_running_states() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
// set to running directly (bypassing claim)
db.finish_run(run_id, "running", None)?;
let specs = vec![dummy_job_spec("job", "job")];
db.insert_jobs(run_id, &specs)?;
let job_id = db.jobs(run_id)?[0].id;
db.set_job_status(job_id, "running", None, None)?;
let _dep_id = db.begin_deployment(run_id, job_id, "staging", "[]")?;
let recovered = db.recover_interrupted()?;
assert!(recovered >= 3, "should recover run + job + deployment");
let run = db.run(run_id)?;
assert_eq!(run.status, "interrupted");
let jobs = db.jobs(run_id)?;
assert_eq!(jobs[0].status, "interrupted");
let deps = db.deployments(run_id)?;
assert_eq!(deps[0].status, "failed");
Ok(())
}
// ── claim_next_queued_run ──
#[test]
fn claim_next_queued_run_returns_none_when_empty() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
assert_eq!(db.claim_next_queued_run()?, None);
Ok(())
}
#[test]
fn claim_next_queued_run_picks_oldest_queued() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let r1 = db.create_run(repo.id, "ref1", None, "manual")?;
let r2 = db.create_run(repo.id, "ref2", None, "manual")?;
// both queued; claim_next should pick oldest (r1)
let claimed = db.claim_next_queued_run()?;
assert_eq!(claimed, Some(r1));
// next claim gets r2
let next = db.claim_next_queued_run()?;
assert_eq!(next, Some(r2));
Ok(())
}
// ── queued_runs count ──
#[test]
fn queued_runs_counts_correctly() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
assert_eq!(db.queued_runs()?, 0);
db.create_run(repo.id, "ref1", None, "manual")?;
db.create_run(repo.id, "ref2", None, "manual")?;
assert_eq!(db.queued_runs()?, 2);
db.claim_next_queued_run()?;
assert_eq!(db.queued_runs()?, 1);
Ok(())
}
// ── detail aggregates ──
#[test]
fn detail_aggregates_run_jobs_artifacts_deployments() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
let specs = vec![dummy_job_spec("build", "build")];
db.insert_jobs(run_id, &specs)?;
let job_id = db.jobs(run_id)?[0].id;
db.add_artifact(
run_id,
job_id,
"bin",
"dist/bin",
"a".repeat(64).as_str(),
42,
)?;
db.begin_deployment(run_id, job_id, "staging", "[]")?;
let detail = db.detail(run_id)?;
assert_eq!(detail.run.id, run_id);
assert_eq!(detail.jobs.len(), 1);
assert_eq!(detail.artifacts.len(), 1);
assert_eq!(detail.deployments.len(), 1);
Ok(())
}
// ── set_run_commit ──
#[test]
fn set_run_commit_updates_sha() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
assert!(db.run(run_id)?.commit_sha.is_none());
db.set_run_commit(run_id, &"a".repeat(40))?;
assert_eq!(
db.run(run_id)?.commit_sha.as_deref(),
Some("a".repeat(40).as_str())
);
Ok(())
}
// ── approve_environment ──
#[test]
fn approve_environment_promotes_waiting_jobs() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
let repo = db.add_repository("app", "https://example.com/app.git", "main")?;
let run_id = db.create_run(repo.id, "main", None, "manual")?;
let mut spec = dummy_job_spec("deploy", "deploy");
spec.approval = true;
spec.environment = Some("production".into());
db.insert_jobs(run_id, &[spec])?;
let job_id = db.jobs(run_id)?[0].id;
// jobs inserted with approval_required are created as 'queued'
// set to waiting to test approval promotion
db.set_job_status(job_id, "waiting", None, None)?;
db.finish_run(run_id, "waiting", None)?;
let changed = db.approve_environment(run_id, "production")?;
assert_eq!(changed, 1);
let jobs = db.jobs(run_id)?;
assert_eq!(jobs[0].status, "queued");
assert!(jobs[0].approved_at.is_some());
Ok(())
}
// ── open empty db readable ──
#[test]
fn check_ready_returns_ok_on_fresh_db() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.check_ready()?;
Ok(())
}
// ── public_repositories ──
#[test]
fn public_repositories_include_real_build_activity() -> Result<()> {
let dir = tempfile::tempdir()?;
let db = open_db(&dir)?;
db.add_repository("priv", "https://example.com/priv.git", "main")?;
let repository = db.add_repository("pub", "https://github.com/olibuijr/pub.git", "main")?;
db.set_repository_visibility("pub", "public")?;
let succeeded = db.create_run(repository.id, "main", None, "manual")?;
db.finish_run(succeeded, "succeeded", None)?;
db.create_run(repository.id, "feature", None, "manual")?;
let public = db.public_repositories()?;
assert_eq!(public.len(), 1);
assert_eq!(public[0].name, "pub");
assert_eq!(public[0].run_count, 2);
assert_eq!(public[0].successful_runs, 1);
assert!(public[0].last_activity_at >= public[0].created_at);
Ok(())
}
// ── original tests preserved ──
#[test]
fn stores_ci_state_and_reopens_details() -> Result<()> {
let database = Database::memory(KEY)?;
let repo = database.add_repository("app", "https://example.com/app.git", "main")?;
assert_eq!(repo.visibility, "private");
assert_eq!(
database
.set_repository_visibility("app", "public")?
.visibility,
"public"
);
let run = database.create_run(repo.id, "main", None, "manual")?;
assert_eq!(database.detail(run)?.run.status, "queued");
Ok(())
}
#[test]
fn queries_runs_and_tracks_workers() -> Result<()> {
let database = Database::memory(KEY)?;
let alpha = database.add_repository("alpha", "https://example.com/alpha.git", "main")?;
let beta = database.add_repository("beta", "https://example.com/beta.git", "main")?;
let succeeded = database.create_run(alpha.id, "main", None, "manual")?;
database.finish_run(succeeded, "succeeded", None)?;
let failed = database.create_run(beta.id, "release", None, "webhook")?;
database.finish_run(failed, "failed", Some("compiler error"))?;
let runs = database.query_runs(&RunQuery {
repository: Some("alpha".into()),
statuses: vec!["succeeded".into()],
search: Some("main".into()),
limit: 20,
..RunQuery::default()
})?;
assert_eq!(runs.len(), 1);
assert_eq!(runs[0].id, succeeded);
let multi = database.query_runs(&RunQuery {
repository: Some("alpha,beta".into()),
statuses: vec!["succeeded".into(), "failed".into()],
limit: 20,
..RunQuery::default()
})?;
assert_eq!(multi.len(), 2);
let repositories = database.query_repositories(&RepositoryQuery {
search: Some("example.com/beta".into()),
visibility: Some("private".into()),
..RepositoryQuery::default()
})?;
assert_eq!(repositories.len(), 1);
assert_eq!(repositories[0].name, "beta");
let queued = database.create_run(alpha.id, "feature", None, "manual")?;
assert_eq!(database.claim_next_queued_run()?, Some(queued));
assert_eq!(database.run(queued)?.status, "running");
assert_eq!(database.claim_next_queued_run()?, None);
database.register_worker("titan-1", "titan", "docker,native", 1)?;
database.heartbeat_worker("titan-1", "running", Some(succeeded), false)?;
database.heartbeat_worker("titan-1", "idle", None, true)?;
let workers = database.workers()?;
assert_eq!(workers.len(), 1);
assert_eq!(workers[0].completed_runs, 1);
assert_eq!(workers[0].status, "idle");
database.reset_workers()?;
assert_eq!(database.workers()?[0].status, "offline");
Ok(())
}
}