use super::{HistoryError, RunRecord, RunStatus, Transience};
use chrono::{DateTime, Utc};
use std::time::Duration;
pub const DDL: &[&str] = &[
"CREATE TABLE IF NOT EXISTS faucet_serve_runs (\
run_id TEXT PRIMARY KEY,\
name TEXT,\
status TEXT NOT NULL,\
submitted_at TEXT NOT NULL,\
finished_at TEXT,\
idempotency_key TEXT,\
owner TEXT,\
lease_expires_at TEXT,\
cancel_requested TEXT,\
body TEXT NOT NULL)",
"CREATE INDEX IF NOT EXISTS faucet_serve_runs_submitted_idx \
ON faucet_serve_runs (submitted_at)",
"CREATE INDEX IF NOT EXISTS faucet_serve_runs_status_lease_idx \
ON faucet_serve_runs (status, lease_expires_at)",
"CREATE INDEX IF NOT EXISTS faucet_serve_runs_pending_idx \
ON faucet_serve_runs (status, submitted_at)",
"CREATE TABLE IF NOT EXISTS faucet_serve_instances (\
instance_id TEXT PRIMARY KEY,\
started_at TEXT NOT NULL,\
last_heartbeat TEXT NOT NULL,\
listen TEXT,\
max_concurrent TEXT,\
in_flight TEXT)",
"CREATE INDEX IF NOT EXISTS faucet_serve_instances_hb_idx \
ON faucet_serve_instances (last_heartbeat)",
"CREATE TABLE IF NOT EXISTS faucet_serve_instance_caps (\
instance_id TEXT PRIMARY KEY,\
state_format TEXT NOT NULL)",
"CREATE TABLE IF NOT EXISTS faucet_serve_idem (\
key TEXT PRIMARY KEY,\
run_id TEXT NOT NULL,\
fingerprint TEXT NOT NULL,\
claimed_at TEXT NOT NULL)",
"CREATE TABLE IF NOT EXISTS faucet_serve_shards (\
run_id TEXT NOT NULL,\
shard_id TEXT NOT NULL,\
descriptor TEXT NOT NULL,\
size_estimate TEXT,\
status TEXT NOT NULL,\
owner TEXT,\
lease_expires_at TEXT,\
attempt TEXT NOT NULL,\
finished_at TEXT,\
PRIMARY KEY (run_id, shard_id))",
"CREATE INDEX IF NOT EXISTS faucet_serve_shards_claim_idx \
ON faucet_serve_shards (status, lease_expires_at)",
"CREATE TABLE IF NOT EXISTS faucet_serve_audit (\
id TEXT PRIMARY KEY,\
ts TEXT NOT NULL,\
principal TEXT NOT NULL,\
role TEXT NOT NULL,\
action TEXT NOT NULL,\
run_id TEXT,\
config_fingerprint TEXT,\
source_ip TEXT,\
result TEXT NOT NULL)",
"CREATE INDEX IF NOT EXISTS faucet_serve_audit_ts_idx \
ON faucet_serve_audit (ts)",
"CREATE TABLE IF NOT EXISTS faucet_serve_audit_tenants (\
id TEXT PRIMARY KEY,\
tenant TEXT NOT NULL)",
"CREATE TABLE IF NOT EXISTS faucet_serve_run_logs (\
run_id TEXT NOT NULL,\
seq TEXT NOT NULL,\
ts TEXT NOT NULL,\
level TEXT NOT NULL,\
line TEXT NOT NULL,\
PRIMARY KEY (run_id, seq))",
"CREATE INDEX IF NOT EXISTS faucet_serve_run_logs_ts_idx \
ON faucet_serve_run_logs (ts)",
"CREATE TABLE IF NOT EXISTS faucet_catalog_datasets (\
id TEXT PRIMARY KEY,\
uri TEXT NOT NULL,\
kind TEXT NOT NULL,\
last_seen TEXT NOT NULL,\
body TEXT NOT NULL)",
"CREATE TABLE IF NOT EXISTS faucet_catalog_schema_versions (\
dataset_id TEXT NOT NULL,\
version TEXT NOT NULL,\
recorded_at TEXT NOT NULL,\
body TEXT NOT NULL,\
PRIMARY KEY (dataset_id, version))",
"CREATE TABLE IF NOT EXISTS faucet_catalog_edges (\
src_id TEXT NOT NULL,\
dst_id TEXT NOT NULL,\
last_seen TEXT NOT NULL,\
body TEXT NOT NULL,\
PRIMARY KEY (src_id, dst_id))",
"CREATE TABLE IF NOT EXISTS faucet_catalog_stats (\
dataset_id TEXT NOT NULL,\
recorded_at TEXT NOT NULL,\
run_id TEXT NOT NULL,\
records TEXT NOT NULL,\
PRIMARY KEY (dataset_id, recorded_at))",
"CREATE TABLE IF NOT EXISTS faucet_catalog_profiles (\
dataset_id TEXT NOT NULL,\
recorded_at TEXT NOT NULL,\
run_id TEXT NOT NULL,\
body TEXT NOT NULL,\
PRIMARY KEY (dataset_id, recorded_at))",
"CREATE TABLE IF NOT EXISTS faucet_config_snapshots (\
pipeline TEXT PRIMARY KEY,\
recorded_at TEXT NOT NULL,\
faucet_version TEXT NOT NULL,\
body TEXT NOT NULL)",
"CREATE TABLE IF NOT EXISTS faucet_templates (\
id TEXT NOT NULL,\
version TEXT NOT NULL,\
name TEXT,\
created_at TEXT NOT NULL,\
body TEXT NOT NULL,\
PRIMARY KEY (id, version))",
"CREATE TABLE IF NOT EXISTS faucet_template_tags (\
id TEXT NOT NULL,\
tag TEXT NOT NULL,\
version TEXT NOT NULL,\
updated_at TEXT NOT NULL,\
PRIMARY KEY (id, tag))",
"CREATE TABLE IF NOT EXISTS faucet_template_launches (\
id TEXT NOT NULL,\
seq TEXT NOT NULL,\
version TEXT NOT NULL,\
launched_at TEXT NOT NULL,\
launched_by TEXT,\
PRIMARY KEY (id, seq))",
"CREATE TABLE IF NOT EXISTS faucet_template_deprecations (\
id TEXT PRIMARY KEY,\
deprecated_at TEXT NOT NULL,\
deprecated_by TEXT,\
reason TEXT)",
"CREATE TABLE IF NOT EXISTS faucet_template_version_deprecations (\
id TEXT NOT NULL,\
version TEXT NOT NULL,\
deprecated_at TEXT NOT NULL,\
deprecated_by TEXT,\
reason TEXT,\
PRIMARY KEY (id, version))",
"CREATE TABLE IF NOT EXISTS faucet_local_outputs (\
id TEXT PRIMARY KEY,\
path TEXT NOT NULL,\
dataset_id TEXT NOT NULL,\
pipeline TEXT NOT NULL,\
last_written_at TEXT NOT NULL,\
deleted_at TEXT,\
body TEXT NOT NULL)",
"CREATE INDEX IF NOT EXISTS faucet_local_outputs_written_idx \
ON faucet_local_outputs (last_written_at)",
"CREATE INDEX IF NOT EXISTS faucet_local_outputs_dataset_idx \
ON faucet_local_outputs (dataset_id)",
"CREATE TABLE IF NOT EXISTS faucet_usage (\
run_id TEXT NOT NULL,\
row_id TEXT NOT NULL,\
pipeline TEXT NOT NULL,\
recorded_at TEXT NOT NULL,\
body TEXT NOT NULL,\
PRIMARY KEY (run_id, row_id))",
"CREATE INDEX IF NOT EXISTS faucet_usage_recorded_idx \
ON faucet_usage (recorded_at)",
"CREATE INDEX IF NOT EXISTS faucet_usage_pipeline_idx \
ON faucet_usage (pipeline, recorded_at)",
"CREATE TABLE IF NOT EXISTS faucet_serve_changes (\
id TEXT PRIMARY KEY,\
kind TEXT NOT NULL,\
status TEXT NOT NULL,\
requester TEXT NOT NULL,\
created_at TEXT NOT NULL,\
expires_at TEXT NOT NULL,\
body TEXT NOT NULL)",
"CREATE INDEX IF NOT EXISTS faucet_serve_changes_status_idx \
ON faucet_serve_changes (status, created_at)",
"CREATE TABLE IF NOT EXISTS faucet_tenants (\
id TEXT PRIMARY KEY,\
updated_at TEXT NOT NULL,\
body TEXT NOT NULL)",
"CREATE TABLE IF NOT EXISTS faucet_tenant_connections (\
tenant TEXT NOT NULL,\
name TEXT NOT NULL,\
status TEXT NOT NULL,\
updated_at TEXT NOT NULL,\
body TEXT NOT NULL,\
PRIMARY KEY (tenant, name))",
"CREATE TABLE IF NOT EXISTS faucet_connect_sessions (\
state TEXT PRIMARY KEY,\
tenant TEXT NOT NULL,\
expires_at TEXT NOT NULL,\
body TEXT NOT NULL)",
"CREATE TABLE IF NOT EXISTS faucet_tenant_runs (\
run_id TEXT PRIMARY KEY,\
tenant TEXT NOT NULL)",
"CREATE INDEX IF NOT EXISTS faucet_tenant_runs_tenant_idx \
ON faucet_tenant_runs (tenant)",
"CREATE TABLE IF NOT EXISTS faucet_tenant_state_refs (\
tenant TEXT NOT NULL,\
state_key TEXT NOT NULL,\
body TEXT NOT NULL,\
PRIMARY KEY (tenant, state_key))",
];
pub fn usage_limit(filter: &crate::usage::UsageFilter) -> usize {
if filter.limit > 0 {
filter.limit
} else {
crate::usage::DEFAULT_LIST_LIMIT
}
}
pub fn change_limit(filter: &crate::serve::changes::ChangeListFilter) -> usize {
if filter.limit > 0 {
filter.limit
} else {
crate::serve::changes::DEFAULT_LIST_LIMIT
}
}
#[derive(Clone, Copy, Debug)]
pub enum Dialect {
Postgres,
Sqlite,
}
pub struct Stmts {
pub upsert: String,
pub select_body: String,
pub select_status: String,
pub select_submitted: String,
pub delete: String,
pub list: String,
pub purge_runs: String,
pub purge_idem: String,
pub select_orphans: String,
pub renew_leases: String,
pub insert_idem: String,
pub select_idem: String,
pub takeover_idem: String,
pub delete_idem_by_run: String,
pub select_pending: String,
pub claim_one: String,
pub reclaim_select: String,
pub reclaim_requeue: String,
pub reclaim_fail: String,
pub finalize_owned: String,
pub cancel_pending: String,
pub request_cancel: String,
pub pending_cancellations: String,
pub heartbeat_instance: String,
pub heartbeat_caps: String,
pub live_instances: String,
pub prune_instances: String,
pub prune_instance_caps: String,
pub insert_shard: String,
pub claim_shards_select: String,
pub claim_shard_one: String,
pub renew_shard_leases: String,
pub reclaim_shards_select: String,
pub reclaim_shard_requeue: String,
pub reclaim_shard_fail: String,
pub finalize_shard: String,
pub shard_progress: String,
pub pending_shard_cancellations: String,
pub select_sharded_parents: String,
pub finalize_sharded_parent: String,
pub delete_shards_by_run: String,
pub purge_orphan_shards: String,
pub insert_audit: String,
pub list_audit: String,
pub purge_audit: String,
pub insert_audit_tenant: String,
pub purge_audit_tenants: String,
pub insert_run_log: String,
pub list_run_logs: String,
pub run_log_truncated: String,
pub purge_run_logs: String,
pub catalog_select_dataset: String,
pub local_output_select: String,
pub local_output_upsert: String,
pub local_output_mark_deleted: String,
pub usage_insert: String,
pub change_upsert: String,
pub change_select: String,
pub change_delete: String,
pub usage_delete_run: String,
pub tenant_upsert: String,
pub tenant_select: String,
pub tenant_list: String,
pub tenant_delete: String,
pub tenant_delete_connections: String,
pub tenant_delete_sessions: String,
pub tenant_delete_runs: String,
pub tenant_delete_state_refs: String,
pub connection_upsert: String,
pub connection_select: String,
pub connection_list: String,
pub connection_delete: String,
pub session_insert: String,
pub session_purge: String,
pub session_take: String,
pub tenant_run_link: String,
pub tenant_run_delete: String,
pub purge_orphan_tenant_runs: String,
pub state_ref_insert: String,
pub state_ref_select: String,
pub catalog_upsert_dataset: String,
pub catalog_select_datasets: String,
pub catalog_insert_schema_version: String,
pub catalog_select_schema_versions: String,
pub catalog_upsert_edge: String,
pub catalog_select_edges: String,
pub catalog_insert_stat: String,
pub catalog_select_stats: String,
pub catalog_prune_stats: String,
pub catalog_insert_profile: String,
pub catalog_select_profiles: String,
pub catalog_prune_profiles: String,
pub catalog_upsert_config_snapshot: String,
pub catalog_select_config_snapshot: String,
pub template_max_version: String,
pub template_insert: String,
pub template_select_version: String,
pub template_select_latest: String,
pub template_select_all: String,
pub template_versions: String,
pub template_delete_version: String,
pub template_delete_all: String,
pub template_upsert_tag: String,
pub template_select_tags: String,
pub template_delete_tag: String,
pub template_delete_tags_all: String,
pub template_delete_tags_for_version: String,
pub template_max_launch_seq: String,
pub template_insert_launch: String,
pub template_select_launches: String,
pub template_delete_launches_all: String,
pub template_delete_launches_for_version: String,
pub template_upsert_deprecation: String,
pub template_select_deprecation: String,
pub template_delete_deprecation: String,
pub template_upsert_version_deprecation: String,
pub template_select_version_deprecations: String,
pub template_delete_version_deprecation: String,
pub template_delete_version_deprecations_all: String,
pub dialect: Dialect,
}
impl Stmts {
pub fn new(dialect: Dialect) -> Self {
match dialect {
Dialect::Postgres => Self::postgres(),
Dialect::Sqlite => Self::sqlite(),
}
}
pub fn local_output_query(
&self,
filter: &crate::local_outputs::LocalOutputFilter,
) -> (String, Vec<String>) {
let mut sql = String::from("SELECT body FROM faucet_local_outputs");
let mut binds: Vec<String> = Vec::new();
let mut clauses: Vec<String> = Vec::new();
if !filter.include_deleted {
clauses.push("deleted_at IS NULL".to_string());
}
if let Some(dataset_id) = &filter.dataset_id {
binds.push(dataset_id.clone());
clauses.push(format!("dataset_id={}", self.placeholder(binds.len())));
}
if let Some(pipeline) = &filter.pipeline {
binds.push(pipeline.clone());
clauses.push(format!("pipeline={}", self.placeholder(binds.len())));
}
if !clauses.is_empty() {
sql.push_str(" WHERE ");
sql.push_str(&clauses.join(" AND "));
}
sql.push_str(" ORDER BY last_written_at DESC, path ASC");
if filter.limit > 0 {
sql.push_str(&format!(" LIMIT {}", filter.limit));
}
(sql, binds)
}
pub fn usage_query(&self, filter: &crate::usage::UsageFilter) -> (String, Vec<String>) {
let mut sql = String::from("SELECT body FROM faucet_usage");
let mut binds: Vec<String> = Vec::new();
let mut clauses: Vec<String> = Vec::new();
if let Some(since) = filter.since {
binds.push(fmt_ts(since));
clauses.push(format!("recorded_at >= {}", self.placeholder(binds.len())));
}
if let Some(until) = filter.until {
binds.push(fmt_ts(until));
clauses.push(format!("recorded_at < {}", self.placeholder(binds.len())));
}
if let Some(pipeline) = &filter.pipeline {
binds.push(pipeline.clone());
clauses.push(format!("pipeline={}", self.placeholder(binds.len())));
}
if !clauses.is_empty() {
sql.push_str(" WHERE ");
sql.push_str(&clauses.join(" AND "));
}
sql.push_str(" ORDER BY recorded_at DESC, run_id ASC, row_id ASC");
if filter.tenant.is_none() {
sql.push_str(&format!(" LIMIT {}", usage_limit(filter)));
}
(sql, binds)
}
pub fn change_query(
&self,
filter: &crate::serve::changes::ChangeListFilter,
) -> (String, Vec<String>) {
let mut sql = String::from("SELECT body FROM faucet_serve_changes");
let mut binds: Vec<String> = Vec::new();
let mut clauses: Vec<String> = Vec::new();
if let Some(status) = filter.status {
binds.push(status.as_str().to_string());
clauses.push(format!("status={}", self.placeholder(binds.len())));
}
if let Some(kind) = filter.kind {
binds.push(kind.as_str().to_string());
clauses.push(format!("kind={}", self.placeholder(binds.len())));
}
if let Some(requester) = &filter.requester {
binds.push(requester.clone());
clauses.push(format!("requester={}", self.placeholder(binds.len())));
}
if !clauses.is_empty() {
sql.push_str(" WHERE ");
sql.push_str(&clauses.join(" AND "));
}
sql.push_str(" ORDER BY created_at DESC, id ASC");
if filter.tenant.is_none() {
sql.push_str(&format!(" LIMIT {}", change_limit(filter)));
}
(sql, binds)
}
fn placeholder(&self, n: usize) -> String {
match self.dialect {
Dialect::Postgres => format!("${n}"),
Dialect::Sqlite => "?".to_string(),
}
}
fn postgres() -> Self {
Self {
upsert: "INSERT INTO faucet_serve_runs \
(run_id,name,status,submitted_at,finished_at,idempotency_key,owner,lease_expires_at,body) \
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9) \
ON CONFLICT (run_id) DO UPDATE SET \
name=excluded.name,status=excluded.status,submitted_at=excluded.submitted_at,\
finished_at=excluded.finished_at,idempotency_key=excluded.idempotency_key,\
owner=excluded.owner,lease_expires_at=excluded.lease_expires_at,\
body=excluded.body"
.into(),
select_body: "SELECT body FROM faucet_serve_runs WHERE run_id=$1".into(),
select_status: "SELECT status FROM faucet_serve_runs WHERE run_id=$1".into(),
select_submitted: "SELECT submitted_at FROM faucet_serve_runs WHERE run_id=$1".into(),
delete: "DELETE FROM faucet_serve_runs WHERE run_id=$1".into(),
list: "SELECT body FROM faucet_serve_runs \
WHERE ($1::text IS NULL OR position(',' || status || ',' in $2::text) > 0) \
AND ($3::text IS NULL OR name = $4::text) \
AND ($5::text IS NULL OR submitted_at >= $6::text) \
AND ($7::text IS NULL OR submitted_at <= $8::text) \
AND ($9::text IS NULL OR (submitted_at < $10::text \
OR (submitted_at = $11::text AND run_id < $12::text))) \
AND ($13::text IS NULL OR run_id IN \
(SELECT run_id FROM faucet_tenant_runs WHERE tenant = $14::text)) \
ORDER BY submitted_at DESC, run_id DESC LIMIT $15"
.into(),
purge_runs: "DELETE FROM faucet_serve_runs \
WHERE status IN ('completed','failed','cancelled') \
AND finished_at IS NOT NULL AND finished_at < $1"
.into(),
purge_idem: "DELETE FROM faucet_serve_idem WHERE claimed_at < $1".into(),
select_orphans: "SELECT body FROM faucet_serve_runs \
WHERE status IN ('queued','running') \
AND (lease_expires_at IS NULL OR lease_expires_at < $1)"
.into(),
renew_leases: "UPDATE faucet_serve_runs SET lease_expires_at = $1 \
WHERE owner = $2 AND status IN ('queued','running')"
.into(),
insert_idem: "INSERT INTO faucet_serve_idem (key,run_id,fingerprint,claimed_at) \
VALUES ($1,$2,$3,$4) ON CONFLICT (key) DO NOTHING"
.into(),
select_idem: "SELECT run_id,fingerprint,claimed_at FROM faucet_serve_idem WHERE key=$1"
.into(),
takeover_idem: "UPDATE faucet_serve_idem \
SET run_id=$1,fingerprint=$2,claimed_at=$3 WHERE key=$4 AND claimed_at=$5"
.into(),
delete_idem_by_run: "DELETE FROM faucet_serve_idem WHERE run_id=$1".into(),
select_pending: "SELECT run_id, body FROM faucet_serve_runs \
WHERE status = 'pending' ORDER BY submitted_at ASC LIMIT $1"
.into(),
claim_one: "UPDATE faucet_serve_runs \
SET owner = $1, status = 'running', lease_expires_at = $2, body = $3 \
WHERE run_id = $4 AND status = 'pending'"
.into(),
reclaim_select: "SELECT body FROM faucet_serve_runs \
WHERE status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < $1)"
.into(),
reclaim_requeue: "UPDATE faucet_serve_runs \
SET status = 'pending', owner = NULL, lease_expires_at = NULL, \
body = $1 \
WHERE run_id = $2 AND status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < $3)"
.into(),
reclaim_fail: "UPDATE faucet_serve_runs \
SET status = 'failed', finished_at = $1, body = $2, owner = NULL \
WHERE run_id = $3 AND status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < $4)"
.into(),
finalize_owned: "UPDATE faucet_serve_runs \
SET status = $1, finished_at = $2, lease_expires_at = $3, body = $4 \
WHERE run_id = $5 AND owner = $6 \
AND status NOT IN ('completed','failed','cancelled')"
.into(),
cancel_pending: "UPDATE faucet_serve_runs \
SET status = 'cancelled', finished_at = $1, body = $2 \
WHERE run_id = $3 AND status = 'pending'"
.into(),
request_cancel: "UPDATE faucet_serve_runs \
SET cancel_requested = $1 WHERE run_id = $2 AND status IN ('running','sharded')"
.into(),
pending_cancellations: "SELECT run_id FROM faucet_serve_runs \
WHERE status = 'running' AND owner = $1 AND cancel_requested IS NOT NULL"
.into(),
heartbeat_instance: "INSERT INTO faucet_serve_instances \
(instance_id, started_at, last_heartbeat, listen, max_concurrent, in_flight) \
VALUES ($1,$2,$3,$4,$5,$6) \
ON CONFLICT (instance_id) DO UPDATE SET \
last_heartbeat = excluded.last_heartbeat, listen = excluded.listen, \
max_concurrent = excluded.max_concurrent, in_flight = excluded.in_flight"
.into(),
heartbeat_caps: "INSERT INTO faucet_serve_instance_caps (instance_id, state_format) \
VALUES ($1,$2) ON CONFLICT (instance_id) DO UPDATE SET \
state_format = excluded.state_format"
.into(),
live_instances: "SELECT i.instance_id, i.started_at, i.last_heartbeat, i.listen, \
i.max_concurrent, i.in_flight, c.state_format FROM faucet_serve_instances i \
LEFT JOIN faucet_serve_instance_caps c ON c.instance_id = i.instance_id \
WHERE i.last_heartbeat >= $1"
.into(),
prune_instances: "DELETE FROM faucet_serve_instances WHERE last_heartbeat < $1".into(),
prune_instance_caps: "DELETE FROM faucet_serve_instance_caps WHERE instance_id NOT IN (SELECT instance_id FROM faucet_serve_instances)".into(),
insert_shard: "INSERT INTO faucet_serve_shards \
(run_id, shard_id, descriptor, size_estimate, status, attempt) \
VALUES ($1,$2,$3,$4,'pending','0') \
ON CONFLICT (run_id, shard_id) DO NOTHING"
.into(),
claim_shards_select: "SELECT s.run_id, s.shard_id, s.descriptor, r.body \
FROM faucet_serve_shards s JOIN faucet_serve_runs r ON r.run_id = s.run_id \
WHERE s.status = 'pending' \
ORDER BY CAST(COALESCE(s.size_estimate, '0') AS BIGINT) DESC, s.run_id, s.shard_id \
LIMIT $1"
.into(),
claim_shard_one: "UPDATE faucet_serve_shards \
SET owner = $1, status = 'running', lease_expires_at = $2 \
WHERE run_id = $3 AND shard_id = $4 AND status = 'pending'"
.into(),
renew_shard_leases: "UPDATE faucet_serve_shards SET lease_expires_at = $1 \
WHERE owner = $2 AND status = 'running'"
.into(),
reclaim_shards_select: "SELECT run_id, shard_id, attempt FROM faucet_serve_shards \
WHERE status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < $1)"
.into(),
reclaim_shard_requeue: "UPDATE faucet_serve_shards \
SET status = 'pending', owner = NULL, lease_expires_at = NULL, attempt = $1 \
WHERE run_id = $2 AND shard_id = $3 AND status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < $4)"
.into(),
reclaim_shard_fail: "UPDATE faucet_serve_shards \
SET status = 'failed', finished_at = $1, owner = NULL \
WHERE run_id = $2 AND shard_id = $3 AND status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < $4)"
.into(),
finalize_shard: "UPDATE faucet_serve_shards \
SET status = $1, finished_at = $2 \
WHERE run_id = $3 AND shard_id = $4 AND owner = $5 AND status = 'running'"
.into(),
shard_progress: "SELECT status, COUNT(*) AS n FROM faucet_serve_shards \
WHERE run_id = $1 GROUP BY status"
.into(),
pending_shard_cancellations: "SELECT DISTINCT s.run_id \
FROM faucet_serve_shards s \
JOIN faucet_serve_runs r ON r.run_id = s.run_id \
WHERE s.owner = $1 AND s.status = 'running' \
AND r.cancel_requested IS NOT NULL"
.into(),
select_sharded_parents: "SELECT run_id FROM faucet_serve_runs \
WHERE status = 'sharded'"
.into(),
finalize_sharded_parent: "UPDATE faucet_serve_runs \
SET status = $1, finished_at = $2, body = $3 \
WHERE run_id = $4 AND status = 'sharded'"
.into(),
delete_shards_by_run: "DELETE FROM faucet_serve_shards WHERE run_id = $1".into(),
purge_orphan_shards: "DELETE FROM faucet_serve_shards \
WHERE run_id NOT IN (SELECT run_id FROM faucet_serve_runs)"
.into(),
insert_audit: "INSERT INTO faucet_serve_audit \
(id, ts, principal, role, action, run_id, config_fingerprint, source_ip, result) \
VALUES ($1,$2,$3,$4,$5,$6,$7,$8,$9)"
.into(),
list_audit: "SELECT a.id AS id, ts, principal, role, action, run_id, \
config_fingerprint, source_ip, result, t.tenant AS tenant \
FROM faucet_serve_audit a LEFT JOIN faucet_serve_audit_tenants t ON t.id = a.id \
WHERE ($1::text IS NULL OR principal = $2::text) \
AND ($3::text IS NULL OR action = $4::text) \
AND ($5::text IS NULL OR ts >= $6::text) \
AND ($7::text IS NULL OR ts <= $8::text) \
AND ($9::text IS NULL OR t.tenant = $10::text) \
ORDER BY a.ts DESC, a.id DESC LIMIT $11"
.into(),
purge_audit: "DELETE FROM faucet_serve_audit WHERE ts < $1".into(),
insert_audit_tenant: "INSERT INTO faucet_serve_audit_tenants (id, tenant) \
VALUES ($1,$2) ON CONFLICT (id) DO NOTHING"
.into(),
purge_audit_tenants: "DELETE FROM faucet_serve_audit_tenants \
WHERE id NOT IN (SELECT id FROM faucet_serve_audit)"
.into(),
insert_run_log: "INSERT INTO faucet_serve_run_logs \
(run_id, seq, ts, level, line) VALUES ($1,$2,$3,$4,$5) \
ON CONFLICT (run_id, seq) DO NOTHING"
.into(),
list_run_logs: "SELECT seq, ts, level, line FROM faucet_serve_run_logs \
WHERE run_id = $1 AND seq <> $2 AND ($3::text IS NULL OR seq > $4::text) \
ORDER BY seq ASC LIMIT $5"
.into(),
run_log_truncated: "SELECT 1 FROM faucet_serve_run_logs \
WHERE run_id = $1 AND seq = $2 LIMIT 1"
.into(),
purge_run_logs: "DELETE FROM faucet_serve_run_logs WHERE ts < $1".into(),
catalog_select_dataset: "SELECT body FROM faucet_catalog_datasets WHERE id=$1".into(),
local_output_select: "SELECT body FROM faucet_local_outputs WHERE id=$1".into(),
local_output_upsert: "INSERT INTO faucet_local_outputs \
(id, path, dataset_id, pipeline, last_written_at, deleted_at, body) \
VALUES ($1,$2,$3,$4,$5,$6,$7) \
ON CONFLICT (id) DO UPDATE SET path=excluded.path, \
dataset_id=excluded.dataset_id, pipeline=excluded.pipeline, \
last_written_at=excluded.last_written_at, deleted_at=excluded.deleted_at, \
body=excluded.body"
.into(),
local_output_mark_deleted: "UPDATE faucet_local_outputs \
SET deleted_at=$2, body=$3 WHERE id=$1"
.into(),
usage_insert: "INSERT INTO faucet_usage \
(run_id, row_id, pipeline, recorded_at, body) VALUES ($1,$2,$3,$4,$5) \
ON CONFLICT (run_id, row_id) DO NOTHING"
.into(),
change_upsert: "INSERT INTO faucet_serve_changes \
(id, kind, status, requester, created_at, expires_at, body) \
VALUES ($1,$2,$3,$4,$5,$6,$7) \
ON CONFLICT (id) DO UPDATE SET kind=excluded.kind, status=excluded.status, \
requester=excluded.requester, created_at=excluded.created_at, \
expires_at=excluded.expires_at, body=excluded.body"
.into(),
change_select: "SELECT body FROM faucet_serve_changes WHERE id=$1".into(),
change_delete: "DELETE FROM faucet_serve_changes WHERE id=$1".into(),
usage_delete_run: "DELETE FROM faucet_usage WHERE run_id=$1".into(),
tenant_upsert: "INSERT INTO faucet_tenants (id, updated_at, body) \
VALUES ($1,$2,$3) \
ON CONFLICT (id) DO UPDATE SET updated_at=excluded.updated_at, body=excluded.body"
.into(),
tenant_select: "SELECT body FROM faucet_tenants WHERE id=$1".into(),
tenant_list: "SELECT body FROM faucet_tenants ORDER BY id".into(),
tenant_delete: "DELETE FROM faucet_tenants WHERE id=$1".into(),
tenant_delete_connections: "DELETE FROM faucet_tenant_connections WHERE tenant=$1"
.into(),
tenant_delete_sessions: "DELETE FROM faucet_connect_sessions WHERE tenant=$1".into(),
tenant_delete_runs: "DELETE FROM faucet_tenant_runs WHERE tenant=$1".into(),
tenant_delete_state_refs: "DELETE FROM faucet_tenant_state_refs WHERE tenant=$1"
.into(),
connection_upsert: "INSERT INTO faucet_tenant_connections \
(tenant, name, status, updated_at, body) VALUES ($1,$2,$3,$4,$5) \
ON CONFLICT (tenant, name) DO UPDATE SET status=excluded.status, \
updated_at=excluded.updated_at, body=excluded.body"
.into(),
connection_select: "SELECT body FROM faucet_tenant_connections \
WHERE tenant=$1 AND name=$2"
.into(),
connection_list: "SELECT body FROM faucet_tenant_connections \
WHERE tenant=$1 ORDER BY name"
.into(),
connection_delete: "DELETE FROM faucet_tenant_connections \
WHERE tenant=$1 AND name=$2"
.into(),
session_insert: "INSERT INTO faucet_connect_sessions (state, tenant, expires_at, body) \
VALUES ($1,$2,$3,$4)"
.into(),
session_purge: "DELETE FROM faucet_connect_sessions WHERE expires_at < $1".into(),
session_take: "DELETE FROM faucet_connect_sessions WHERE state=$1 RETURNING body"
.into(),
tenant_run_link: "INSERT INTO faucet_tenant_runs (run_id, tenant) VALUES ($1,$2) \
ON CONFLICT (run_id) DO UPDATE SET tenant=excluded.tenant"
.into(),
tenant_run_delete: "DELETE FROM faucet_tenant_runs WHERE run_id=$1".into(),
purge_orphan_tenant_runs: "DELETE FROM faucet_tenant_runs \
WHERE run_id NOT IN (SELECT run_id FROM faucet_serve_runs)"
.into(),
state_ref_insert: "INSERT INTO faucet_tenant_state_refs (tenant, state_key, body) \
VALUES ($1,$2,$3) ON CONFLICT (tenant, state_key) DO NOTHING"
.into(),
state_ref_select: "SELECT body FROM faucet_tenant_state_refs \
WHERE tenant=$1 ORDER BY state_key"
.into(),
catalog_upsert_dataset: "INSERT INTO faucet_catalog_datasets \
(id, uri, kind, last_seen, body) VALUES ($1,$2,$3,$4,$5) \
ON CONFLICT (id) DO UPDATE SET uri=excluded.uri, kind=excluded.kind, \
last_seen=excluded.last_seen, body=excluded.body"
.into(),
catalog_select_datasets: "SELECT body FROM faucet_catalog_datasets".into(),
catalog_insert_schema_version: "INSERT INTO faucet_catalog_schema_versions \
(dataset_id, version, recorded_at, body) VALUES ($1,$2,$3,$4) \
ON CONFLICT (dataset_id, version) DO NOTHING"
.into(),
catalog_select_schema_versions: "SELECT body FROM faucet_catalog_schema_versions \
WHERE dataset_id=$1 ORDER BY CAST(version AS BIGINT) ASC"
.into(),
catalog_upsert_edge: "INSERT INTO faucet_catalog_edges \
(src_id, dst_id, last_seen, body) VALUES ($1,$2,$3,$4) \
ON CONFLICT (src_id, dst_id) DO UPDATE SET \
last_seen=excluded.last_seen, body=excluded.body"
.into(),
catalog_select_edges: "SELECT body FROM faucet_catalog_edges \
ORDER BY last_seen DESC, src_id, dst_id"
.into(),
catalog_insert_stat: "INSERT INTO faucet_catalog_stats \
(dataset_id, recorded_at, run_id, records) VALUES ($1,$2,$3,$4) \
ON CONFLICT (dataset_id, recorded_at) DO NOTHING"
.into(),
catalog_select_stats: "SELECT recorded_at, run_id, records \
FROM faucet_catalog_stats WHERE dataset_id=$1 \
ORDER BY recorded_at DESC LIMIT $2"
.into(),
catalog_prune_stats: "DELETE FROM faucet_catalog_stats \
WHERE dataset_id=$1 AND recorded_at NOT IN (\
SELECT recorded_at FROM faucet_catalog_stats WHERE dataset_id=$2 \
ORDER BY recorded_at DESC LIMIT $3)"
.into(),
catalog_insert_profile: "INSERT INTO faucet_catalog_profiles \
(dataset_id, recorded_at, run_id, body) VALUES ($1,$2,$3,$4) \
ON CONFLICT (dataset_id, recorded_at) DO NOTHING"
.into(),
catalog_select_profiles: "SELECT body FROM faucet_catalog_profiles \
WHERE dataset_id=$1 ORDER BY recorded_at DESC LIMIT $2"
.into(),
catalog_prune_profiles: "DELETE FROM faucet_catalog_profiles \
WHERE dataset_id=$1 AND recorded_at NOT IN (\
SELECT recorded_at FROM faucet_catalog_profiles WHERE dataset_id=$2 \
ORDER BY recorded_at DESC LIMIT $3)"
.into(),
catalog_upsert_config_snapshot: "INSERT INTO faucet_config_snapshots \
(pipeline, recorded_at, faucet_version, body) VALUES ($1,$2,$3,$4) \
ON CONFLICT (pipeline) DO UPDATE SET recorded_at=excluded.recorded_at, \
faucet_version=excluded.faucet_version, body=excluded.body"
.into(),
catalog_select_config_snapshot:
"SELECT body FROM faucet_config_snapshots WHERE pipeline=$1".into(),
template_max_version: "SELECT COALESCE(MAX(CAST(version AS BIGINT)), 0) AS v \
FROM faucet_templates WHERE id=$1"
.into(),
template_insert: "INSERT INTO faucet_templates \
(id, version, name, created_at, body) VALUES ($1,$2,$3,$4,$5)"
.into(),
template_select_version:
"SELECT body FROM faucet_templates WHERE id=$1 AND version=$2".into(),
template_select_latest: "SELECT body FROM faucet_templates WHERE id=$1 \
ORDER BY CAST(version AS BIGINT) DESC LIMIT 1"
.into(),
template_select_all: "SELECT body FROM faucet_templates".into(),
template_versions: "SELECT version FROM faucet_templates WHERE id=$1 \
ORDER BY CAST(version AS BIGINT) DESC"
.into(),
template_delete_version: "DELETE FROM faucet_templates WHERE id=$1 AND version=$2"
.into(),
template_delete_all: "DELETE FROM faucet_templates WHERE id=$1".into(),
template_upsert_tag: "INSERT INTO faucet_template_tags \
(id, tag, version, updated_at) VALUES ($1,$2,$3,$4) \
ON CONFLICT (id, tag) DO UPDATE SET version=excluded.version, \
updated_at=excluded.updated_at"
.into(),
template_select_tags: "SELECT tag, version FROM faucet_template_tags \
WHERE id=$1 ORDER BY tag"
.into(),
template_delete_tag: "DELETE FROM faucet_template_tags WHERE id=$1 AND tag=$2".into(),
template_delete_tags_all: "DELETE FROM faucet_template_tags WHERE id=$1".into(),
template_delete_tags_for_version:
"DELETE FROM faucet_template_tags WHERE id=$1 AND version=$2".into(),
template_max_launch_seq: "SELECT COALESCE(MAX(CAST(seq AS BIGINT)), 0) AS v \
FROM faucet_template_launches WHERE id=$1"
.into(),
template_insert_launch: "INSERT INTO faucet_template_launches \
(id, seq, version, launched_at, launched_by) VALUES ($1,$2,$3,$4,$5)"
.into(),
template_select_launches: "SELECT seq, version, launched_at, launched_by \
FROM faucet_template_launches WHERE id=$1 ORDER BY CAST(seq AS BIGINT) DESC"
.into(),
template_delete_launches_all: "DELETE FROM faucet_template_launches WHERE id=$1".into(),
template_delete_launches_for_version:
"DELETE FROM faucet_template_launches WHERE id=$1 AND version=$2".into(),
template_upsert_deprecation: "INSERT INTO faucet_template_deprecations \
(id, deprecated_at, deprecated_by, reason) VALUES ($1,$2,$3,$4) \
ON CONFLICT (id) DO UPDATE SET deprecated_at=excluded.deprecated_at, \
deprecated_by=excluded.deprecated_by, reason=excluded.reason"
.into(),
template_select_deprecation: "SELECT deprecated_at, deprecated_by, reason \
FROM faucet_template_deprecations WHERE id=$1"
.into(),
template_delete_deprecation: "DELETE FROM faucet_template_deprecations WHERE id=$1"
.into(),
template_upsert_version_deprecation: "INSERT INTO faucet_template_version_deprecations \
(id, version, deprecated_at, deprecated_by, reason) VALUES ($1,$2,$3,$4,$5) \
ON CONFLICT (id, version) DO UPDATE SET deprecated_at=excluded.deprecated_at, \
deprecated_by=excluded.deprecated_by, reason=excluded.reason"
.into(),
template_select_version_deprecations: "SELECT version, deprecated_at, deprecated_by, reason \
FROM faucet_template_version_deprecations WHERE id=$1 \
ORDER BY CAST(version AS BIGINT) DESC"
.into(),
template_delete_version_deprecation:
"DELETE FROM faucet_template_version_deprecations WHERE id=$1 AND version=$2".into(),
template_delete_version_deprecations_all:
"DELETE FROM faucet_template_version_deprecations WHERE id=$1".into(),
dialect: Dialect::Postgres,
}
}
fn sqlite() -> Self {
Self {
upsert: "INSERT INTO faucet_serve_runs \
(run_id,name,status,submitted_at,finished_at,idempotency_key,owner,lease_expires_at,body) \
VALUES (?,?,?,?,?,?,?,?,?) \
ON CONFLICT (run_id) DO UPDATE SET \
name=excluded.name,status=excluded.status,submitted_at=excluded.submitted_at,\
finished_at=excluded.finished_at,idempotency_key=excluded.idempotency_key,\
owner=excluded.owner,lease_expires_at=excluded.lease_expires_at,\
body=excluded.body"
.into(),
select_body: "SELECT body FROM faucet_serve_runs WHERE run_id=?".into(),
select_status: "SELECT status FROM faucet_serve_runs WHERE run_id=?".into(),
select_submitted: "SELECT submitted_at FROM faucet_serve_runs WHERE run_id=?".into(),
delete: "DELETE FROM faucet_serve_runs WHERE run_id=?".into(),
list: "SELECT body FROM faucet_serve_runs \
WHERE (? IS NULL OR instr(?, ',' || status || ',') > 0) \
AND (? IS NULL OR name = ?) \
AND (? IS NULL OR submitted_at >= ?) \
AND (? IS NULL OR submitted_at <= ?) \
AND (? IS NULL OR (submitted_at < ? \
OR (submitted_at = ? AND run_id < ?))) \
AND (? IS NULL OR run_id IN \
(SELECT run_id FROM faucet_tenant_runs WHERE tenant = ?)) \
ORDER BY submitted_at DESC, run_id DESC LIMIT ?"
.into(),
purge_runs: "DELETE FROM faucet_serve_runs \
WHERE status IN ('completed','failed','cancelled') \
AND finished_at IS NOT NULL AND finished_at < ?"
.into(),
purge_idem: "DELETE FROM faucet_serve_idem WHERE claimed_at < ?".into(),
select_orphans: "SELECT body FROM faucet_serve_runs \
WHERE status IN ('queued','running') \
AND (lease_expires_at IS NULL OR lease_expires_at < ?)"
.into(),
renew_leases: "UPDATE faucet_serve_runs SET lease_expires_at = ? \
WHERE owner = ? AND status IN ('queued','running')"
.into(),
insert_idem: "INSERT INTO faucet_serve_idem (key,run_id,fingerprint,claimed_at) \
VALUES (?,?,?,?) ON CONFLICT (key) DO NOTHING"
.into(),
select_idem: "SELECT run_id,fingerprint,claimed_at FROM faucet_serve_idem WHERE key=?"
.into(),
takeover_idem: "UPDATE faucet_serve_idem \
SET run_id=?,fingerprint=?,claimed_at=? WHERE key=? AND claimed_at=?"
.into(),
delete_idem_by_run: "DELETE FROM faucet_serve_idem WHERE run_id=?".into(),
select_pending: "SELECT run_id, body FROM faucet_serve_runs \
WHERE status = 'pending' ORDER BY submitted_at ASC LIMIT ?"
.into(),
claim_one: "UPDATE faucet_serve_runs \
SET owner = ?, status = 'running', lease_expires_at = ?, body = ? \
WHERE run_id = ? AND status = 'pending'"
.into(),
reclaim_select: "SELECT body FROM faucet_serve_runs \
WHERE status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < ?)"
.into(),
reclaim_requeue: "UPDATE faucet_serve_runs \
SET status = 'pending', owner = NULL, lease_expires_at = NULL, \
body = ? \
WHERE run_id = ? AND status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < ?)"
.into(),
reclaim_fail: "UPDATE faucet_serve_runs \
SET status = 'failed', finished_at = ?, body = ?, owner = NULL \
WHERE run_id = ? AND status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < ?)"
.into(),
finalize_owned: "UPDATE faucet_serve_runs \
SET status = ?, finished_at = ?, lease_expires_at = ?, body = ? \
WHERE run_id = ? AND owner = ? \
AND status NOT IN ('completed','failed','cancelled')"
.into(),
cancel_pending: "UPDATE faucet_serve_runs \
SET status = 'cancelled', finished_at = ?, body = ? \
WHERE run_id = ? AND status = 'pending'"
.into(),
request_cancel: "UPDATE faucet_serve_runs \
SET cancel_requested = ? WHERE run_id = ? AND status IN ('running','sharded')"
.into(),
pending_cancellations: "SELECT run_id FROM faucet_serve_runs \
WHERE status = 'running' AND owner = ? AND cancel_requested IS NOT NULL"
.into(),
heartbeat_instance: "INSERT INTO faucet_serve_instances \
(instance_id, started_at, last_heartbeat, listen, max_concurrent, in_flight) \
VALUES (?,?,?,?,?,?) \
ON CONFLICT (instance_id) DO UPDATE SET \
last_heartbeat = excluded.last_heartbeat, listen = excluded.listen, \
max_concurrent = excluded.max_concurrent, in_flight = excluded.in_flight"
.into(),
heartbeat_caps: "INSERT INTO faucet_serve_instance_caps (instance_id, state_format) \
VALUES (?,?) ON CONFLICT (instance_id) DO UPDATE SET \
state_format = excluded.state_format"
.into(),
live_instances: "SELECT i.instance_id, i.started_at, i.last_heartbeat, i.listen, \
i.max_concurrent, i.in_flight, c.state_format FROM faucet_serve_instances i \
LEFT JOIN faucet_serve_instance_caps c ON c.instance_id = i.instance_id \
WHERE i.last_heartbeat >= ?"
.into(),
prune_instances: "DELETE FROM faucet_serve_instances WHERE last_heartbeat < ?".into(),
prune_instance_caps: "DELETE FROM faucet_serve_instance_caps WHERE instance_id NOT IN (SELECT instance_id FROM faucet_serve_instances)".into(),
insert_shard: "INSERT INTO faucet_serve_shards \
(run_id, shard_id, descriptor, size_estimate, status, attempt) \
VALUES (?,?,?,?,'pending','0') \
ON CONFLICT (run_id, shard_id) DO NOTHING"
.into(),
claim_shards_select: "SELECT s.run_id, s.shard_id, s.descriptor, r.body \
FROM faucet_serve_shards s JOIN faucet_serve_runs r ON r.run_id = s.run_id \
WHERE s.status = 'pending' \
ORDER BY CAST(COALESCE(s.size_estimate, '0') AS INTEGER) DESC, s.run_id, s.shard_id \
LIMIT ?"
.into(),
claim_shard_one: "UPDATE faucet_serve_shards \
SET owner = ?, status = 'running', lease_expires_at = ? \
WHERE run_id = ? AND shard_id = ? AND status = 'pending'"
.into(),
renew_shard_leases: "UPDATE faucet_serve_shards SET lease_expires_at = ? \
WHERE owner = ? AND status = 'running'"
.into(),
reclaim_shards_select: "SELECT run_id, shard_id, attempt FROM faucet_serve_shards \
WHERE status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < ?)"
.into(),
reclaim_shard_requeue: "UPDATE faucet_serve_shards \
SET status = 'pending', owner = NULL, lease_expires_at = NULL, attempt = ? \
WHERE run_id = ? AND shard_id = ? AND status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < ?)"
.into(),
reclaim_shard_fail: "UPDATE faucet_serve_shards \
SET status = 'failed', finished_at = ?, owner = NULL \
WHERE run_id = ? AND shard_id = ? AND status = 'running' \
AND (lease_expires_at IS NULL OR lease_expires_at < ?)"
.into(),
finalize_shard: "UPDATE faucet_serve_shards \
SET status = ?, finished_at = ? \
WHERE run_id = ? AND shard_id = ? AND owner = ? AND status = 'running'"
.into(),
shard_progress: "SELECT status, COUNT(*) AS n FROM faucet_serve_shards \
WHERE run_id = ? GROUP BY status"
.into(),
pending_shard_cancellations: "SELECT DISTINCT s.run_id \
FROM faucet_serve_shards s \
JOIN faucet_serve_runs r ON r.run_id = s.run_id \
WHERE s.owner = ? AND s.status = 'running' \
AND r.cancel_requested IS NOT NULL"
.into(),
select_sharded_parents: "SELECT run_id FROM faucet_serve_runs \
WHERE status = 'sharded'"
.into(),
finalize_sharded_parent: "UPDATE faucet_serve_runs \
SET status = ?, finished_at = ?, body = ? \
WHERE run_id = ? AND status = 'sharded'"
.into(),
delete_shards_by_run: "DELETE FROM faucet_serve_shards WHERE run_id = ?".into(),
purge_orphan_shards: "DELETE FROM faucet_serve_shards \
WHERE run_id NOT IN (SELECT run_id FROM faucet_serve_runs)"
.into(),
insert_audit: "INSERT INTO faucet_serve_audit \
(id, ts, principal, role, action, run_id, config_fingerprint, source_ip, result) \
VALUES (?,?,?,?,?,?,?,?,?)"
.into(),
list_audit: "SELECT a.id AS id, ts, principal, role, action, run_id, \
config_fingerprint, source_ip, result, t.tenant AS tenant \
FROM faucet_serve_audit a LEFT JOIN faucet_serve_audit_tenants t ON t.id = a.id \
WHERE (? IS NULL OR principal = ?) \
AND (? IS NULL OR action = ?) \
AND (? IS NULL OR ts >= ?) \
AND (? IS NULL OR ts <= ?) \
AND (? IS NULL OR t.tenant = ?) \
ORDER BY a.ts DESC, a.id DESC LIMIT ?"
.into(),
purge_audit: "DELETE FROM faucet_serve_audit WHERE ts < ?".into(),
insert_audit_tenant: "INSERT INTO faucet_serve_audit_tenants (id, tenant) \
VALUES (?,?) ON CONFLICT (id) DO NOTHING"
.into(),
purge_audit_tenants: "DELETE FROM faucet_serve_audit_tenants \
WHERE id NOT IN (SELECT id FROM faucet_serve_audit)"
.into(),
insert_run_log: "INSERT INTO faucet_serve_run_logs \
(run_id, seq, ts, level, line) VALUES (?,?,?,?,?) \
ON CONFLICT (run_id, seq) DO NOTHING"
.into(),
list_run_logs: "SELECT seq, ts, level, line FROM faucet_serve_run_logs \
WHERE run_id = ? AND seq <> ? AND (? IS NULL OR seq > ?) \
ORDER BY seq ASC LIMIT ?"
.into(),
run_log_truncated: "SELECT 1 FROM faucet_serve_run_logs \
WHERE run_id = ? AND seq = ? LIMIT 1"
.into(),
purge_run_logs: "DELETE FROM faucet_serve_run_logs WHERE ts < ?".into(),
catalog_select_dataset: "SELECT body FROM faucet_catalog_datasets WHERE id=?".into(),
local_output_select: "SELECT body FROM faucet_local_outputs WHERE id=?".into(),
local_output_upsert: "INSERT INTO faucet_local_outputs \
(id, path, dataset_id, pipeline, last_written_at, deleted_at, body) \
VALUES (?,?,?,?,?,?,?) \
ON CONFLICT (id) DO UPDATE SET path=excluded.path, \
dataset_id=excluded.dataset_id, pipeline=excluded.pipeline, \
last_written_at=excluded.last_written_at, deleted_at=excluded.deleted_at, \
body=excluded.body"
.into(),
local_output_mark_deleted: "UPDATE faucet_local_outputs \
SET deleted_at=?2, body=?3 WHERE id=?1"
.into(),
usage_insert: "INSERT INTO faucet_usage \
(run_id, row_id, pipeline, recorded_at, body) VALUES (?,?,?,?,?) \
ON CONFLICT (run_id, row_id) DO NOTHING"
.into(),
change_upsert: "INSERT INTO faucet_serve_changes \
(id, kind, status, requester, created_at, expires_at, body) \
VALUES (?,?,?,?,?,?,?) \
ON CONFLICT (id) DO UPDATE SET kind=excluded.kind, status=excluded.status, \
requester=excluded.requester, created_at=excluded.created_at, \
expires_at=excluded.expires_at, body=excluded.body"
.into(),
change_select: "SELECT body FROM faucet_serve_changes WHERE id=?".into(),
change_delete: "DELETE FROM faucet_serve_changes WHERE id=?".into(),
usage_delete_run: "DELETE FROM faucet_usage WHERE run_id=?".into(),
tenant_upsert: "INSERT INTO faucet_tenants (id, updated_at, body) \
VALUES (?,?,?) \
ON CONFLICT (id) DO UPDATE SET updated_at=excluded.updated_at, body=excluded.body"
.into(),
tenant_select: "SELECT body FROM faucet_tenants WHERE id=?".into(),
tenant_list: "SELECT body FROM faucet_tenants ORDER BY id".into(),
tenant_delete: "DELETE FROM faucet_tenants WHERE id=?".into(),
tenant_delete_connections: "DELETE FROM faucet_tenant_connections WHERE tenant=?"
.into(),
tenant_delete_sessions: "DELETE FROM faucet_connect_sessions WHERE tenant=?".into(),
tenant_delete_runs: "DELETE FROM faucet_tenant_runs WHERE tenant=?".into(),
tenant_delete_state_refs: "DELETE FROM faucet_tenant_state_refs WHERE tenant=?"
.into(),
connection_upsert: "INSERT INTO faucet_tenant_connections \
(tenant, name, status, updated_at, body) VALUES (?,?,?,?,?) \
ON CONFLICT (tenant, name) DO UPDATE SET status=excluded.status, \
updated_at=excluded.updated_at, body=excluded.body"
.into(),
connection_select: "SELECT body FROM faucet_tenant_connections \
WHERE tenant=? AND name=?"
.into(),
connection_list: "SELECT body FROM faucet_tenant_connections \
WHERE tenant=? ORDER BY name"
.into(),
connection_delete: "DELETE FROM faucet_tenant_connections \
WHERE tenant=? AND name=?"
.into(),
session_insert: "INSERT INTO faucet_connect_sessions (state, tenant, expires_at, body) \
VALUES (?,?,?,?)"
.into(),
session_purge: "DELETE FROM faucet_connect_sessions WHERE expires_at < ?".into(),
session_take: "DELETE FROM faucet_connect_sessions WHERE state=? RETURNING body"
.into(),
tenant_run_link: "INSERT INTO faucet_tenant_runs (run_id, tenant) VALUES (?,?) \
ON CONFLICT (run_id) DO UPDATE SET tenant=excluded.tenant"
.into(),
tenant_run_delete: "DELETE FROM faucet_tenant_runs WHERE run_id=?".into(),
purge_orphan_tenant_runs: "DELETE FROM faucet_tenant_runs \
WHERE run_id NOT IN (SELECT run_id FROM faucet_serve_runs)"
.into(),
state_ref_insert: "INSERT INTO faucet_tenant_state_refs (tenant, state_key, body) \
VALUES (?,?,?) ON CONFLICT (tenant, state_key) DO NOTHING"
.into(),
state_ref_select: "SELECT body FROM faucet_tenant_state_refs \
WHERE tenant=? ORDER BY state_key"
.into(),
catalog_upsert_dataset: "INSERT INTO faucet_catalog_datasets \
(id, uri, kind, last_seen, body) VALUES (?,?,?,?,?) \
ON CONFLICT (id) DO UPDATE SET uri=excluded.uri, kind=excluded.kind, \
last_seen=excluded.last_seen, body=excluded.body"
.into(),
catalog_select_datasets: "SELECT body FROM faucet_catalog_datasets".into(),
catalog_insert_schema_version: "INSERT INTO faucet_catalog_schema_versions \
(dataset_id, version, recorded_at, body) VALUES (?,?,?,?) \
ON CONFLICT (dataset_id, version) DO NOTHING"
.into(),
catalog_select_schema_versions: "SELECT body FROM faucet_catalog_schema_versions \
WHERE dataset_id=? ORDER BY CAST(version AS INTEGER) ASC"
.into(),
catalog_upsert_edge: "INSERT INTO faucet_catalog_edges \
(src_id, dst_id, last_seen, body) VALUES (?,?,?,?) \
ON CONFLICT (src_id, dst_id) DO UPDATE SET \
last_seen=excluded.last_seen, body=excluded.body"
.into(),
catalog_select_edges: "SELECT body FROM faucet_catalog_edges \
ORDER BY last_seen DESC, src_id, dst_id"
.into(),
catalog_insert_stat: "INSERT INTO faucet_catalog_stats \
(dataset_id, recorded_at, run_id, records) VALUES (?,?,?,?) \
ON CONFLICT (dataset_id, recorded_at) DO NOTHING"
.into(),
catalog_select_stats: "SELECT recorded_at, run_id, records \
FROM faucet_catalog_stats WHERE dataset_id=? \
ORDER BY recorded_at DESC LIMIT ?"
.into(),
catalog_prune_stats: "DELETE FROM faucet_catalog_stats \
WHERE dataset_id=? AND recorded_at NOT IN (\
SELECT recorded_at FROM faucet_catalog_stats WHERE dataset_id=? \
ORDER BY recorded_at DESC LIMIT ?)"
.into(),
catalog_insert_profile: "INSERT INTO faucet_catalog_profiles \
(dataset_id, recorded_at, run_id, body) VALUES (?,?,?,?) \
ON CONFLICT (dataset_id, recorded_at) DO NOTHING"
.into(),
catalog_select_profiles: "SELECT body FROM faucet_catalog_profiles \
WHERE dataset_id=? ORDER BY recorded_at DESC LIMIT ?"
.into(),
catalog_prune_profiles: "DELETE FROM faucet_catalog_profiles \
WHERE dataset_id=? AND recorded_at NOT IN (\
SELECT recorded_at FROM faucet_catalog_profiles WHERE dataset_id=? \
ORDER BY recorded_at DESC LIMIT ?)"
.into(),
catalog_upsert_config_snapshot: "INSERT INTO faucet_config_snapshots \
(pipeline, recorded_at, faucet_version, body) VALUES (?,?,?,?) \
ON CONFLICT (pipeline) DO UPDATE SET recorded_at=excluded.recorded_at, \
faucet_version=excluded.faucet_version, body=excluded.body"
.into(),
catalog_select_config_snapshot:
"SELECT body FROM faucet_config_snapshots WHERE pipeline=?".into(),
template_max_version: "SELECT COALESCE(MAX(CAST(version AS INTEGER)), 0) AS v \
FROM faucet_templates WHERE id=?"
.into(),
template_insert: "INSERT INTO faucet_templates \
(id, version, name, created_at, body) VALUES (?,?,?,?,?)"
.into(),
template_select_version: "SELECT body FROM faucet_templates WHERE id=? AND version=?"
.into(),
template_select_latest: "SELECT body FROM faucet_templates WHERE id=? \
ORDER BY CAST(version AS INTEGER) DESC LIMIT 1"
.into(),
template_select_all: "SELECT body FROM faucet_templates".into(),
template_versions: "SELECT version FROM faucet_templates WHERE id=? \
ORDER BY CAST(version AS INTEGER) DESC"
.into(),
template_delete_version: "DELETE FROM faucet_templates WHERE id=? AND version=?".into(),
template_delete_all: "DELETE FROM faucet_templates WHERE id=?".into(),
template_upsert_tag: "INSERT INTO faucet_template_tags \
(id, tag, version, updated_at) VALUES (?,?,?,?) \
ON CONFLICT (id, tag) DO UPDATE SET version=excluded.version, \
updated_at=excluded.updated_at"
.into(),
template_select_tags: "SELECT tag, version FROM faucet_template_tags \
WHERE id=? ORDER BY tag"
.into(),
template_delete_tag: "DELETE FROM faucet_template_tags WHERE id=? AND tag=?".into(),
template_delete_tags_all: "DELETE FROM faucet_template_tags WHERE id=?".into(),
template_delete_tags_for_version:
"DELETE FROM faucet_template_tags WHERE id=? AND version=?".into(),
template_max_launch_seq: "SELECT COALESCE(MAX(CAST(seq AS INTEGER)), 0) AS v \
FROM faucet_template_launches WHERE id=?"
.into(),
template_insert_launch: "INSERT INTO faucet_template_launches \
(id, seq, version, launched_at, launched_by) VALUES (?,?,?,?,?)"
.into(),
template_select_launches: "SELECT seq, version, launched_at, launched_by \
FROM faucet_template_launches WHERE id=? ORDER BY CAST(seq AS INTEGER) DESC"
.into(),
template_delete_launches_all: "DELETE FROM faucet_template_launches WHERE id=?".into(),
template_delete_launches_for_version:
"DELETE FROM faucet_template_launches WHERE id=? AND version=?".into(),
template_upsert_deprecation: "INSERT INTO faucet_template_deprecations \
(id, deprecated_at, deprecated_by, reason) VALUES (?,?,?,?) \
ON CONFLICT (id) DO UPDATE SET deprecated_at=excluded.deprecated_at, \
deprecated_by=excluded.deprecated_by, reason=excluded.reason"
.into(),
template_select_deprecation: "SELECT deprecated_at, deprecated_by, reason \
FROM faucet_template_deprecations WHERE id=?"
.into(),
template_delete_deprecation: "DELETE FROM faucet_template_deprecations WHERE id=?"
.into(),
template_upsert_version_deprecation: "INSERT INTO faucet_template_version_deprecations \
(id, version, deprecated_at, deprecated_by, reason) VALUES (?,?,?,?,?) \
ON CONFLICT (id, version) DO UPDATE SET deprecated_at=excluded.deprecated_at, \
deprecated_by=excluded.deprecated_by, reason=excluded.reason"
.into(),
template_select_version_deprecations: "SELECT version, deprecated_at, deprecated_by, reason \
FROM faucet_template_version_deprecations WHERE id=? \
ORDER BY CAST(version AS BIGINT) DESC"
.into(),
template_delete_version_deprecation:
"DELETE FROM faucet_template_version_deprecations WHERE id=? AND version=?".into(),
template_delete_version_deprecations_all:
"DELETE FROM faucet_template_version_deprecations WHERE id=?".into(),
dialect: Dialect::Sqlite,
}
}
}
pub const CLAIM_ATTEMPTS: usize = 8;
static RETRY_SEQ: std::sync::atomic::AtomicUsize = std::sync::atomic::AtomicUsize::new(0);
pub async fn retry_backoff(attempt: usize) {
if attempt <= 1 {
return;
}
const BASE_MS: u64 = 5;
const CAP_MS: u64 = 160;
let exp = BASE_MS
.saturating_mul(1u64 << (attempt - 2).min(6))
.min(CAP_MS);
let stagger =
(RETRY_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed) as u64) % (BASE_MS + 1);
tokio::time::sleep(std::time::Duration::from_millis(exp + stagger)).await;
}
pub fn fmt_ts(dt: DateTime<Utc>) -> String {
dt.to_rfc3339_opts(chrono::SecondsFormat::Nanos, true)
}
pub fn pad_seq(seq: u64) -> String {
format!("{seq:020}")
}
pub fn unpad_seq(raw: &str) -> u64 {
raw.trim_start_matches('0').parse().unwrap_or(0)
}
pub fn parse_ts(raw: &str) -> DateTime<Utc> {
DateTime::parse_from_rfc3339(raw)
.map(|d| d.to_utc())
.unwrap_or_else(|_| Utc::now())
}
pub fn is_expired(claimed_at: &str, now: DateTime<Utc>, window: Duration) -> bool {
match DateTime::parse_from_rfc3339(claimed_at) {
Ok(t) => now
.signed_duration_since(t.with_timezone(&Utc))
.to_std()
.map(|age| age >= window)
.unwrap_or(false),
Err(_) => false,
}
}
pub fn threshold(now: DateTime<Utc>, window: Duration) -> String {
let delta =
chrono::Duration::from_std(window).unwrap_or_else(|_| chrono::Duration::days(36_500));
fmt_ts(now - delta)
}
pub fn classify_backend_error(e: sqlx::Error) -> HistoryError {
HistoryError::BackendClassified {
message: e.to_string(),
transience: transience_of(&e),
}
}
pub fn classify_backend_error_with_context(context: &str, e: sqlx::Error) -> HistoryError {
HistoryError::BackendClassified {
message: format!("{context}: {e}"),
transience: transience_of(&e),
}
}
pub fn transience_of(e: &sqlx::Error) -> Transience {
use sqlx::Error as E;
match e {
E::Database(db) => transience_of_code(db.code().as_deref()),
E::PoolTimedOut | E::Io(_) => Transience::Transient,
E::Configuration(_)
| E::InvalidArgument(_)
| E::Protocol(_)
| E::RowNotFound
| E::TypeNotFound { .. }
| E::ColumnIndexOutOfBounds { .. }
| E::ColumnNotFound(_)
| E::ColumnDecode { .. }
| E::Encode(_)
| E::Decode(_) => Transience::Permanent,
_ => Transience::Unknown,
}
}
const SQLSTATE_LEN: usize = 5;
pub fn transience_of_code(code: Option<&str>) -> Transience {
let Some(code) = code else {
return Transience::Unknown;
};
if code.len() == SQLSTATE_LEN {
return match code {
"57P03" | "53300" | "40001" | "40P01" => Transience::Transient,
_ if code.starts_with("08") => Transience::Transient,
_ => Transience::Permanent,
};
}
if let Ok(numeric) = code.parse::<i32>() {
return match numeric & 0xFF {
5 | 6 => Transience::Transient,
_ => Transience::Permanent,
};
}
Transience::Unknown
}
pub fn encode_body(rec: &RunRecord) -> Result<String, HistoryError> {
serde_json::to_string(rec).map_err(|e| HistoryError::Backend(format!("encode run record: {e}")))
}
pub fn decode_body(body: &str) -> Result<RunRecord, HistoryError> {
serde_json::from_str(body).map_err(|e| HistoryError::Backend(format!("decode run record: {e}")))
}
pub fn encode_json<T: serde::Serialize>(value: &T, what: &str) -> Result<String, HistoryError> {
serde_json::to_string(value).map_err(|e| HistoryError::Backend(format!("encode {what}: {e}")))
}
pub fn decode_json<T: serde::de::DeserializeOwned>(
body: &str,
what: &str,
) -> Result<T, HistoryError> {
serde_json::from_str(body).map_err(|e| HistoryError::Backend(format!("decode {what}: {e}")))
}
pub fn parse_status(s: &str) -> RunStatus {
match s {
"queued" => RunStatus::Queued,
"pending" => RunStatus::Pending,
"running" => RunStatus::Running,
"sharded" => RunStatus::Sharded,
"completed" => RunStatus::Completed,
"cancelled" => RunStatus::Cancelled,
_ => RunStatus::Failed,
}
}
macro_rules! impl_sql_history {
($name:ident, $pool:ty, $begin_write:expr) => {
pub struct $name {
pool: $pool,
idem_retention: std::time::Duration,
instance_id: String,
lease_ttl: std::time::Duration,
stmts: $crate::serve::history::sql::Stmts,
}
impl $name {
pub fn from_parts(
pool: $pool,
idem_retention: std::time::Duration,
lease_ttl: std::time::Duration,
instance_id: String,
stmts: $crate::serve::history::sql::Stmts,
) -> Self {
Self {
pool,
idem_retention,
instance_id,
lease_ttl,
stmts,
}
}
async fn select_one_body<T: serde::de::DeserializeOwned>(
&self,
stmt: &str,
binds: &[&str],
what: &str,
) -> Result<Option<T>, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
let backend = $crate::serve::history::sql::classify_backend_error;
let mut query = sqlx::query(stmt);
for b in binds {
query = query.bind(*b);
}
let Some(row) = query.fetch_optional(&self.pool).await.map_err(backend)? else {
return Ok(None);
};
let body: String = row.try_get("body").map_err(backend)?;
Ok(Some($crate::serve::history::sql::decode_json(&body, what)?))
}
async fn select_bodies<T: serde::de::DeserializeOwned>(
&self,
stmt: &str,
binds: &[&str],
what: &str,
) -> Result<Vec<T>, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
let backend = $crate::serve::history::sql::classify_backend_error;
let mut query = sqlx::query(stmt);
for b in binds {
query = query.bind(*b);
}
let rows = query.fetch_all(&self.pool).await.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let body: String = row.try_get("body").map_err(backend)?;
out.push($crate::serve::history::sql::decode_json(&body, what)?);
}
Ok(out)
}
pub fn pool(&self) -> &$pool {
&self.pool
}
}
#[async_trait::async_trait]
impl $crate::serve::history::RunHistory for $name {
async fn claim_idempotency(
&self,
key: &str,
fingerprint: &str,
run_id: &str,
window: std::time::Duration,
) -> Result<$crate::serve::history::Claim, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::serve::history::Claim;
use $crate::serve::history::sql;
let now = chrono::Utc::now();
let now_s = sql::fmt_ts(now);
let backend = $crate::serve::history::sql::classify_backend_error;
for _ in 0..sql::CLAIM_ATTEMPTS {
let inserted = sqlx::query(&self.stmts.insert_idem)
.bind(key)
.bind(run_id)
.bind(fingerprint)
.bind(&now_s)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if inserted == 1 {
return Ok(Claim::Fresh);
}
let Some(row) = sqlx::query(&self.stmts.select_idem)
.bind(key)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
continue;
};
let existing_run: String = row.try_get("run_id").map_err(backend)?;
let existing_fp: String = row.try_get("fingerprint").map_err(backend)?;
let claimed_at: String = row.try_get("claimed_at").map_err(backend)?;
if sql::is_expired(&claimed_at, now, window) {
let took = sqlx::query(&self.stmts.takeover_idem)
.bind(run_id)
.bind(fingerprint)
.bind(&now_s)
.bind(key)
.bind(&claimed_at)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if took == 1 {
return Ok(Claim::Fresh);
}
continue; }
return Ok(if existing_fp == fingerprint {
Claim::Replay(existing_run)
} else {
Claim::Conflict
});
}
tracing::warn!(
key,
"idempotency claim exhausted retries; reporting conflict"
);
Ok(Claim::Conflict)
}
async fn upsert(
&self,
rec: &$crate::serve::history::RunRecord,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let body = sql::encode_body(rec)?;
let submitted = sql::fmt_ts(rec.submitted_at);
let finished = rec.finished_at.map(sql::fmt_ts);
let lease = sql::fmt_ts(chrono::Utc::now() + self.lease_ttl);
sqlx::query(&self.stmts.upsert)
.bind(&rec.run_id)
.bind(rec.name.as_deref())
.bind(rec.status.as_str())
.bind(&submitted)
.bind(finished.as_deref())
.bind(rec.idempotency_key.as_deref())
.bind(&self.instance_id)
.bind(&lease)
.bind(&body)
.execute(&self.pool)
.await
.map_err($crate::serve::history::sql::classify_backend_error)?;
Ok(())
}
async fn get(
&self,
id: &str,
) -> Result<
Option<$crate::serve::history::RunRecord>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let row = sqlx::query(&self.stmts.select_body)
.bind(id)
.fetch_optional(&self.pool)
.await
.map_err($crate::serve::history::sql::classify_backend_error)?;
match row {
None => Ok(None),
Some(r) => {
let body: String = r
.try_get("body")
.map_err($crate::serve::history::sql::classify_backend_error)?;
Ok(Some(sql::decode_body(&body)?))
}
}
}
async fn list(
&self,
filter: &$crate::serve::history::ListFilter,
) -> Result<$crate::serve::history::ListPage, $crate::serve::history::HistoryError>
{
use sqlx::Row as _;
use $crate::serve::history::ListPage;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let cursor_ts: Option<String> = match &filter.cursor {
None => None,
Some(c) => sqlx::query(&self.stmts.select_submitted)
.bind(c)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
.map(|r| r.try_get::<String, _>("submitted_at"))
.transpose()
.map_err(backend)?,
};
let cur_id = if cursor_ts.is_some() {
filter.cursor.as_deref()
} else {
None
};
let status_s: Option<String> = if filter.status.is_empty() {
None
} else {
Some(format!(
",{},",
filter
.status
.iter()
.map(|s| s.as_str())
.collect::<Vec<_>>()
.join(",")
))
};
let status_s = status_s.as_deref();
let name_s = filter.name.as_deref();
let since_s = filter.since.map(sql::fmt_ts);
let until_s = filter.until.map(sql::fmt_ts);
let limit = filter.limit.max(1);
let fetch_n = limit as i64 + 1;
let rows = sqlx::query(&self.stmts.list)
.bind(status_s)
.bind(status_s)
.bind(name_s)
.bind(name_s)
.bind(since_s.as_deref())
.bind(since_s.as_deref())
.bind(until_s.as_deref())
.bind(until_s.as_deref())
.bind(cursor_ts.as_deref())
.bind(cursor_ts.as_deref())
.bind(cursor_ts.as_deref())
.bind(cur_id)
.bind(filter.tenant.as_deref())
.bind(filter.tenant.as_deref())
.bind(fetch_n)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut runs = Vec::with_capacity(rows.len());
for r in &rows {
let body: String = r.try_get("body").map_err(backend)?;
runs.push(sql::decode_body(&body)?);
}
let next_cursor = if runs.len() > limit {
Some(runs[limit - 1].run_id.clone())
} else {
None
};
runs.truncate(limit);
Ok(ListPage { runs, next_cursor })
}
async fn delete(
&self,
id: &str,
) -> Result<$crate::serve::history::DeleteOutcome, $crate::serve::history::HistoryError>
{
use sqlx::Row as _;
use $crate::serve::history::DeleteOutcome;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let status: Option<String> = sqlx::query(&self.stmts.select_status)
.bind(id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
.map(|r| r.try_get::<String, _>("status"))
.transpose()
.map_err(backend)?;
match status {
None => Ok(DeleteOutcome::NotFound),
Some(s) if !sql::parse_status(&s).is_terminal() => {
Ok(DeleteOutcome::StillRunning)
}
Some(_) => {
sqlx::query(&self.stmts.delete)
.bind(id)
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.tenant_run_delete)
.bind(id)
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.delete_idem_by_run)
.bind(id)
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.delete_shards_by_run)
.bind(id)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(DeleteOutcome::Deleted)
}
}
}
async fn release_idempotency(
&self,
run_id: &str,
) -> Result<(), $crate::serve::history::HistoryError> {
sqlx::query(&self.stmts.delete_idem_by_run)
.bind(run_id)
.execute(&self.pool)
.await
.map_err($crate::serve::history::sql::classify_backend_error)?;
Ok(())
}
async fn purge_expired(
&self,
retain_for: std::time::Duration,
) -> Result<usize, $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let now = chrono::Utc::now();
let removed = sqlx::query(&self.stmts.purge_runs)
.bind(sql::threshold(now, retain_for))
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected() as usize;
let _ = sqlx::query(&self.stmts.purge_idem)
.bind(sql::threshold(now, self.idem_retention))
.execute(&self.pool)
.await;
let _ = sqlx::query(&self.stmts.prune_instances)
.bind(sql::threshold(now, retain_for))
.execute(&self.pool)
.await;
let _ = sqlx::query(&self.stmts.prune_instance_caps)
.execute(&self.pool)
.await;
let _ = sqlx::query(&self.stmts.purge_orphan_shards)
.execute(&self.pool)
.await;
let _ = sqlx::query(&self.stmts.purge_orphan_tenant_runs)
.execute(&self.pool)
.await;
let _ = sqlx::query(&self.stmts.purge_audit)
.bind(sql::threshold(now, retain_for))
.execute(&self.pool)
.await;
let _ = sqlx::query(&self.stmts.purge_audit_tenants)
.execute(&self.pool)
.await;
Ok(removed)
}
async fn recover_orphans(&self) -> Result<usize, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::serve::history::RunStatus;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let now = chrono::Utc::now();
let rows = sqlx::query(&self.stmts.select_orphans)
.bind(sql::fmt_ts(now))
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut count = 0usize;
for r in &rows {
let body: String = r.try_get("body").map_err(backend)?;
let mut rec = sql::decode_body(&body)?;
rec.status = RunStatus::Failed;
rec.finished_at = Some(now);
rec.error = Some(
"owning serve instance's lease expired before the run finished".into(),
);
if rec.elapsed_secs.is_none()
&& let Some(started) = rec.started_at
{
rec.elapsed_secs = (now - started).to_std().ok().map(|d| d.as_secs_f64());
}
self.upsert(&rec).await?;
count += 1;
}
Ok(count)
}
async fn renew_leases(&self) -> Result<usize, $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let new_lease = sql::fmt_ts(chrono::Utc::now() + self.lease_ttl);
let renewed = sqlx::query(&self.stmts.renew_leases)
.bind(&new_lease)
.bind(&self.instance_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected() as usize;
Ok(renewed)
}
async fn claim_pending(
&self,
limit: usize,
) -> Result<Vec<$crate::serve::history::RunRecord>, $crate::serve::history::HistoryError>
{
use sqlx::Row as _;
use $crate::serve::history::RunStatus;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
if limit == 0 {
return Ok(Vec::new());
}
let now = chrono::Utc::now();
let lease = sql::fmt_ts(now + self.lease_ttl);
let rows = sqlx::query(&self.stmts.select_pending)
.bind(limit as i64)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut claimed = Vec::new();
for row in &rows {
let run_id: String = row.try_get("run_id").map_err(backend)?;
let body: String = row.try_get("body").map_err(backend)?;
let mut r = sql::decode_body(&body)?;
r.status = RunStatus::Running;
let new_body = sql::encode_body(&r)?;
let won = sqlx::query(&self.stmts.claim_one)
.bind(&self.instance_id)
.bind(&lease)
.bind(&new_body)
.bind(&run_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if won == 1 {
claimed.push(r);
}
}
Ok(claimed)
}
async fn reclaim_orphans(
&self,
max_attempts: u32,
) -> Result<$crate::serve::history::ReclaimReport, $crate::serve::history::HistoryError>
{
use sqlx::Row as _;
use $crate::serve::history::ReclaimReport;
use $crate::serve::history::RunStatus;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let now = chrono::Utc::now();
let now_s = sql::fmt_ts(now);
let rows = sqlx::query(&self.stmts.reclaim_select)
.bind(&now_s)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut report = ReclaimReport::default();
for row in &rows {
let body: String = row.try_get("body").map_err(backend)?;
let mut rec = sql::decode_body(&body)?;
let next_attempt = rec.attempt + 1;
if rec.attempt < max_attempts {
rec.attempt = next_attempt;
rec.status = RunStatus::Pending;
let new_body = sql::encode_body(&rec)?;
let n = sqlx::query(&self.stmts.reclaim_requeue)
.bind(&new_body)
.bind(&rec.run_id)
.bind(&now_s)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if n == 1 {
report.requeued += 1;
}
} else {
rec.attempt = next_attempt;
rec.status = RunStatus::Failed;
rec.finished_at = Some(now);
rec.error = Some(format!(
"run reclaimed {next_attempt} times after its owning instance's \
lease expired; giving up (poison run)"
));
if rec.elapsed_secs.is_none()
&& let Some(started) = rec.started_at
{
rec.elapsed_secs =
(now - started).to_std().ok().map(|d| d.as_secs_f64());
}
let new_body = sql::encode_body(&rec)?;
let n = sqlx::query(&self.stmts.reclaim_fail)
.bind(&now_s)
.bind(&new_body)
.bind(&rec.run_id)
.bind(&now_s)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if n == 1 {
report.failed += 1;
}
}
}
Ok(report)
}
async fn finalize_owned(
&self,
rec: &$crate::serve::history::RunRecord,
) -> Result<bool, $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let mut rec = rec.clone();
if rec.status.is_terminal() && rec.finished_at.is_none() {
rec.finished_at = Some(chrono::Utc::now());
}
let body = sql::encode_body(&rec)?;
let finished = rec.finished_at.map(sql::fmt_ts);
let lease = sql::fmt_ts(chrono::Utc::now() + self.lease_ttl);
let n = sqlx::query(&self.stmts.finalize_owned)
.bind(rec.status.as_str())
.bind(finished.as_deref())
.bind(&lease)
.bind(&body)
.bind(&rec.run_id)
.bind(&self.instance_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
Ok(n == 1)
}
async fn finalize_sharded_parent(
&self,
run_id: &str,
status: $crate::serve::history::RunStatus,
finished_at: chrono::DateTime<chrono::Utc>,
error: Option<String>,
) -> Result<bool, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::serve::history::RunStatus;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(row) = sqlx::query(&self.stmts.select_body)
.bind(run_id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
return Ok(false);
};
let body: String = row.try_get("body").map_err(backend)?;
let mut rec = sql::decode_body(&body)?;
if rec.status != RunStatus::Sharded {
return Ok(false);
}
rec.status = status;
rec.finished_at = Some(finished_at);
rec.error = error;
let new_body = sql::encode_body(&rec)?;
let n = sqlx::query(&self.stmts.finalize_sharded_parent)
.bind(status.as_str())
.bind(sql::fmt_ts(finished_at))
.bind(&new_body)
.bind(run_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
Ok(n == 1)
}
async fn cancel_pending(
&self,
run_id: &str,
) -> Result<bool, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::serve::history::RunStatus;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(row) = sqlx::query(&self.stmts.select_body)
.bind(run_id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
return Ok(false);
};
let body: String = row.try_get("body").map_err(backend)?;
let mut rec = sql::decode_body(&body)?;
if rec.status != RunStatus::Pending {
return Ok(false);
}
let now = chrono::Utc::now();
rec.status = RunStatus::Cancelled;
rec.finished_at = Some(now);
let new_body = sql::encode_body(&rec)?;
let n = sqlx::query(&self.stmts.cancel_pending)
.bind(sql::fmt_ts(now))
.bind(&new_body)
.bind(run_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
Ok(n == 1)
}
async fn request_cancel(
&self,
run_id: &str,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.request_cancel)
.bind(sql::fmt_ts(chrono::Utc::now()))
.bind(run_id)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn pending_cancellations(
&self,
) -> Result<Vec<String>, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.pending_cancellations)
.bind(&self.instance_id)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut ids = Vec::with_capacity(rows.len());
for r in &rows {
ids.push(r.try_get::<String, _>("run_id").map_err(backend)?);
}
Ok(ids)
}
async fn heartbeat_instance(
&self,
beat: &$crate::serve::history::InstanceHeartbeat,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let now = sql::fmt_ts(chrono::Utc::now());
sqlx::query(&self.stmts.heartbeat_instance)
.bind(&self.instance_id)
.bind(sql::fmt_ts(beat.started_at))
.bind(&now)
.bind(beat.listen.as_deref())
.bind(beat.max_concurrent.to_string())
.bind(beat.in_flight.to_string())
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.heartbeat_caps)
.bind(&self.instance_id)
.bind(beat.state_format.to_string())
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn live_instances(
&self,
ttl: std::time::Duration,
) -> Result<Vec<$crate::serve::history::InstanceRecord>, $crate::serve::history::HistoryError>
{
use sqlx::Row as _;
use $crate::serve::history::InstanceRecord;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let now = chrono::Utc::now();
let rows = sqlx::query(&self.stmts.live_instances)
.bind(sql::threshold(now, ttl))
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let parse_dt = |s: &str| {
chrono::DateTime::parse_from_rfc3339(s)
.map(|d| d.to_utc())
.unwrap_or(now)
};
let mut out = Vec::with_capacity(rows.len());
for r in &rows {
let started: String = r.try_get("started_at").map_err(backend)?;
let hb: String = r.try_get("last_heartbeat").map_err(backend)?;
let mc: Option<String> = r.try_get("max_concurrent").map_err(backend)?;
let inf: Option<String> = r.try_get("in_flight").map_err(backend)?;
let fmt: Option<String> = r.try_get("state_format").map_err(backend)?;
out.push(InstanceRecord {
instance_id: r.try_get("instance_id").map_err(backend)?,
started_at: parse_dt(&started),
last_heartbeat: parse_dt(&hb),
listen: r.try_get("listen").map_err(backend)?,
max_concurrent: mc.and_then(|s| s.parse().ok()).unwrap_or(0),
in_flight: inf.and_then(|s| s.parse().ok()).unwrap_or(0),
state_format: fmt.and_then(|s| s.parse().ok()).unwrap_or(0),
});
}
Ok(out)
}
async fn insert_shards(
&self,
run_id: &str,
shards: &[$crate::serve::history::ShardInsert],
) -> Result<usize, $crate::serve::history::HistoryError> {
use $crate::serve::history::HistoryError;
let backend = $crate::serve::history::sql::classify_backend_error;
let mut inserted = 0usize;
for s in shards {
let descriptor = serde_json::to_string(&s.descriptor).map_err(|e| {
HistoryError::Backend(format!("encode shard descriptor: {e}"))
})?;
let size = s.size_estimate.map(|n| n.to_string());
let n = sqlx::query(&self.stmts.insert_shard)
.bind(run_id)
.bind(&s.shard_id)
.bind(&descriptor)
.bind(size.as_deref())
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
inserted += n as usize;
}
Ok(inserted)
}
async fn claim_shards(
&self,
limit: usize,
) -> Result<
Vec<$crate::serve::history::ClaimedShard>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::ClaimedShard;
use $crate::serve::history::HistoryError;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
if limit == 0 {
return Ok(Vec::new());
}
let lease = sql::fmt_ts(chrono::Utc::now() + self.lease_ttl);
let rows = sqlx::query(&self.stmts.claim_shards_select)
.bind(limit as i64)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut claimed = Vec::new();
for row in &rows {
let run_id: String = row.try_get("run_id").map_err(backend)?;
let shard_id: String = row.try_get("shard_id").map_err(backend)?;
let descriptor_s: String = row.try_get("descriptor").map_err(backend)?;
let body: String = row.try_get("body").map_err(backend)?;
let won = sqlx::query(&self.stmts.claim_shard_one)
.bind(&self.instance_id)
.bind(&lease)
.bind(&run_id)
.bind(&shard_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if won == 1 {
let descriptor: serde_json::Value = serde_json::from_str(&descriptor_s)
.map_err(|e| {
HistoryError::Backend(format!("decode shard descriptor: {e}"))
})?;
let run = sql::decode_body(&body)?;
claimed.push(ClaimedShard {
run_id,
shard_id,
descriptor,
run,
});
}
}
Ok(claimed)
}
async fn renew_shard_leases(
&self,
) -> Result<usize, $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let lease = sql::fmt_ts(chrono::Utc::now() + self.lease_ttl);
let n = sqlx::query(&self.stmts.renew_shard_leases)
.bind(&lease)
.bind(&self.instance_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected() as usize;
Ok(n)
}
async fn reclaim_shards(
&self,
max_attempts: u32,
) -> Result<$crate::serve::history::ReclaimReport, $crate::serve::history::HistoryError>
{
use sqlx::Row as _;
use $crate::serve::history::ReclaimReport;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let now_s = sql::fmt_ts(chrono::Utc::now());
let rows = sqlx::query(&self.stmts.reclaim_shards_select)
.bind(&now_s)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut report = ReclaimReport::default();
for row in &rows {
let run_id: String = row.try_get("run_id").map_err(backend)?;
let shard_id: String = row.try_get("shard_id").map_err(backend)?;
let attempt_s: String = row.try_get("attempt").map_err(backend)?;
let attempt: u32 = attempt_s.parse().unwrap_or(0);
if attempt < max_attempts {
let next = (attempt + 1).to_string();
let n = sqlx::query(&self.stmts.reclaim_shard_requeue)
.bind(&next)
.bind(&run_id)
.bind(&shard_id)
.bind(&now_s)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if n == 1 {
report.requeued += 1;
}
} else {
let n = sqlx::query(&self.stmts.reclaim_shard_fail)
.bind(&now_s)
.bind(&run_id)
.bind(&shard_id)
.bind(&now_s)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if n == 1 {
report.failed += 1;
}
}
}
Ok(report)
}
async fn finalize_shard(
&self,
run_id: &str,
shard_id: &str,
success: bool,
) -> Result<bool, $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let status = if success { "completed" } else { "failed" };
let now_s = sql::fmt_ts(chrono::Utc::now());
let n = sqlx::query(&self.stmts.finalize_shard)
.bind(status)
.bind(&now_s)
.bind(run_id)
.bind(shard_id)
.bind(&self.instance_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
Ok(n == 1)
}
async fn shard_progress(
&self,
run_id: &str,
) -> Result<$crate::serve::history::ShardProgress, $crate::serve::history::HistoryError>
{
use sqlx::Row as _;
use $crate::serve::history::ShardProgress;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.shard_progress)
.bind(run_id)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut p = ShardProgress::default();
for row in &rows {
let status: String = row.try_get("status").map_err(backend)?;
let n: i64 = row.try_get("n").map_err(backend)?;
let n = n.max(0) as usize;
p.total += n;
match status.as_str() {
"completed" => p.completed += n,
"failed" => p.failed += n,
"running" => p.running += n,
_ => p.pending += n,
}
}
Ok(p)
}
async fn pending_shard_cancellations(
&self,
) -> Result<Vec<String>, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.pending_shard_cancellations)
.bind(&self.instance_id)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut ids = Vec::with_capacity(rows.len());
for r in &rows {
ids.push(r.try_get::<String, _>("run_id").map_err(backend)?);
}
Ok(ids)
}
async fn finalize_completed_sharded_parents(
&self,
) -> Result<usize, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::serve::history::RunStatus;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.select_sharded_parents)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut finalized = 0usize;
for row in &rows {
let run_id: String = row.try_get("run_id").map_err(backend)?;
let progress = self.shard_progress(&run_id).await?;
if !progress.all_terminal() {
continue;
}
let success = progress.failed == 0;
let Some(body_row) = sqlx::query(&self.stmts.select_body)
.bind(&run_id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
continue;
};
let body: String = body_row.try_get("body").map_err(backend)?;
let mut rec = sql::decode_body(&body)?;
if rec.status != RunStatus::Sharded {
continue;
}
let now = chrono::Utc::now();
rec.status = if success {
RunStatus::Completed
} else {
RunStatus::Failed
};
rec.finished_at = Some(now);
if !success {
rec.error = Some(format!(
"{}/{} shard(s) failed",
progress.failed, progress.total
));
}
let new_body = sql::encode_body(&rec)?;
let n = sqlx::query(&self.stmts.finalize_sharded_parent)
.bind(rec.status.as_str())
.bind(sql::fmt_ts(now))
.bind(&new_body)
.bind(&run_id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
if n == 1 {
finalized += 1;
$crate::serve::metrics::record_run_finished(
rec.status,
if success { "ok" } else { "error" },
);
tracing::info!(
run_id,
shards = progress.total,
failed = progress.failed,
"sharded run finalized by sweep (F11)"
);
}
}
Ok(finalized)
}
async fn record_audit(
&self,
entry: &$crate::serve::history::AuditEntry,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.insert_audit)
.bind(&entry.id)
.bind(sql::fmt_ts(entry.timestamp))
.bind(&entry.principal)
.bind(&entry.role)
.bind(&entry.action)
.bind(entry.run_id.as_deref())
.bind(entry.config_fingerprint.as_deref())
.bind(entry.source_ip.as_deref())
.bind(&entry.result)
.execute(&self.pool)
.await
.map_err(backend)?;
if let Some(tenant) = &entry.tenant {
sqlx::query(&self.stmts.insert_audit_tenant)
.bind(&entry.id)
.bind(tenant)
.execute(&self.pool)
.await
.map_err(backend)?;
}
Ok(())
}
async fn list_audit(
&self,
filter: &$crate::serve::history::AuditFilter,
) -> Result<
Vec<$crate::serve::history::AuditEntry>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::AuditEntry;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let principal = filter.principal.as_deref();
let action = filter.action.as_deref();
let since = filter.since.map(sql::fmt_ts);
let until = filter.until.map(sql::fmt_ts);
let limit = filter.limit.max(1) as i64;
let rows = sqlx::query(&self.stmts.list_audit)
.bind(principal)
.bind(principal)
.bind(action)
.bind(action)
.bind(since.as_deref())
.bind(since.as_deref())
.bind(until.as_deref())
.bind(until.as_deref())
.bind(filter.tenant.as_deref())
.bind(filter.tenant.as_deref())
.bind(limit)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for r in &rows {
let ts: String = r.try_get("ts").map_err(backend)?;
let timestamp = $crate::serve::history::sql::parse_ts(&ts);
out.push(AuditEntry {
id: r.try_get("id").map_err(backend)?,
timestamp,
principal: r.try_get("principal").map_err(backend)?,
role: r.try_get("role").map_err(backend)?,
action: r.try_get("action").map_err(backend)?,
run_id: r.try_get("run_id").map_err(backend)?,
config_fingerprint: r.try_get("config_fingerprint").map_err(backend)?,
source_ip: r.try_get("source_ip").map_err(backend)?,
tenant: r.try_get("tenant").map_err(backend)?,
result: r.try_get("result").map_err(backend)?,
});
}
Ok(out)
}
async fn record_run_logs(
&self,
run_id: &str,
lines: &[$crate::serve::history::RunLogLine],
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
if lines.is_empty() {
return Ok(());
}
let backend = $crate::serve::history::sql::classify_backend_error;
let mut tx = self.pool.begin().await.map_err(backend)?;
for l in lines {
sqlx::query(&self.stmts.insert_run_log)
.bind(run_id)
.bind(sql::pad_seq(l.seq))
.bind(&l.ts)
.bind(&l.level)
.bind(&l.line)
.execute(&mut *tx)
.await
.map_err(backend)?;
}
tx.commit().await.map_err(backend)?;
Ok(())
}
async fn list_run_logs(
&self,
run_id: &str,
after_seq: Option<u64>,
limit: usize,
) -> Result<
$crate::serve::history::RunLogPage,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
use $crate::serve::history::{RunLogLine, RunLogPage, RUN_LOG_TRUNCATED_SEQ};
let backend = $crate::serve::history::sql::classify_backend_error;
let sentinel = sql::pad_seq(RUN_LOG_TRUNCATED_SEQ);
let after = after_seq.map(sql::pad_seq);
let limit = limit.max(1) as i64;
let rows = sqlx::query(&self.stmts.list_run_logs)
.bind(run_id)
.bind(&sentinel)
.bind(after.as_deref())
.bind(after.as_deref())
.bind(limit)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for r in &rows {
let seq: String = r.try_get("seq").map_err(backend)?;
out.push(RunLogLine {
seq: sql::unpad_seq(&seq),
ts: r.try_get("ts").map_err(backend)?,
level: r.try_get("level").map_err(backend)?,
line: r.try_get("line").map_err(backend)?,
});
}
let truncated = sqlx::query(&self.stmts.run_log_truncated)
.bind(run_id)
.bind(&sentinel)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
.is_some();
Ok(RunLogPage { lines: out, truncated })
}
async fn purge_run_logs(
&self,
older_than: std::time::Duration,
) -> Result<usize, $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let cutoff = sql::threshold(chrono::Utc::now(), older_than);
let res = sqlx::query(&self.stmts.purge_run_logs)
.bind(cutoff)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(res.rows_affected() as usize)
}
async fn catalog_record(
&self,
update: &$crate::serve::history::catalog::CatalogUpdate,
) -> Result<(), $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::serve::history::catalog;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let now_s = sql::fmt_ts(update.recorded_at);
for obs in update.sources.iter().chain(std::iter::once(&update.sink)) {
let id = catalog::dataset_id(&obs.uri);
let existing = sqlx::query(&self.stmts.catalog_select_dataset)
.bind(&id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
.map(|r| r.try_get::<String, _>("body"))
.transpose()
.map_err(backend)?
.map(|b| {
sql::decode_json::<catalog::CatalogDataset>(&b, "catalog dataset")
})
.transpose()?;
let (ds, new_version) = catalog::apply_observation(
existing.as_ref(),
obs,
&update.run_id,
&update.pipeline,
&update.row,
update.recorded_at,
);
sqlx::query(&self.stmts.catalog_upsert_dataset)
.bind(&ds.id)
.bind(&ds.uri)
.bind(&ds.kind)
.bind(&now_s)
.bind(sql::encode_json(&ds, "catalog dataset")?)
.execute(&self.pool)
.await
.map_err(backend)?;
if let Some(v) = new_version {
sqlx::query(&self.stmts.catalog_insert_schema_version)
.bind(&v.dataset_id)
.bind(v.version.to_string())
.bind(sql::fmt_ts(v.recorded_at))
.bind(sql::encode_json(&v, "catalog schema version")?)
.execute(&self.pool)
.await
.map_err(backend)?;
}
sqlx::query(&self.stmts.catalog_insert_stat)
.bind(&id)
.bind(&now_s)
.bind(&update.run_id)
.bind(obs.records.to_string())
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.catalog_prune_stats)
.bind(&id)
.bind(&id)
.bind(catalog::STATS_RETAIN as i64)
.execute(&self.pool)
.await
.map_err(backend)?;
}
let dst_id = catalog::dataset_id(&update.sink.uri);
let existing_edges = self.catalog_all_edges().await?;
for source in &update.sources {
let src_id = catalog::dataset_id(&source.uri);
let existing = existing_edges
.iter()
.find(|e| e.src_id == src_id && e.dst_id == dst_id);
let edge = catalog::apply_edge(existing, update, source);
sqlx::query(&self.stmts.catalog_upsert_edge)
.bind(&edge.src_id)
.bind(&edge.dst_id)
.bind(&now_s)
.bind(sql::encode_json(&edge, "catalog edge")?)
.execute(&self.pool)
.await
.map_err(backend)?;
}
Ok(())
}
async fn catalog_annotate(
&self,
dataset_id: &str,
annotation: &$crate::serve::history::catalog::CatalogAnnotation,
) -> Result<bool, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::serve::history::catalog;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(row) = sqlx::query(&self.stmts.catalog_select_dataset)
.bind(dataset_id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
return Ok(false);
};
let body: String = row.try_get("body").map_err(backend)?;
let mut ds: catalog::CatalogDataset =
sql::decode_json(&body, "catalog dataset")?;
catalog::apply_annotation(&mut ds, annotation);
sqlx::query(&self.stmts.catalog_upsert_dataset)
.bind(&ds.id)
.bind(&ds.uri)
.bind(&ds.kind)
.bind(sql::fmt_ts(ds.last_seen))
.bind(sql::encode_json(&ds, "catalog dataset")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(true)
}
async fn catalog_list_datasets(
&self,
filter: &$crate::serve::history::catalog::CatalogListFilter,
) -> Result<
$crate::serve::history::catalog::CatalogDatasetPage,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::catalog;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.catalog_select_datasets)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut all = Vec::with_capacity(rows.len());
for r in &rows {
let body: String = r.try_get("body").map_err(backend)?;
all.push(sql::decode_json(&body, "catalog dataset")?);
}
Ok(catalog::filter_datasets(all, filter))
}
async fn catalog_get_dataset(
&self,
id: &str,
) -> Result<
Option<$crate::serve::history::catalog::CatalogDatasetDetail>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::catalog;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(row) = sqlx::query(&self.stmts.catalog_select_dataset)
.bind(id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
return Ok(None);
};
let body: String = row.try_get("body").map_err(backend)?;
let dataset: catalog::CatalogDataset =
sql::decode_json(&body, "catalog dataset")?;
let rows = sqlx::query(&self.stmts.catalog_select_schema_versions)
.bind(id)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut schema_timeline = Vec::with_capacity(rows.len());
for r in &rows {
let body: String = r.try_get("body").map_err(backend)?;
schema_timeline.push(sql::decode_json(&body, "catalog schema version")?);
}
let rows = sqlx::query(&self.stmts.catalog_select_stats)
.bind(id)
.bind(catalog::STATS_DETAIL_LIMIT as i64)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut stats = Vec::with_capacity(rows.len());
for r in &rows {
let recorded: String = r.try_get("recorded_at").map_err(backend)?;
let run_id: String = r.try_get("run_id").map_err(backend)?;
let records: String = r.try_get("records").map_err(backend)?;
stats.push(catalog::CatalogStatsPoint {
recorded_at: chrono::DateTime::parse_from_rfc3339(&recorded)
.map(|d| d.to_utc())
.unwrap_or_else(|_| chrono::Utc::now()),
run_id,
records: records.parse().unwrap_or(0),
});
}
let edges = self.catalog_all_edges().await?;
let (downstream, rest): (Vec<_>, Vec<_>) =
edges.into_iter().partition(|e| e.src_id == id);
let upstream = rest.into_iter().filter(|e| e.dst_id == id).collect();
let profile = catalog::profile_view(
self.catalog_profile_history(id, catalog::PROFILE_DETAIL_LIMIT)
.await?,
);
Ok(Some(catalog::CatalogDatasetDetail {
dataset,
schema_timeline,
stats,
upstream,
downstream,
profile,
}))
}
async fn catalog_record_profile(
&self,
dataset_id: &str,
record: &$crate::serve::history::catalog::CatalogProfileRecord,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::catalog;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.catalog_insert_profile)
.bind(dataset_id)
.bind(sql::fmt_ts(record.recorded_at))
.bind(&record.run_id)
.bind(sql::encode_json(record, "catalog profile")?)
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.catalog_prune_profiles)
.bind(dataset_id)
.bind(dataset_id)
.bind(catalog::PROFILE_RETAIN as i64)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn catalog_profile_history(
&self,
dataset_id: &str,
limit: usize,
) -> Result<
Vec<$crate::serve::history::catalog::CatalogProfileRecord>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.catalog_select_profiles)
.bind(dataset_id)
.bind(limit as i64)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for r in &rows {
let body: String = r.try_get("body").map_err(backend)?;
out.push(sql::decode_json(&body, "catalog profile")?);
}
Ok(out)
}
async fn catalog_lineage(
&self,
root: Option<&str>,
depth: u32,
) -> Result<
Vec<$crate::serve::history::catalog::CatalogLineageEdge>,
$crate::serve::history::HistoryError,
> {
use $crate::serve::history::catalog;
let edges = self.catalog_all_edges().await?;
Ok(catalog::lineage_slice(edges, root, depth))
}
async fn catalog_record_config_snapshot(
&self,
snapshot: &$crate::serve::history::catalog::ConfigSnapshot,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.catalog_upsert_config_snapshot)
.bind(&snapshot.pipeline)
.bind(sql::fmt_ts(snapshot.recorded_at))
.bind(&snapshot.faucet_version)
.bind(sql::encode_json(snapshot, "config snapshot")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn catalog_last_config_snapshot(
&self,
pipeline: &str,
) -> Result<
Option<$crate::serve::history::catalog::ConfigSnapshot>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(row) = sqlx::query(&self.stmts.catalog_select_config_snapshot)
.bind(pipeline)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
return Ok(None);
};
let body: String = row.try_get("body").map_err(backend)?;
Ok(Some(sql::decode_json(&body, "config snapshot")?))
}
async fn local_output_record(
&self,
obs: &$crate::local_outputs::LocalOutputObservation,
) -> Result<(), $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::local_outputs::LocalOutputRecord;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let id = $crate::local_outputs::ledger::output_id(&obs.path);
let existing: Option<LocalOutputRecord> =
match sqlx::query(&self.stmts.local_output_select)
.bind(&id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
{
Some(row) => {
let body: String = row.try_get("body").map_err(backend)?;
Some(sql::decode_json(&body, "local output")?)
}
None => None,
};
let rec = match existing {
Some(mut rec) => {
rec.observe(obs);
rec
}
None => LocalOutputRecord::new(obs),
};
sqlx::query(&self.stmts.local_output_upsert)
.bind(&rec.id)
.bind(&rec.path)
.bind(&rec.dataset_id)
.bind(&rec.pipeline)
.bind(sql::fmt_ts(rec.last_written_at))
.bind(rec.deleted_at.map(sql::fmt_ts))
.bind(sql::encode_json(&rec, "local output")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn local_output_list(
&self,
filter: &$crate::local_outputs::LocalOutputFilter,
) -> Result<
Vec<$crate::local_outputs::LocalOutputRecord>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let (sql, binds) = self.stmts.local_output_query(filter);
let mut query = sqlx::query(&sql);
for bind in &binds {
query = query.bind(bind);
}
let rows = query.fetch_all(&self.pool).await.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let body: String = row.try_get("body").map_err(backend)?;
let rec: $crate::local_outputs::LocalOutputRecord =
sql::decode_json(&body, "local output")?;
if $crate::local_outputs::ledger::matches(&rec, filter) {
out.push(rec);
}
}
Ok(out)
}
async fn local_output_get(
&self,
id: &str,
) -> Result<
Option<$crate::local_outputs::LocalOutputRecord>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(row) = sqlx::query(&self.stmts.local_output_select)
.bind(id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
return Ok(None);
};
let body: String = row.try_get("body").map_err(backend)?;
Ok(Some(sql::decode_json(&body, "local output")?))
}
async fn local_output_mark_deleted(
&self,
id: &str,
at: chrono::DateTime<chrono::Utc>,
bytes: u64,
) -> Result<bool, $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(mut rec) =
$crate::serve::history::RunHistory::local_output_get(self, id).await?
else {
return Ok(false);
};
rec.deleted_at = Some(at);
rec.deleted_bytes = Some(bytes);
let n = sqlx::query(&self.stmts.local_output_mark_deleted)
.bind(id)
.bind(sql::fmt_ts(at))
.bind(sql::encode_json(&rec, "local output")?)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
Ok(n > 0)
}
async fn tenant_upsert(
&self,
tenant: &$crate::serve::history::tenants::TenantRecord,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.tenant_upsert)
.bind(&tenant.id)
.bind(sql::fmt_ts(tenant.updated_at))
.bind(sql::encode_json(tenant, "tenant")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn tenant_get(
&self,
id: &str,
) -> Result<Option<$crate::serve::history::tenants::TenantRecord>, $crate::serve::history::HistoryError> {
self.select_one_body(&self.stmts.tenant_select, &[id], "tenant")
.await
}
async fn tenant_list(&self) -> Result<Vec<$crate::serve::history::tenants::TenantRecord>, $crate::serve::history::HistoryError> {
self.select_bodies(&self.stmts.tenant_list, &[], "tenant").await
}
async fn tenant_delete(&self, id: &str) -> Result<bool, $crate::serve::history::HistoryError> {
let backend = $crate::serve::history::sql::classify_backend_error;
let mut tx = self.pool.begin_with($begin_write).await.map_err(backend)?;
for stmt in [
&self.stmts.tenant_delete_connections,
&self.stmts.tenant_delete_sessions,
&self.stmts.tenant_delete_runs,
&self.stmts.tenant_delete_state_refs,
] {
sqlx::query(stmt)
.bind(id)
.execute(&mut *tx)
.await
.map_err(backend)?;
}
let n = sqlx::query(&self.stmts.tenant_delete)
.bind(id)
.execute(&mut *tx)
.await
.map_err(backend)?
.rows_affected();
tx.commit().await.map_err(backend)?;
Ok(n > 0)
}
async fn connection_upsert(
&self,
connection: &$crate::serve::history::tenants::ConnectionRecord,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.connection_upsert)
.bind(&connection.tenant)
.bind(&connection.name)
.bind(connection.status.as_str())
.bind(sql::fmt_ts(connection.updated_at))
.bind(sql::encode_json(connection, "connection")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn connection_get(
&self,
tenant: &str,
name: &str,
) -> Result<Option<$crate::serve::history::tenants::ConnectionRecord>, $crate::serve::history::HistoryError> {
self.select_one_body(&self.stmts.connection_select, &[tenant, name], "connection")
.await
}
async fn connection_list(
&self,
tenant: &str,
) -> Result<Vec<$crate::serve::history::tenants::ConnectionRecord>, $crate::serve::history::HistoryError> {
self.select_bodies(&self.stmts.connection_list, &[tenant], "connection")
.await
}
async fn connection_delete(&self, tenant: &str, name: &str) -> Result<bool, $crate::serve::history::HistoryError> {
let backend = $crate::serve::history::sql::classify_backend_error;
let n = sqlx::query(&self.stmts.connection_delete)
.bind(tenant)
.bind(name)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
Ok(n > 0)
}
async fn connect_session_put(
&self,
session: &$crate::serve::history::tenants::ConnectSession,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.session_purge)
.bind(sql::fmt_ts(chrono::Utc::now()))
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.session_insert)
.bind(&session.state)
.bind(&session.tenant)
.bind(sql::fmt_ts(session.expires_at))
.bind(sql::encode_json(session, "connect session")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn connect_session_take(
&self,
state: &str,
) -> Result<Option<$crate::serve::history::tenants::ConnectSession>, $crate::serve::history::HistoryError> {
self.select_one_body(&self.stmts.session_take, &[state], "connect session")
.await
}
async fn tenant_run_link(&self, run_id: &str, tenant: &str) -> Result<(), $crate::serve::history::HistoryError> {
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.tenant_run_link)
.bind(run_id)
.bind(tenant)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn tenant_state_ref_add(
&self,
state_ref: &$crate::serve::history::tenants::TenantStateRef,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.state_ref_insert)
.bind(&state_ref.tenant)
.bind(&state_ref.key)
.bind(sql::encode_json(state_ref, "tenant state ref")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn tenant_state_refs(
&self,
tenant: &str,
) -> Result<Vec<$crate::serve::history::tenants::TenantStateRef>, $crate::serve::history::HistoryError> {
self.select_bodies(&self.stmts.state_ref_select, &[tenant], "tenant state ref")
.await
}
async fn change_delete(
&self,
id: &str,
) -> Result<bool, $crate::serve::history::HistoryError> {
let backend = $crate::serve::history::sql::classify_backend_error;
let n = sqlx::query(&self.stmts.change_delete)
.bind(id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected();
Ok(n > 0)
}
async fn usage_delete_runs(
&self,
run_ids: &[String],
) -> Result<usize, $crate::serve::history::HistoryError> {
let backend = $crate::serve::history::sql::classify_backend_error;
let mut total = 0usize;
for id in run_ids {
total += sqlx::query(&self.stmts.usage_delete_run)
.bind(id)
.execute(&self.pool)
.await
.map_err(backend)?
.rows_affected() as usize;
}
Ok(total)
}
async fn change_upsert(
&self,
change: &$crate::serve::changes::ChangeRequest,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.change_upsert)
.bind(&change.id)
.bind(change.kind.as_str())
.bind(change.status.as_str())
.bind(&change.requester)
.bind(sql::fmt_ts(change.created_at))
.bind(sql::fmt_ts(change.expires_at))
.bind(sql::encode_json(change, "change request")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn change_get(
&self,
id: &str,
) -> Result<
Option<$crate::serve::changes::ChangeRequest>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(row) = sqlx::query(&self.stmts.change_select)
.bind(id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
return Ok(None);
};
let body: String = row.try_get("body").map_err(backend)?;
Ok(Some(sql::decode_json(&body, "change request")?))
}
async fn change_list(
&self,
filter: &$crate::serve::changes::ChangeListFilter,
) -> Result<
Vec<$crate::serve::changes::ChangeRequest>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let (sql, binds) = self.stmts.change_query(filter);
let mut query = sqlx::query(&sql);
for bind in &binds {
query = query.bind(bind);
}
let rows = query.fetch_all(&self.pool).await.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let body: String = row.try_get("body").map_err(backend)?;
let rec: $crate::serve::changes::ChangeRequest =
sql::decode_json(&body, "change request")?;
if filter.matches(&rec) {
out.push(rec);
}
}
out.truncate(sql::change_limit(filter));
Ok(out)
}
async fn usage_record(
&self,
record: &$crate::usage::UsageRecord,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.usage_insert)
.bind(&record.run_id)
.bind(&record.row)
.bind(&record.pipeline)
.bind(sql::fmt_ts(record.recorded_at))
.bind(sql::encode_json(record, "usage record")?)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn usage_list(
&self,
filter: &$crate::usage::UsageFilter,
) -> Result<Vec<$crate::usage::UsageRecord>, $crate::serve::history::HistoryError>
{
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let (sql, binds) = self.stmts.usage_query(filter);
let mut query = sqlx::query(&sql);
for bind in &binds {
query = query.bind(bind);
}
let rows = query.fetch_all(&self.pool).await.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let body: String = row.try_get("body").map_err(backend)?;
let rec: $crate::usage::UsageRecord = sql::decode_json(&body, "usage record")?;
if filter.matches(&rec) {
out.push(rec);
}
}
out.truncate(sql::usage_limit(filter));
Ok(out)
}
async fn template_register(
&self,
draft: &$crate::serve::history::templates::TemplateDraft,
) -> Result<
$crate::serve::history::templates::TemplateRecord,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::HistoryError;
use $crate::serve::history::{sql, templates};
let backend = $crate::serve::history::sql::classify_backend_error;
let id = draft.id.to_string();
for attempt in 1..=sql::CLAIM_ATTEMPTS {
sql::retry_backoff(attempt).await;
let mut tx = match self.pool.begin_with($begin_write).await {
Ok(tx) => tx,
Err(e) if attempt < sql::CLAIM_ATTEMPTS => {
tracing::debug!(
template = %id, attempt, error = %e,
"template version transaction lost a race; retrying"
);
continue;
}
Err(e) => return Err(backend(e)),
};
let row = match sqlx::query(&self.stmts.template_max_version)
.bind(&id)
.fetch_one(&mut *tx)
.await
{
Ok(row) => row,
Err(e) if attempt < sql::CLAIM_ATTEMPTS => {
let _ = tx.rollback().await;
tracing::debug!(
template = %id, attempt, error = %e,
"template version read lost a race; retrying"
);
continue;
}
Err(e) => {
let _ = tx.rollback().await;
return Err(backend(e));
}
};
let max: i64 = row.try_get("v").map_err(backend)?;
let next = (max as u32).saturating_add(1);
let record = templates::TemplateRecord {
id: id.clone(),
version: next,
kind: draft.kind,
name: draft.name.clone(),
description: draft.description.clone(),
body: draft.body.clone(),
format: draft.format,
params: draft.params.clone(),
created_at: chrono::Utc::now(),
created_by: draft.created_by.clone(),
};
let insert = sqlx::query(&self.stmts.template_insert)
.bind(&id)
.bind(next.to_string())
.bind(&record.name)
.bind(sql::fmt_ts(record.created_at))
.bind(sql::encode_json(&record, "pipeline template")?)
.execute(&mut *tx)
.await;
match insert {
Ok(_) => {
if let Err(e) = tx.commit().await {
if attempt < sql::CLAIM_ATTEMPTS {
tracing::debug!(
template = %id, attempt, error = %e,
"template version commit lost a race; retrying"
);
continue;
}
return Err(backend(e));
}
let keep = self.template_versions(&id).await?;
for stale in templates::versions_to_prune(keep) {
let _ = sqlx::query(&self.stmts.template_delete_version)
.bind(&id)
.bind(stale.to_string())
.execute(&self.pool)
.await;
}
return Ok(record);
}
Err(e) if attempt < sql::CLAIM_ATTEMPTS => {
let _ = tx.rollback().await;
tracing::debug!(
template = %id, attempt, error = %e,
"template version insert lost a race; retrying with the next version"
);
}
Err(e) => {
let _ = tx.rollback().await;
return Err(backend(e));
}
}
}
Err(HistoryError::Backend(format!(
"could not assign a version for template '{id}' after {} attempts",
sql::CLAIM_ATTEMPTS
)))
}
async fn template_get(
&self,
id: &str,
version: Option<u32>,
) -> Result<
Option<$crate::serve::history::templates::TemplateRecord>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let row = match version {
Some(v) => sqlx::query(&self.stmts.template_select_version)
.bind(id)
.bind(v.to_string())
.fetch_optional(&self.pool)
.await,
None => sqlx::query(&self.stmts.template_select_latest)
.bind(id)
.fetch_optional(&self.pool)
.await,
}
.map_err(backend)?;
let Some(row) = row else {
return Ok(None);
};
let body: String = row.try_get("body").map_err(backend)?;
Ok(Some(sql::decode_json(&body, "pipeline template")?))
}
async fn template_list(
&self,
) -> Result<
Vec<$crate::serve::history::templates::TemplateSummary>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::{sql, templates};
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.template_select_all)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut all = Vec::with_capacity(rows.len());
for r in &rows {
let body: String = r.try_get("body").map_err(backend)?;
all.push(sql::decode_json(&body, "pipeline template")?);
}
Ok(templates::latest_per_id(all))
}
async fn template_versions(
&self,
id: &str,
) -> Result<Vec<u32>, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.template_versions)
.bind(id)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for r in &rows {
let v: String = r.try_get("version").map_err(backend)?;
if let Ok(n) = v.parse::<u32>() {
out.push(n);
}
}
Ok(out)
}
async fn template_delete(
&self,
id: &str,
version: Option<u32>,
) -> Result<usize, $crate::serve::history::HistoryError> {
let backend = $crate::serve::history::sql::classify_backend_error;
let result = match version {
Some(v) => {
sqlx::query(&self.stmts.template_delete_tags_for_version)
.bind(id)
.bind(v.to_string())
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.template_delete_launches_for_version)
.bind(id)
.bind(v.to_string())
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.template_delete_version_deprecation)
.bind(id)
.bind(v.to_string())
.execute(&self.pool)
.await
.map_err(backend)?;
sqlx::query(&self.stmts.template_delete_version)
.bind(id)
.bind(v.to_string())
.execute(&self.pool)
.await
}
None => {
for stmt in [
&self.stmts.template_delete_tags_all,
&self.stmts.template_delete_launches_all,
&self.stmts.template_delete_deprecation,
&self.stmts.template_delete_version_deprecations_all,
] {
sqlx::query(stmt)
.bind(id)
.execute(&self.pool)
.await
.map_err(backend)?;
}
sqlx::query(&self.stmts.template_delete_all)
.bind(id)
.execute(&self.pool)
.await
}
}
.map_err(backend)?;
Ok(result.rows_affected() as usize)
}
async fn template_set_tag(
&self,
id: &str,
tag: &str,
version: u32,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
sqlx::query(&self.stmts.template_upsert_tag)
.bind(id)
.bind(tag)
.bind(version.to_string())
.bind(sql::fmt_ts(chrono::Utc::now()))
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(())
}
async fn template_tags(
&self,
id: &str,
) -> Result<
std::collections::BTreeMap<String, u32>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.template_select_tags)
.bind(id)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut out = std::collections::BTreeMap::new();
for r in &rows {
let tag: String = r.try_get("tag").map_err(backend)?;
let version: String = r.try_get("version").map_err(backend)?;
if let Ok(n) = version.parse::<u32>() {
out.insert(tag, n);
}
}
Ok(out)
}
async fn template_delete_tag(
&self,
id: &str,
tag: &str,
) -> Result<bool, $crate::serve::history::HistoryError> {
let backend = $crate::serve::history::sql::classify_backend_error;
let result = sqlx::query(&self.stmts.template_delete_tag)
.bind(id)
.bind(tag)
.execute(&self.pool)
.await
.map_err(backend)?;
Ok(result.rows_affected() > 0)
}
async fn template_launch(
&self,
id: &str,
version: u32,
launched_by: Option<&str>,
) -> Result<Option<u32>, $crate::serve::history::HistoryError> {
use sqlx::Row as _;
use $crate::serve::history::HistoryError;
use $crate::serve::history::{sql, templates};
let backend = $crate::serve::history::sql::classify_backend_error;
let existing = self.template_launches(id).await?;
if templates::stable_version(&existing) == Some(version) {
return Ok(None);
}
for attempt in 1..=sql::CLAIM_ATTEMPTS {
sql::retry_backoff(attempt).await;
let row = match sqlx::query(&self.stmts.template_max_launch_seq)
.bind(id)
.fetch_one(&self.pool)
.await
{
Ok(row) => row,
Err(e) if attempt < sql::CLAIM_ATTEMPTS => {
tracing::debug!(
template = %id, attempt, error = %e,
"launch-log seq read lost a race; retrying"
);
continue;
}
Err(e) => return Err(backend(e)),
};
let max: i64 = row.try_get("v").map_err(backend)?;
let seq = (max as u32).saturating_add(1);
let insert = sqlx::query(&self.stmts.template_insert_launch)
.bind(id)
.bind(seq.to_string())
.bind(version.to_string())
.bind(sql::fmt_ts(chrono::Utc::now()))
.bind(launched_by)
.execute(&self.pool)
.await;
match insert {
Ok(_) => return Ok(Some(seq)),
Err(e) if attempt < sql::CLAIM_ATTEMPTS => {
tracing::debug!(
template = %id, attempt, error = %e,
"launch-log insert lost a race; retrying with the next seq"
);
}
Err(e) => return Err(backend(e)),
}
}
Err(HistoryError::Backend(format!(
"could not append a launch for template '{id}' after {} attempts",
sql::CLAIM_ATTEMPTS
)))
}
async fn template_launches(
&self,
id: &str,
) -> Result<
Vec<$crate::serve::history::templates::LaunchRecord>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::{sql, templates};
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.template_select_launches)
.bind(id)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for r in &rows {
let seq: String = r.try_get("seq").map_err(backend)?;
let version: String = r.try_get("version").map_err(backend)?;
let launched_at: String = r.try_get("launched_at").map_err(backend)?;
let launched_by: Option<String> = r.try_get("launched_by").map_err(backend)?;
let (Ok(seq), Ok(version)) = (seq.parse::<u32>(), version.parse::<u32>()) else {
continue;
};
out.push(templates::LaunchRecord {
seq,
version,
launched_at: sql::parse_ts(&launched_at),
launched_by,
});
}
Ok(out)
}
async fn template_set_deprecation(
&self,
id: &str,
record: Option<&$crate::serve::history::templates::DeprecationRecord>,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
match record {
Some(r) => {
sqlx::query(&self.stmts.template_upsert_deprecation)
.bind(id)
.bind(sql::fmt_ts(r.deprecated_at))
.bind(r.deprecated_by.as_deref())
.bind(r.reason.as_deref())
.execute(&self.pool)
.await
.map_err(backend)?;
}
None => {
sqlx::query(&self.stmts.template_delete_deprecation)
.bind(id)
.execute(&self.pool)
.await
.map_err(backend)?;
}
}
Ok(())
}
async fn template_set_version_deprecation(
&self,
id: &str,
version: u32,
record: Option<&$crate::serve::history::templates::DeprecationRecord>,
) -> Result<(), $crate::serve::history::HistoryError> {
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
match record {
Some(r) => {
sqlx::query(&self.stmts.template_upsert_version_deprecation)
.bind(id)
.bind(version.to_string())
.bind(sql::fmt_ts(r.deprecated_at))
.bind(r.deprecated_by.as_deref())
.bind(r.reason.as_deref())
.execute(&self.pool)
.await
.map_err(backend)?;
}
None => {
sqlx::query(&self.stmts.template_delete_version_deprecation)
.bind(id)
.bind(version.to_string())
.execute(&self.pool)
.await
.map_err(backend)?;
}
}
Ok(())
}
async fn template_version_deprecations(
&self,
id: &str,
) -> Result<
Vec<$crate::serve::history::templates::VersionDeprecation>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::{sql, templates};
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.template_select_version_deprecations)
.bind(id)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut out = Vec::with_capacity(rows.len());
for row in rows {
let version: String = row.try_get("version").map_err(backend)?;
let Ok(version) = version.parse::<u32>() else {
continue;
};
let at: String = row.try_get("deprecated_at").map_err(backend)?;
out.push(templates::VersionDeprecation {
version,
record: templates::DeprecationRecord {
deprecated_at: sql::parse_ts(&at),
deprecated_by: row.try_get("deprecated_by").map_err(backend)?,
reason: row.try_get("reason").map_err(backend)?,
},
});
}
Ok(out)
}
async fn template_deprecation(
&self,
id: &str,
) -> Result<
Option<$crate::serve::history::templates::DeprecationRecord>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::{sql, templates};
let backend = $crate::serve::history::sql::classify_backend_error;
let Some(row) = sqlx::query(&self.stmts.template_select_deprecation)
.bind(id)
.fetch_optional(&self.pool)
.await
.map_err(backend)?
else {
return Ok(None);
};
let at: String = row.try_get("deprecated_at").map_err(backend)?;
Ok(Some(templates::DeprecationRecord {
deprecated_at: sql::parse_ts(&at),
deprecated_by: row.try_get("deprecated_by").map_err(backend)?,
reason: row.try_get("reason").map_err(backend)?,
}))
}
fn degraded(&self) -> bool {
false
}
}
impl $name {
async fn catalog_all_edges(
&self,
) -> Result<
Vec<$crate::serve::history::catalog::CatalogLineageEdge>,
$crate::serve::history::HistoryError,
> {
use sqlx::Row as _;
use $crate::serve::history::sql;
let backend = $crate::serve::history::sql::classify_backend_error;
let rows = sqlx::query(&self.stmts.catalog_select_edges)
.fetch_all(&self.pool)
.await
.map_err(backend)?;
let mut edges = Vec::with_capacity(rows.len());
for r in &rows {
let body: String = r.try_get("body").map_err(backend)?;
edges.push(sql::decode_json(&body, "catalog edge")?);
}
Ok(edges)
}
}
};
}
pub(crate) use impl_sql_history;
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn postgres_shard_statements_are_built() {
let s = Stmts::new(Dialect::Postgres);
assert!(s.insert_shard.contains("faucet_serve_shards"));
assert!(s.insert_shard.contains("ON CONFLICT"));
assert!(s.claim_shards_select.contains("JOIN faucet_serve_runs"));
assert!(s.claim_shard_one.contains("'running'"));
assert!(s.renew_shard_leases.contains("lease_expires_at"));
assert!(s.reclaim_shards_select.contains("'running'"));
assert!(s.reclaim_shard_requeue.contains("'pending'"));
assert!(s.reclaim_shard_fail.contains("'failed'"));
assert!(s.finalize_shard.contains("owner"));
assert!(s.shard_progress.contains("GROUP BY"));
}
#[test]
fn fmt_ts_is_fixed_width_and_sortable() {
let a = fmt_ts(
DateTime::parse_from_rfc3339("2026-01-01T00:00:00Z")
.unwrap()
.to_utc(),
);
let b = fmt_ts(
DateTime::parse_from_rfc3339("2026-01-01T00:00:01Z")
.unwrap()
.to_utc(),
);
assert!(a.ends_with('Z'));
assert_eq!(a.len(), b.len(), "fixed width");
assert!(a < b, "lexicographic order matches chronological order");
}
#[test]
fn is_expired_respects_window() {
let now = Utc::now();
let old = fmt_ts(now - chrono::Duration::seconds(120));
assert!(is_expired(&old, now, Duration::from_secs(60)));
assert!(!is_expired(&old, now, Duration::from_secs(600)));
assert!(!is_expired("not-a-timestamp", now, Duration::ZERO));
}
#[test]
fn parse_status_round_trips_known_and_defaults_failed() {
for s in [
RunStatus::Queued,
RunStatus::Pending,
RunStatus::Running,
RunStatus::Completed,
RunStatus::Failed,
RunStatus::Cancelled,
] {
assert_eq!(parse_status(s.as_str()), s);
}
assert_eq!(parse_status("garbage"), RunStatus::Failed);
}
#[test]
fn body_round_trips() {
let rec = RunRecord::queued(
"r1".into(),
Some("n".into()),
Default::default(),
Some("idem".into()),
Utc::now(),
);
let encoded = encode_body(&rec).unwrap();
let decoded = decode_body(&encoded).unwrap();
assert_eq!(decoded.run_id, "r1");
assert_eq!(decoded.idempotency_key.as_deref(), Some("idem"));
}
#[derive(Debug)]
struct FakeDbError(&'static str);
impl std::fmt::Display for FakeDbError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "db error {}", self.0)
}
}
impl std::error::Error for FakeDbError {}
impl sqlx::error::DatabaseError for FakeDbError {
fn message(&self) -> &str {
"db error"
}
fn code(&self) -> Option<std::borrow::Cow<'_, str>> {
Some(std::borrow::Cow::Borrowed(self.0))
}
fn as_error(&self) -> &(dyn std::error::Error + Send + Sync + 'static) {
self
}
fn as_error_mut(&mut self) -> &mut (dyn std::error::Error + Send + Sync + 'static) {
self
}
fn into_error(self: Box<Self>) -> Box<dyn std::error::Error + Send + Sync + 'static> {
self
}
fn kind(&self) -> sqlx::error::ErrorKind {
sqlx::error::ErrorKind::Other
}
}
fn db(code: &'static str) -> sqlx::Error {
sqlx::Error::Database(Box::new(FakeDbError(code)))
}
#[test]
fn sqlite_busy_and_locked_codes_are_transient_including_extended_forms() {
for code in ["5", "6", "517", "262", "261"] {
assert_eq!(
transience_of(&db(code)),
Transience::Transient,
"sqlite code {code}"
);
}
for code in ["1", "19", "2067"] {
assert_eq!(
transience_of(&db(code)),
Transience::Permanent,
"sqlite code {code}"
);
}
}
#[test]
fn postgres_sqlstates_split_transient_from_permanent() {
for code in [
"57P03", "53300", "40001", "40P01", "08006", "08000", "08003",
] {
assert_eq!(
transience_of(&db(code)),
Transience::Transient,
"sqlstate {code}"
);
}
for code in ["42601", "42P01", "23505"] {
assert_eq!(
transience_of(&db(code)),
Transience::Permanent,
"sqlstate {code}"
);
}
}
#[test]
fn non_database_variants_classify_from_the_variant() {
assert_eq!(
transience_of(&sqlx::Error::PoolTimedOut),
Transience::Transient
);
assert_eq!(
transience_of(&sqlx::Error::Io(std::io::Error::new(
std::io::ErrorKind::ConnectionReset,
"reset"
))),
Transience::Transient
);
assert_eq!(
transience_of(&sqlx::Error::RowNotFound),
Transience::Permanent
);
assert_eq!(
transience_of(&sqlx::Error::ColumnNotFound("nope".into())),
Transience::Permanent
);
assert_eq!(
transience_of(&sqlx::Error::WorkerCrashed),
Transience::Unknown
);
assert_eq!(transience_of_code(None), Transience::Unknown);
assert_eq!(transience_of_code(Some("nonsense")), Transience::Unknown);
}
#[test]
fn code_length_separates_the_two_dialects() {
assert_eq!(transience_of_code(Some("53300")), Transience::Transient);
assert_eq!(transience_of_code(Some("517")), Transience::Transient);
}
#[test]
fn classify_backend_error_carries_the_verdict_and_keeps_the_message() {
let err = classify_backend_error(db("57P03"));
assert_eq!(err.transience(), Transience::Transient);
assert!(err.is_retriable());
assert!(
err.to_string().starts_with("run-history backend error:"),
"display is unchanged for users: {err}"
);
let err = classify_backend_error(db("42601"));
assert_eq!(err.transience(), Transience::Permanent);
assert!(
!err.is_retriable(),
"a syntax error must not be retried at the connect gate"
);
}
#[test]
fn context_variant_names_the_phase_and_keeps_the_verdict() {
let err = classify_backend_error_with_context("creating run-history schema", db("5"));
assert_eq!(err.transience(), Transience::Transient);
assert!(
err.to_string()
.contains("creating run-history schema: error returned from database"),
"the phase that failed must survive classification: {err}"
);
}
#[test]
fn postgres_and_sqlite_statements_differ_only_in_placeholders() {
let pg = Stmts::new(Dialect::Postgres);
let lite = Stmts::new(Dialect::Sqlite);
assert!(pg.upsert.contains("$1") && lite.upsert.contains('?'));
assert!(pg.list.contains("$13") && lite.list.contains('?'));
assert!(pg.insert_idem.contains("ON CONFLICT (key) DO NOTHING"));
assert!(lite.insert_idem.contains("ON CONFLICT (key) DO NOTHING"));
assert!(pg.claim_one.contains("$3") && lite.claim_one.contains('?'));
assert!(pg.heartbeat_instance.contains("faucet_serve_instances"));
assert!(lite.heartbeat_instance.contains("faucet_serve_instances"));
}
}