Merge pull request #573 from fabro-sh/sqlite-storage-migration

Migrate secrets and automations storage to SQLite
This commit is contained in:
Scott Werner 2026-07-22 13:53:44 -04:00 • committed by GitHub
commit a71bd1ad33
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
37 changed files with 1898 additions and 634 deletions

5
Cargo.lock generated
View file

@ -2350,10 +2350,14 @@ dependencies = [
name = "fabro-automation"
version = "0.302.0-nightly.1"
dependencies = [
"anyhow",
"chrono",
"croner",
"fabro-db",
"hex",
"serde",
"sha2 0.10.9",
"sqlx",
"tempfile",
"thiserror 2.0.18",
"tokio",
@ -3047,6 +3051,7 @@ dependencies = [
"serde_json",
"serde_yaml",
"sha2 0.10.9",
"sqlx",
"strum 0.28.0",
"sysinfo",
"tempfile",

View file

@ -0,0 +1,318 @@
# Parallel fan-out / fan-in strategy
Status: proposed (not yet implemented). Line numbers reference the tree at the
time of writing and will drift.
This document specifies the intended product behavior of parallel fan-out
(`shape=component`) and fan-in (`shape=tripleoctagon`). Appendix A catalogs
pre-existing bugs this design resolves. Appendix B lists behavior changes
relative to today's implementation.
## 1. Introduction: goals and current weaknesses
The goal of this design is a parallel execution model that is **simple**,
**coherent**, and **correct**:
- **Simple** — a user should be able to predict what a fan-out/fan-in does
from the graph alone. One mental model ("branches produce candidates; fan-in
picks a commit; downstream sees all responses"), no new syntax for
synthesis, no incantations (special fidelity settings, escape-hatch
attributes) to make the flagship patterns work.
- **Coherent** — the same rules apply regardless of node type or channel.
What holds for a sequential command node's output should hold for a branch
command node's output; the workspace (git) facet and the text (context)
facet should follow parallel logic; selection should live in exactly one
place.
- **Correct** — the engine must do what the graph, the docs, and the recorded
run state say it does. Nodes drawn in the graph must run; documented outputs
must exist; a judge asked to pick the best candidate must be shown the
candidates.
Today's implementation misses all three. Issue #490 is the visible symptom:
a synthesis node after fan-in never sees branch responses — only
`{id, status, head_sha}` metadata — while two tutorials promise the opposite.
Investigation showed a broader incoherence:
- **Inherited spec gap.** The Attractor spec (fabro's ancestor) deliberately
isolates branch context and never defines a channel for branch output text.
Its fan-in pseudocode judges candidates it cannot see (`llm_evaluate`
receives only statuses) and sorts by a `score` field nothing sets. Fabro
inherited this gap faithfully.
- **Missing compensation.** Kilroy (a sibling Attractor implementation)
compensates with a file/git handoff: the post-merge node's prompt is
injected with each branch's `worktree_dir`, `logs_root`, and `head_sha` plus
instructions to read/merge them. Fabro has no equivalent, but its tutorials
promise kilroy-like behavior ("the synth node receives all four perspectives
in its preamble").
- **Accidental semantics.** Several adjacent behaviors are unprincipled
accidents rather than decisions: branches execute a single node and silently
skip chained nodes; two handlers fast-forward to two different definitions
of "winner"; branch nodes run with a stale preamble built for the fan-out
node; selection is vacuous in both modes; nested fan-outs silently lose git
isolation. Appendix A catalogs these.
## 2. Core model: candidates
A parallel branch is an isolated unit of execution that produces a
**candidate**. A candidate has exactly three facets:
| Facet | Content | Carried by |
|---|---|---|
| Commit | Workspace state produced by the branch | Per-branch git commit (`head_sha`) |
| Response | Text output (LLM response, command output) | Branch stage records + context keys |
| Verdict | Terminal status plus optional numeric `score` | `parallel.results` entries |
Fan-out produces N candidates in isolation. Fan-in selects **one commit** to
continue on. Downstream nodes (synthesis) get **all responses**. Selection and
synthesis are distinct concerns: selection needs a node (`tripleoctagon`);
synthesis is any ordinary downstream node, because responses propagate.
The post-fan-in contract, in one sentence: **after fan-in, the run looks as if
every branch had run sequentially, and the winner ran last.**
## 3. Topology: branches are subgraphs
Each outgoing edge of a fan-out node starts a branch. A branch executes as a
subgraph walk: the engine traverses nodes and edges from the branch entry node
until it reaches the join node (the fan-in). Multi-node chains
(`fork -> plan_a; plan_a -> impl_a; impl_a -> merge`) run every node in the
chain. This matches the Attractor spec (`execute_subgraph`) and kilroy
(`runSubgraphUntil`); today fabro executes exactly one node per branch and
silently skips the rest (Appendix A.2).
Structural validation (lint):
- Every path from a branch entry must converge on the run's join node.
A branch path that escapes the join (reaches exit, or a node outside the
fan-out region) is a validation error.
- A nested `component` node inside a branch is rejected until
worktree-from-worktree isolation is implemented (today it silently runs with
no git isolation; Appendix A.6).
- `fidelity="full"` on a fan-in's outgoing edge gets a lint warning: branches
run on different threads, so full fidelity can never carry branch outputs
(it drops the preamble entirely).
## 4. Node types in branches
No type-based restrictions. Restriction is structural (§3), not by allowlist:
- **Agent / prompt nodes** — the primary case.
- **Command / script / tool nodes** — first-class. Deterministic fan-out
(test matrices, benchmark bake-offs across worktrees) is a supported pattern
with no LLM anywhere: branches emit `score` via status fields, heuristic
selection picks the winner by measurement.
- **Conditionals** — meaningful under subgraph branches: they route *within*
the branch.
- **Human gates** — allowed; each branch may pause independently.
- **Nested parallel** — rejected by lint until isolation composes (§3).
## 5. Git isolation
Unchanged mechanics, with ownership fixed:
1. Before fan-out, checkpoint the sandbox to produce `base_sha`.
2. Each branch gets a worktree on a branch ref
(`fabro/run/parallel/<run>/<node>/pass<N>/<branch>`), rooted at `base_sha`.
The branch's `internal.work_dir` points at the worktree.
3. After a branch completes, `git add -A` + commit (`--allow-empty`) yields the
candidate's `head_sha`.
4. Worktrees are removed after the join. Loser branch refs are **kept** so
downstream nodes and humans can `git show`/`git diff` any candidate.
5. **Fan-in exclusively owns the fast-forward.** The parallel handler performs
no merge. After selection, fan-in fast-forwards the primary workspace to the
winner's `head_sha`. (Today both handlers fast-forward, to potentially
different winners; Appendix A.3.)
Degradation without git (no repo, or git isolation disabled): branches share
the primary sandbox with no workspace isolation, `head_sha` is absent from
candidates, and fan-in performs no merge. Response and verdict facets work
unchanged — prompt-only ensembles do not require git.
## 6. Execution and stage recording
Branch nodes execute as **real stages**, recorded through the normal
`ExecutionState::record` path and namespaced under the fan-out
(e.g. stage `a@1` within `fork@1`). Consequences (all fixes to current
behavior):
- Branch prompts/responses appear in events, `fabro dump`, and the web UI as
ordinary stages.
- Each branch node gets a **freshly built preamble** for its own position, via
the standard lifecycle, instead of reusing the fan-out node's stale preamble.
- Branch stages participate in the standard retry policy per node.
## 7. Context merge-back at fan-in
When fan-in completes, branch context updates are applied to the parent
context with a collision rule:
- **Per-node keys** (`response.<branch_node_id>`, structured-output fields
namespaced by node) apply for **all** branches. Branch node IDs are unique,
so no collisions.
- **Singleton keys** (`last_stage`, `last_response`, `command.output`,
un-namespaced status fields) are taken from the **winner only**.
- Failed branches' `response.<id>` values are still applied (a synthesis node
analyzing disagreement wants to see the failure text). Their singletons are
never applied.
Fan-in additionally writes (as today):
- `parallel.results` — one entry per candidate: `{id, status, head_sha?,
score?}`.
- `parallel.branch_count`, `parallel.fan_in.best_id`,
`parallel.fan_in.best_outcome`, `parallel.fan_in.best_head_sha`.
No file is materialized into the run workspace. `parallel.results` reaches LLM
consumers through the preamble's context section, and agents can `git show`
any candidate via its `head_sha`. (The current docs claim
`parallel_results.json` is available to downstream nodes; that claim is false
today and should be corrected rather than implemented — see open question 4
for the one consumer this leaves unserved.)
## 8. Selection
Fan-in selects the winning candidate. Two modes, as today, but with real
signal:
- **Heuristic** (no prompt on the fan-in node): rank by status
(succeeded < partially_succeeded < failed), then `score` descending, then
lexical id. `score` becomes settable: branches emit it via structured-output
/ status fields, which now survive into `parallel.results` (§7).
- **LLM judge** (fan-in node has a `prompt` and a backend is configured): the
judge prompt includes, per candidate: id, status, score, a bounded response
excerpt, and `git diff --stat` vs. `base_sha` when git isolation is active.
Today the judge sees only `{id, status, head_sha}` and cannot possibly
discriminate (Appendix A.4).
Selection determines the commit facet only. It does not suppress loser
responses (§7) or loser refs (§5).
## 9. Downstream visibility (preamble)
Prompt templates render once at manifest build time with `{goal, inputs}`
only; the preamble is the sole channel for runtime context into a
fresh-session node. Therefore:
- At **`compact`** (default) fidelity, branch stage summaries render their
responses **inline**, bounded, with a `See: <blob path>` reference when
truncated — the same treatment `command.output` already receives at compact.
Rationale: post-fan-in branch responses are unrecoverable through any
fidelity setting (different threads), exactly like command output.
- The per-branch response budget is larger than command output's 25-line tail
(ensemble analyses front-load their substance; a small tail amputates it).
Exact budget TBD at implementation; must remain bounded so N long branches
cannot blow the downstream context. Agent-type synthesis nodes can read the
full text from the blob/artifact reference.
- `summary:high` renders the same with its larger budget. `truncate` carries
goal only (explicit opt-out). `full` remains the degenerate case and lints
(§3).
## 10. Join policies
- `wait_all` (default): all branches run to completion; join proceeds when all
are terminal. Succeeds if no branch failed, else partially succeeds (fan-in
fails only when *all* candidates failed).
- `first_success`: join proceeds at the first successful branch. Remaining
branches are cancelled; cancelled branches record a terminal cancelled stage
(their partial responses are not merged back). The sole successful branch is
the winner.
- `k_of_n` / `quorum` (kilroy has them): deliberately **not** added now. The
surface stays minimal until a concrete need appears.
## Open questions
1. Implementation phasing: subgraph branches (§3) are the largest lift.
Response propagation (§6–§9) fixes #490 and both tutorials on its own and
can ship first.
2. `first_success` cancellation semantics for in-flight agent sessions
(graceful stop vs. abort; what the cancelled stage records).
3. Exact preamble budget per branch response (§9).
4. Context access for deterministic post-fan-in consumers. Command/script
nodes have no channel to context (no preamble, no template rendering in
`script`), so a deterministic aggregator after fan-in cannot learn
candidate `head_sha`s. Candidate mechanisms: a results file under a
checkpoint-excluded workspace path, or an env var (e.g.
`FABRO_PARALLEL_RESULTS`) pointing at a file outside the workspace. A bare
workspace file is ruled out: checkpoint commits `git add -A`, so it would
leak into run history and PRs. Design alongside the deterministic fan-out
pattern (§4).
---
## Appendix A: pre-existing bugs
Cataloged against the current tree; line numbers will drift.
1. **Branch outputs dropped (#490).** The fan-out task reads only
`outcome.status` and `head_sha` from each branch; branch `context_updates`
(including `response.<id>`) are discarded with the forked context
(`fabro-workflow/src/handler/parallel.rs:389-397,465-470`). No channel
carries branch text to downstream nodes. Confirmed by live repro on
fabro-testing (run `01KX148ZAMMJRAADHK1HBF7PC3`, server 0.287.0-nightly.0).
2. **Chained branch nodes silently skipped.** Branches execute exactly one
node; the engine then jumps to the join (`parallel.rs:388-397,612,628`;
`fabro-core/src/executor.rs:421`). In `fork -> a; a -> a2; a2 -> merge`,
`a2` never runs and nothing warns.
3. **Double fast-forward with two different winner definitions.** The parallel
handler fast-forwards the *lexically first* successful branch
(`parallel.rs:511-538`); fan-in then fast-forwards *its* selected winner
(`fan_in.rs:120-131`). If selection ever picks a non-lexical-first branch,
the second `--ff-only` merge cannot succeed (sibling commits diverge).
Masked today only because selection is vacuous (A.4).
4. **Selection is vacuous.** Heuristic tie-breaks on a `score` field nothing
can set (scores would arrive via branch context updates, which are dropped
per A.1). The LLM judge prompt is `serde_json::to_string_pretty` of
`parallel.results` — id/status/head_sha only (`fan_in.rs:247-250`). Both
modes reduce to "first successful branch, alphabetically."
5. **Branch preambles are stale.** Branches run via `dispatch_handler`,
bypassing the lifecycle's per-node preamble rebuild; each branch node
inherits `current.preamble` as computed for the fan-out node itself.
6. **Nested parallel silently loses git isolation.** Branch `EngineServices`
are built with `git_state: RwLock::new(None)` (`parallel.rs:380`), so a
`component` node inside a branch runs its own branches with no worktrees
and no warning.
7. **Docs contradict the engine.** `tutorials/ensemble.mdx:85` and
`tutorials/parallel-review.mdx:81` claim the post-merge node receives all
branch perspectives in its preamble (false, per A.1).
`workflows/stages-and-nodes.mdx:195` claims merged results are available to
downstream nodes as `parallel_results.json` (the file exists only under
`stages/<fork>@1/` in dumps, not in any node's working directory).
8. **`fidelity="full"` across a fan-in is a trap.** It drops the preamble
(metadata included) and cannot attach to any branch thread; raising
fidelity strictly reduces what the downstream node sees. No lint warns.
## Appendix B: behavior changes vs. today
Changes a user could observe if this spec is implemented as written.
1. **Chained branch nodes execute.** Graphs that (unknowingly) relied on
single-node branch semantics will now run the full chain (fixes A.2; may
lengthen existing runs).
2. **New validation errors.** Branch paths that don't converge on the join,
and nested `component` nodes inside branches, become lint failures for
graphs that previously ran (with wrong or silently degraded semantics).
3. **Fan-in owns the fast-forward.** The workspace after fan-in may land on a
different commit than today whenever selection (scores, LLM judge)
disagrees with lexical-first order. The parallel handler no longer merges.
4. **Post-fan-in context is richer.** `response.<id>` for every branch,
winner-sourced singletons (`last_stage`, `last_response`,
`command.output`), and `score` in `parallel.results`. Today those
singletons retain their pre-fork values; workflows with edge conditions
over them could route differently.
5. **Preambles after fan-in grow.** Branch responses render inline at
`compact` fidelity (bounded). Downstream nodes see more tokens per run;
snapshot tests over preambles will churn.
6. **Branch executions become visible stages.** Events, dumps, the web UI, and
the stage list gain per-branch stages (`a@1` under `fork@1`). Consumers of
`events.jsonl` / the API will see new stage records.
7. **`parallel_results.json` docs claim is corrected, not implemented.**
`workflows/stages-and-nodes.mdx:195` is updated to describe the real
channels (context key + preamble); no file appears in the workspace
(open question 4 covers deterministic consumers).
8. **Loser branch refs are documented as retained** and become part of the
product contract instead of an accident of not deleting them.
9. **`first_success` cancels losers explicitly** and records cancelled stages;
today's exact cancellation behavior is unspecified.
10. **LLM judge prompts change shape.** Fan-in nodes with prompts now send
candidate excerpts and diff stats to the judge — more tokens, different
(better) selections than today's id-only prompt.

View file

@ -0,0 +1,91 @@
# Automations SQLite schema design
## Outcome
Use two tables: `automations` owns the aggregate, public revision, and API enablement; `automation_triggers` stores schedule triggers keyed by automation and trigger ID. Keep scheduler cursors and executions out of this migration.
## Proposed migration
`lib/crates/fabro-db/migrations/2026071102_automations.sql` after the secrets migration.
```sql
CREATE TABLE automations (
id TEXT PRIMARY KEY NOT NULL,
revision TEXT NOT NULL,
name TEXT NOT NULL,
description TEXT,
api_enabled INTEGER NOT NULL,
target_repository TEXT NOT NULL,
target_ref TEXT NOT NULL,
target_workflow TEXT NOT NULL,
CHECK (length(id) BETWEEN 1 AND 63),
CHECK (substr(id, 1, 1) GLOB '[a-z0-9]'),
CHECK (id NOT GLOB '*[^a-z0-9-]*'),
CHECK (length(revision) = 64),
CHECK (revision NOT GLOB '*[^0-9a-f]*'),
CHECK (length(trim(name)) > 0),
CHECK (api_enabled IN (0, 1)),
CHECK (length(target_repository) BETWEEN 3 AND 140),
CHECK (length(target_ref) BETWEEN 1 AND 255),
CHECK (length(target_workflow) BETWEEN 1 AND 255)
);
CREATE TABLE automation_triggers (
automation_id TEXT NOT NULL,
id TEXT NOT NULL,
enabled INTEGER NOT NULL,
expression TEXT NOT NULL,
PRIMARY KEY (automation_id, id),
FOREIGN KEY (automation_id) REFERENCES automations(id) ON DELETE CASCADE,
CHECK (length(id) BETWEEN 1 AND 63),
CHECK (substr(id, 1, 1) GLOB '[a-z0-9]'),
CHECK (id NOT GLOB '*[^a-z0-9_-]*'),
CHECK (enabled IN (0, 1)),
CHECK (length(trim(expression)) > 0)
);
```
## Decisions
- Trigger order is non-semantic. Load schedule triggers by ID and hash that canonical order so equivalent definitions share a revision.
- Trigger IDs are unique only within an automation, matching current validation.
- `api_enabled` on `automations` models the manual/API capability, so the trigger table needs no API-row special case or partial unique index.
- `automation_triggers` stores schedule configuration only. Rust validates the five-field UTC cron grammar; SQLite enforces required schedule fields.
- `revision` remains the 64-character SHA-256 ETag for the complete aggregate, including API enablement and schedules in canonical trigger-ID order. Trigger rows do not have independent revisions.
- No timestamps: the current model and API expose none, and revisions already own concurrency.
- No JSON blob: typed columns give useful constraints and queries. A future trigger kind should get an explicit schema migration and the most natural relational shape for its cardinality.
- No schedule index yet. The scheduler reads all definitions and evaluates cron in Rust; the primary key `(automation_id, id)` supports deterministic child loading.
## Store behavior
- Make `AutomationStore` pool-backed with no process-wide cache. `list` and `get` become async and fallible so database failures never look like empty state.
- Load an automation and its schedules with one join ordered by trigger ID, then re-run the existing Rust validation when constructing domain values.
- Create the parent and schedule rows in one transaction.
- Replace validates and computes the aggregate revision first, conditionally updates the parent with `WHERE id = ? AND revision = ?`, then replaces all children in the same transaction.
- Delete uses `WHERE id = ? AND revision = ?`; child deletion cascades.
- On a zero-row replace/delete, query the current revision inside the transaction to distinguish `NotFound` from `StaleRevision`.
- Canonicalize API enablement and schedules sorted by trigger ID before computing the SHA-256 revision. During import, persist the raw-file revision so existing ETags do not change merely because storage moved.
## Legacy import
- Read every `automations/*.toml` beside the active `settings.toml`, ignoring non-TOML files as today.
- Parse and validate the full directory before opening the write transaction.
- Insert each aggregate transactionally. `ON CONFLICT(id) DO NOTHING`; an existing SQL automation wins as a whole, including its triggers.
- After commit, rename the directory to `automations.imported-<timestamp>.bak`.
- Missing directory is a no-op. Invalid input leaves the directory untouched. A rename failure is retry-safe because the next run skips already imported IDs and retries the rename.
## Scheduler boundary
This schema stores definitions only. `next_due_at`, `last_fired_at`, leases, and execution claims remain out of scope because the current scheduler intentionally keeps cursors in memory and skips missed occurrences across restarts. Multi-node exactly-once scheduling would require a separately designed durable claim/execution table, not extra mutable fields on definitions.
## Tests
- Schema constraints: IDs, revisions, booleans, API enablement, schedule fields, foreign key cascade, and duplicate schedule-trigger IDs.
- Store: sorted list, deterministic schedule order, empty schedule list, API enablement, create conflict, conditional replace/delete, two-pool visibility, and row revalidation.
- Atomicity: failed child insert leaves the old aggregate unchanged.
- Import: absent, success, SQL wins, invalid directory unchanged, backup rename, retry after destination commit.
- API/scheduler regression: ETag behavior unchanged; scheduler sees SQL changes and does not treat read failures as an empty list.
## Unresolved questions
- None for definition storage. Durable/multi-node scheduling remains a separate design.

View file

@ -7,7 +7,11 @@ An **automation** is a saved run configuration — a repository, ref, and workfl
## Defining automations
The server keeps one TOML file per automation in an `automations/` directory next to its `settings.toml`. Manage them in the web UI at `/automations`, through the `/api/v1/automations` REST API, or by editing the files directly:
The server stores automations in its SQLite database. Manage them in the web UI at `/automations` or through the `/api/v1/automations` REST API.
When upgrading from file-backed automation storage, startup imports every valid `automations/*.toml` file next to the active `settings.toml`. Existing SQLite definitions win on ID conflicts. After a successful import, Fabro renames the directory to a timestamped backup such as `automations.imported-20260711T180000000000Z.bak`. Invalid TOML leaves the original directory untouched for operator repair.
The legacy files use this shape:
```toml title="automations/nightly-release.toml"
name = "Nightly release"

View file

@ -13,15 +13,19 @@ doctest = false
workspace = true
[dependencies]
chrono.workspace = true
croner.workspace = true
fabro-db = { path = "../fabro-db" }
hex.workspace = true
serde.workspace = true
sha2.workspace = true
sqlx.workspace = true
thiserror.workspace = true
tokio.workspace = true
toml.workspace = true
tracing.workspace = true
[dev-dependencies]
anyhow.workspace = true
tempfile = "3"
tokio = { workspace = true, features = ["macros", "test-util"] }

View file

@ -0,0 +1,123 @@
//! Imports the pre-SQLite `automations/*.toml` directory into the relational
//! automation store. Remove this compatibility migration after 2026-10-11,
//! once supported upgrades no longer span a file-backed Fabro release.
use std::path::{Path, PathBuf};
use chrono::Utc;
use fabro_db::{DbPool, ImportReport};
use tokio::fs;
use tracing::info;
use crate::{Automation, AutomationId, AutomationStoreError, store};
pub(crate) const REMOVAL_DEADLINE: &str = "2026-10-11";
pub async fn import_legacy_directory_once(
pool: &DbPool,
source_dir: impl AsRef<Path>,
) -> Result<Option<ImportReport>, AutomationStoreError> {
let source_dir = source_dir.as_ref();
let Some(paths) = legacy_automation_paths(source_dir).await? else {
return Ok(None);
};
let mut automations = Vec::with_capacity(paths.len());
for (id, path) in paths {
let bytes = fs::read(&path)
.await
.map_err(|source| AutomationStoreError::io(&path, source))?;
automations.push(Automation::from_persisted_path(id, &bytes, path)?);
}
let mut transaction = pool.begin().await?;
let mut imported_ids = Vec::new();
let mut skipped_rows = 0;
for automation in &automations {
if store::insert_automation_ignoring_conflict(&mut transaction, automation).await? {
imported_ids.push(automation.id.to_string());
} else {
skipped_rows += 1;
}
}
transaction.commit().await?;
let backup_path = rename_imported_legacy_directory(source_dir).await?;
let report = ImportReport {
source_path: source_dir.to_path_buf(),
backup_path,
imported_rows: imported_ids.len(),
skipped_rows,
names: imported_ids,
};
info!(
source_path = %report.source_path.display(),
backup_path = %report.backup_path.display(),
imported_rows = report.imported_rows,
skipped_rows = report.skipped_rows,
automation_ids = ?report.names,
removal_deadline = REMOVAL_DEADLINE,
"Imported legacy automations directory into SQLite"
);
Ok(Some(report))
}
async fn legacy_automation_paths(
source_dir: &Path,
) -> Result<Option<Vec<(AutomationId, PathBuf)>>, AutomationStoreError> {
let mut entries = match fs::read_dir(source_dir).await {
Ok(entries) => entries,
Err(source) if source.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(source) => return Err(AutomationStoreError::io(source_dir, source)),
};
let mut paths = Vec::new();
while let Some(entry) = entries
.next_entry()
.await
.map_err(|source| AutomationStoreError::io(source_dir, source))?
{
let path = entry.path();
let file_type = entry
.file_type()
.await
.map_err(|source| AutomationStoreError::io(&path, source))?;
if file_type.is_file() && is_toml_file(&path) {
paths.push((id_from_path(&path)?, path));
}
}
paths.sort_by(|left, right| left.1.cmp(&right.1));
Ok(Some(paths))
}
fn id_from_path(path: &Path) -> Result<AutomationId, AutomationStoreError> {
let stem = path
.file_stem()
.and_then(|stem| stem.to_str())
.ok_or_else(|| AutomationStoreError::InvalidFilename {
path: path.to_path_buf(),
reason: "filename is not valid UTF-8".to_string(),
})?;
AutomationId::new(stem).map_err(|source| AutomationStoreError::InvalidFilename {
path: path.to_path_buf(),
reason: source.to_string(),
})
}
fn is_toml_file(path: &Path) -> bool {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| extension == "toml")
}
async fn rename_imported_legacy_directory(
source_dir: &Path,
) -> Result<PathBuf, AutomationStoreError> {
let backup_path = fabro_db::legacy_backup_path(source_dir, "automations", Utc::now());
fs::rename(source_dir, &backup_path)
.await
.map_err(|source| AutomationStoreError::LegacyBackup {
source_path: source_dir.to_path_buf(),
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}

View file

@ -4,7 +4,7 @@ use croner::errors::CronError;
use toml::de::Error as TomlDeError;
use toml::ser::Error as TomlSerError;
use crate::{AutomationId, AutomationRevision};
use crate::{AutomationId, AutomationRevision, AutomationRevisionParseError};
#[derive(Debug, thiserror::Error)]
pub enum AutomationValidationError {
@ -57,6 +57,31 @@ pub enum AutomationStoreError {
#[from]
source: AutomationValidationError,
},
#[error("stored automation {id} failed validation")]
StoredValidation {
id: AutomationId,
#[source]
source: AutomationValidationError,
},
#[error("stored automation has invalid id {value:?}")]
StoredId {
value: String,
#[source]
source: AutomationValidationError,
},
#[error("stored automation {id} has an invalid trigger row")]
StoredTriggerShape { id: AutomationId },
#[error("stored automation {id} has an invalid revision")]
InvalidRevision {
id: AutomationId,
#[source]
source: AutomationRevisionParseError,
},
#[error("database error")]
Db {
#[from]
source: sqlx::Error,
},
#[error("invalid automation filename at {path:?}")]
InvalidFilename { path: PathBuf, reason: String },
#[error("failed to parse automation TOML at {path:?}")]
@ -82,6 +107,13 @@ pub enum AutomationStoreError {
#[source]
source: std::io::Error,
},
#[error("renaming legacy automations directory {source_path:?} to {backup_path:?}")]
LegacyBackup {
source_path: PathBuf,
backup_path: PathBuf,
#[source]
source: std::io::Error,
},
}
impl AutomationStoreError {
@ -114,10 +146,16 @@ impl AutomationStoreError {
Self::MissingRevision { .. } => "missing_revision",
Self::StaleRevision { .. } => "stale_revision",
Self::Validation { .. } => "validation",
Self::StoredValidation { .. } => "stored_validation",
Self::StoredId { .. } => "stored_id",
Self::StoredTriggerShape { .. } => "stored_trigger_shape",
Self::InvalidRevision { .. } => "invalid_revision",
Self::Db { .. } => "db",
Self::InvalidFilename { .. } => "invalid_filename",
Self::Parse { .. } | Self::InvalidUtf8 { .. } => "parse",
Self::Serialize { .. } => "serialize",
Self::Io { .. } => "io",
Self::LegacyBackup { .. } => "legacy_backup",
}
}
}

View file

@ -1,10 +1,12 @@
mod error;
mod id;
mod migrations;
mod model;
mod store;
pub use error::{AutomationStoreError, AutomationValidationError};
pub use id::{AutomationId, AutomationRevision, AutomationRevisionParseError, AutomationTriggerId};
pub use migrations::{ImportReport, import_legacy_directory_once};
pub use model::{
ApiTrigger, Automation, AutomationDraft, AutomationReplace, AutomationTarget,
AutomationTrigger, GitHubRepositorySlug, ScheduleTrigger, parse_github_repository_slug,

View file

@ -0,0 +1,5 @@
#[path = "../migrations/2026071101_file_definitions_to_sqlite.rs"]
mod file_definitions_to_sqlite;
pub use fabro_db::ImportReport;
pub use file_definitions_to_sqlite::import_legacy_directory_once;

View file

@ -21,6 +21,8 @@ static SCHEDULE_CRON_PARSER: LazyLock<CronParser> = LazyLock::new(|| {
.build()
});
pub(crate) const MANUAL_TRIGGER_ID: &str = "manual";
/// Parse an automation schedule trigger expression with the canonical
/// configuration (no seconds, no year). Returned `Cron` instances can be cached
/// and used to find next occurrences.
@ -61,7 +63,7 @@ impl Automation {
id: AutomationId,
draft: AutomationReplace,
) -> Result<(Self, Vec<u8>), AutomationStoreError> {
validate_fields(&draft)?;
let draft = normalize_replace(draft)?;
let persisted = PersistedAutomation::from(draft.clone());
let bytes = canonical_bytes(&persisted)?;
let revision = AutomationRevision::from_bytes(&bytes);
@ -69,6 +71,15 @@ impl Automation {
Ok((automation, bytes))
}
pub(crate) fn from_stored(
id: AutomationId,
revision: AutomationRevision,
value: AutomationReplace,
) -> Result<Self, AutomationValidationError> {
let value = normalize_replace(value)?;
Ok(Self::from_validated_replace(id, revision, value))
}
pub(crate) fn to_persisted(&self) -> PersistedAutomation {
PersistedAutomation {
name: self.name.clone(),
@ -102,13 +113,23 @@ impl Automation {
})
}
pub(crate) fn schedule_triggers(&self) -> impl Iterator<Item = &ScheduleTrigger> {
self.triggers.iter().filter_map(|trigger| match trigger {
AutomationTrigger::Schedule(trigger) => Some(trigger),
AutomationTrigger::Api(_) => None,
})
}
pub(crate) fn api_enabled(&self) -> bool {
self.enabled_api_trigger().is_some()
}
fn from_persisted(
id: AutomationId,
revision: AutomationRevision,
persisted: PersistedAutomation,
) -> Result<Self, AutomationValidationError> {
let replace = AutomationReplace::from(persisted);
validate_fields(&replace)?;
let replace = normalize_replace(AutomationReplace::from(persisted))?;
Ok(Self::from_validated_replace(id, revision, replace))
}
@ -187,6 +208,18 @@ pub struct ApiTrigger {
pub enabled: bool,
}
impl ApiTrigger {
/// The canonical enabled API trigger. Automations store API enablement as a
/// flag and re-materialize it as this trigger with the fixed `manual` id.
pub(crate) fn manual() -> Self {
Self {
id: AutomationTriggerId::new(MANUAL_TRIGGER_ID)
.expect("manual automation trigger id is valid"),
enabled: true,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct ScheduleTrigger {
@ -291,6 +324,46 @@ fn validate_fields(value: &AutomationReplace) -> Result<(), AutomationValidation
validate_triggers(&value.triggers)
}
fn normalize_replace(
mut value: AutomationReplace,
) -> Result<AutomationReplace, AutomationValidationError> {
validate_fields(&value)?;
let api_enabled = value
.triggers
.iter()
.any(|trigger| matches!(trigger, AutomationTrigger::Api(trigger) if trigger.enabled));
let mut schedules = value
.triggers
.into_iter()
.filter_map(|trigger| match trigger {
AutomationTrigger::Schedule(trigger) => Some(trigger),
AutomationTrigger::Api(_) => None,
})
.collect::<Vec<_>>();
schedules.sort_by(|left, right| left.id.cmp(&right.id));
// Canonicalization renames the enabled API trigger to `manual`, which can
// collide with a schedule trigger id even when the input ids were unique.
if api_enabled
&& schedules
.iter()
.any(|schedule| schedule.id.as_str() == MANUAL_TRIGGER_ID)
{
return Err(AutomationValidationError::DuplicateTriggerId {
id: MANUAL_TRIGGER_ID.to_string(),
});
}
let mut triggers = Vec::with_capacity(schedules.len() + usize::from(api_enabled));
if api_enabled {
triggers.push(AutomationTrigger::Api(ApiTrigger::manual()));
}
triggers.extend(schedules.into_iter().map(AutomationTrigger::Schedule));
value.triggers = triggers;
Ok(value)
}
pub fn parse_github_repository_slug(
value: &str,
) -> Result<GitHubRepositorySlug, AutomationValidationError> {

View file

@ -1,65 +1,88 @@
use std::collections::HashMap;
use std::io::ErrorKind;
use std::path::{Path, PathBuf};
use std::time::{SystemTime, UNIX_EPOCH};
use std::str::FromStr as _;
use tokio::fs;
use tokio::io::AsyncWriteExt as _;
use tokio::sync::{Mutex, RwLock};
use fabro_db::DbPool;
use sqlx::sqlite::SqliteRow;
use sqlx::{Row as _, Sqlite, Transaction};
use crate::{
Automation, AutomationDraft, AutomationId, AutomationReplace, AutomationRevision,
AutomationStoreError,
ApiTrigger, Automation, AutomationDraft, AutomationId, AutomationReplace, AutomationRevision,
AutomationStoreError, AutomationTarget, AutomationTrigger, AutomationTriggerId,
ScheduleTrigger,
};
#[derive(Debug)]
/// Shared projection for loading automations with their schedule triggers.
/// A macro rather than a `const` because sqlx requires `&'static str` SQL.
macro_rules! select_automations_sql {
($suffix:expr) => {
concat!(
"SELECT
a.id,
a.revision,
a.name,
a.description,
a.api_enabled,
a.target_repository,
a.target_ref,
a.target_workflow,
t.id AS trigger_id,
t.enabled AS trigger_enabled,
t.expression AS trigger_expression
FROM automations AS a
LEFT JOIN automation_triggers AS t ON t.automation_id = a.id
",
$suffix
)
};
}
#[derive(Clone)]
pub struct AutomationStore {
dir: PathBuf,
mutations: Mutex<()>,
automations: RwLock<HashMap<AutomationId, Automation>>,
pool: DbPool,
}
impl std::fmt::Debug for AutomationStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AutomationStore").finish_non_exhaustive()
}
}
impl AutomationStore {
/// Synchronously load every persisted automation in `dir`. Returns an error
/// if any file fails to parse or validate; the caller decides startup
/// failure policy. Synchronous because it runs once at construction time
/// (typically during server startup) and is invoked from non-async code.
pub fn load(dir: impl Into<PathBuf>) -> Result<Self, AutomationStoreError> {
let dir = dir.into();
let automations = load_automations(&dir)?;
Ok(Self {
dir,
mutations: Mutex::new(()),
automations: RwLock::new(automations),
})
#[must_use]
pub fn new(pool: DbPool) -> Self {
Self { pool }
}
pub async fn list(&self) -> Vec<Automation> {
let automations = self.automations.read().await;
let mut values = automations.values().cloned().collect::<Vec<_>>();
values.sort_by(|left, right| left.id.cmp(&right.id));
values
pub async fn list(&self) -> Result<Vec<Automation>, AutomationStoreError> {
let rows = sqlx::query(select_automations_sql!("ORDER BY a.id, t.id"))
.fetch_all(&self.pool)
.await?;
automations_from_rows(&rows)
}
pub async fn get(&self, id: &AutomationId) -> Option<Automation> {
self.automations.read().await.get(id).cloned()
pub async fn get(&self, id: &AutomationId) -> Result<Option<Automation>, AutomationStoreError> {
let rows = sqlx::query(select_automations_sql!("WHERE a.id = ? ORDER BY t.id"))
.bind(id.as_str())
.fetch_all(&self.pool)
.await?;
Ok(automations_from_rows(&rows)?.into_iter().next())
}
pub async fn exists(&self, id: &AutomationId) -> Result<bool, AutomationStoreError> {
let row = sqlx::query("SELECT 1 FROM automations WHERE id = ?")
.bind(id.as_str())
.fetch_optional(&self.pool)
.await?;
Ok(row.is_some())
}
pub async fn create(&self, draft: AutomationDraft) -> Result<Automation, AutomationStoreError> {
let (id, replace) = draft.into();
let (automation, bytes) = Automation::from_replace(id.clone(), replace)?;
let _mutation = self.mutations.lock().await;
if self.automations.read().await.contains_key(&id) {
let (automation, _) = Automation::from_replace(id.clone(), replace)?;
let mut transaction = self.pool.begin().await?;
if !insert_automation_ignoring_conflict(&mut transaction, &automation).await? {
return Err(AutomationStoreError::AlreadyExists { id });
}
let path = automation_path(&self.dir, &id);
write_new(&self.dir, &path, &bytes)
.await
.map_err(|err| create_error_for(id.clone(), err))?;
let mut automations = self.automations.write().await;
automations.insert(id, automation.clone());
transaction.commit().await?;
Ok(automation)
}
@ -69,25 +92,42 @@ impl AutomationStore {
expected: &AutomationRevision,
draft: AutomationReplace,
) -> Result<Automation, AutomationStoreError> {
let (automation, bytes) = Automation::from_replace(id.clone(), draft)?;
let _mutation = self.mutations.lock().await;
{
let automations = self.automations.read().await;
let current = automations
.get(id)
.ok_or_else(|| AutomationStoreError::NotFound { id: id.clone() })?;
if &current.revision != expected {
return Err(AutomationStoreError::StaleRevision {
id: id.clone(),
expected: expected.clone(),
actual: current.revision.clone(),
});
}
let (automation, _) = Automation::from_replace(id.clone(), draft)?;
let mut transaction = self.pool.begin().await?;
let result = sqlx::query(
r"
UPDATE automations SET
revision = ?,
name = ?,
description = ?,
api_enabled = ?,
target_repository = ?,
target_ref = ?,
target_workflow = ?
WHERE id = ? AND revision = ?
",
)
.bind(automation.revision.as_str())
.bind(&automation.name)
.bind(automation.description.as_deref())
.bind(automation.api_enabled())
.bind(&automation.target.repository)
.bind(&automation.target.ref_selector)
.bind(&automation.target.workflow)
.bind(id.as_str())
.bind(expected.as_str())
.execute(&mut *transaction)
.await?;
if result.rows_affected() == 0 {
return Err(revision_mismatch_error(&mut transaction, id, expected).await?);
}
write_atomic(&self.dir, &automation_path(&self.dir, id), &bytes).await?;
let mut automations = self.automations.write().await;
automations.insert(id.clone(), automation.clone());
sqlx::query("DELETE FROM automation_triggers WHERE automation_id = ?")
.bind(id.as_str())
.execute(&mut *transaction)
.await?;
insert_schedule_triggers(&mut transaction, &automation).await?;
transaction.commit().await?;
Ok(automation)
}
@ -96,333 +136,226 @@ impl AutomationStore {
id: &AutomationId,
expected: &AutomationRevision,
) -> Result<(), AutomationStoreError> {
let _mutation = self.mutations.lock().await;
{
let automations = self.automations.read().await;
let current = automations
.get(id)
.ok_or_else(|| AutomationStoreError::NotFound { id: id.clone() })?;
if &current.revision != expected {
return Err(AutomationStoreError::StaleRevision {
id: id.clone(),
expected: expected.clone(),
actual: current.revision.clone(),
});
}
let mut transaction = self.pool.begin().await?;
let result = sqlx::query("DELETE FROM automations WHERE id = ? AND revision = ?")
.bind(id.as_str())
.bind(expected.as_str())
.execute(&mut *transaction)
.await?;
if result.rows_affected() == 0 {
return Err(revision_mismatch_error(&mut transaction, id, expected).await?);
}
let path = automation_path(&self.dir, id);
fs::remove_file(&path)
.await
.map_err(|err| AutomationStoreError::io(path, err))?;
let mut automations = self.automations.write().await;
automations.remove(id);
transaction.commit().await?;
Ok(())
}
}
#[expect(
clippy::disallowed_methods,
reason = "Automation directory scan runs once at startup, before the runtime needs to make progress; std::fs avoids needing a Tokio runtime for the caller."
)]
fn load_automations(dir: &Path) -> Result<HashMap<AutomationId, Automation>, AutomationStoreError> {
let entries = match std::fs::read_dir(dir) {
Ok(entries) => entries,
Err(err) if err.kind() == ErrorKind::NotFound => return Ok(HashMap::new()),
Err(err) => return Err(AutomationStoreError::io(dir, err)),
};
struct StoredAutomation {
id: AutomationId,
revision: AutomationRevision,
name: String,
description: Option<String>,
api_enabled: bool,
target: AutomationTarget,
schedule_triggers: Vec<ScheduleTrigger>,
}
let mut automations = HashMap::new();
for entry in entries {
let entry = entry.map_err(|err| AutomationStoreError::io(dir, err))?;
let path = entry.path();
let file_type = entry
.file_type()
.map_err(|err| AutomationStoreError::io(&path, err))?;
if !file_type.is_file() || !is_toml_file(&path) {
continue;
impl StoredAutomation {
fn from_row(row: &SqliteRow) -> Result<Self, AutomationStoreError> {
let id_value = row.try_get::<String, _>("id")?;
let id = AutomationId::new(id_value.clone()).map_err(|source| {
AutomationStoreError::StoredId {
value: id_value,
source,
}
})?;
let revision = AutomationRevision::from_str(&row.try_get::<String, _>("revision")?)
.map_err(|source| AutomationStoreError::InvalidRevision {
id: id.clone(),
source,
})?;
Ok(Self {
id,
revision,
name: row.try_get("name")?,
description: row.try_get("description")?,
api_enabled: row.try_get("api_enabled")?,
target: AutomationTarget {
repository: row.try_get("target_repository")?,
ref_selector: row.try_get("target_ref")?,
workflow: row.try_get("target_workflow")?,
},
schedule_triggers: Vec::new(),
})
}
fn push_trigger_row(&mut self, row: &SqliteRow) -> Result<(), AutomationStoreError> {
let Some(id_value) = row.try_get::<Option<String>, _>("trigger_id")? else {
return Ok(());
};
let id = AutomationTriggerId::new(id_value).map_err(|source| {
AutomationStoreError::StoredValidation {
id: self.id.clone(),
source,
}
})?;
self.schedule_triggers.push(ScheduleTrigger {
id,
enabled: row
.try_get::<Option<bool>, _>("trigger_enabled")?
.ok_or_else(|| AutomationStoreError::StoredTriggerShape {
id: self.id.clone(),
})?,
expression: row
.try_get::<Option<String>, _>("trigger_expression")?
.ok_or_else(|| AutomationStoreError::StoredTriggerShape {
id: self.id.clone(),
})?,
});
Ok(())
}
fn finish(self) -> Result<Automation, AutomationStoreError> {
// `from_stored` canonicalizes trigger order and the manual API trigger.
let mut triggers = self
.schedule_triggers
.into_iter()
.map(AutomationTrigger::Schedule)
.collect::<Vec<_>>();
if self.api_enabled {
triggers.push(AutomationTrigger::Api(ApiTrigger::manual()));
}
let automation = load_automation_file(&path)?;
automations.insert(automation.id.clone(), automation);
let id = self.id;
Automation::from_stored(id.clone(), self.revision, AutomationReplace {
name: self.name,
description: self.description,
target: self.target,
triggers,
})
.map_err(|source| AutomationStoreError::StoredValidation { id, source })
}
}
fn automations_from_rows(rows: &[SqliteRow]) -> Result<Vec<Automation>, AutomationStoreError> {
let mut automations = Vec::new();
let mut current: Option<StoredAutomation> = None;
for row in rows {
let row_id = row.try_get::<String, _>("id")?;
if current
.as_ref()
.is_some_and(|automation| automation.id.as_str() != row_id)
{
automations.push(
current
.take()
.expect("current automation exists")
.finish()?,
);
}
if current.is_none() {
current = Some(StoredAutomation::from_row(row)?);
}
current
.as_mut()
.expect("current automation exists")
.push_trigger_row(row)?;
}
if let Some(automation) = current {
automations.push(automation.finish()?);
}
Ok(automations)
}
#[expect(
clippy::disallowed_methods,
reason = "Sync sibling of `load_automations`; only invoked from the synchronous startup load path."
)]
fn load_automation_file(path: &Path) -> Result<Automation, AutomationStoreError> {
let id = id_from_path(path)?;
let bytes = std::fs::read(path).map_err(|err| AutomationStoreError::io(path, err))?;
Automation::from_persisted_path(id, &bytes, path)
pub(crate) async fn insert_automation_ignoring_conflict(
transaction: &mut Transaction<'_, Sqlite>,
automation: &Automation,
) -> Result<bool, AutomationStoreError> {
let result = sqlx::query(
r"
INSERT INTO automations (
id,
revision,
name,
description,
api_enabled,
target_repository,
target_ref,
target_workflow
) VALUES (?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(id) DO NOTHING
",
)
.bind(automation.id.as_str())
.bind(automation.revision.as_str())
.bind(&automation.name)
.bind(automation.description.as_deref())
.bind(automation.api_enabled())
.bind(&automation.target.repository)
.bind(&automation.target.ref_selector)
.bind(&automation.target.workflow)
.execute(&mut **transaction)
.await?;
if result.rows_affected() == 0 {
return Ok(false);
}
insert_schedule_triggers(transaction, automation).await?;
Ok(true)
}
fn id_from_path(path: &Path) -> Result<AutomationId, AutomationStoreError> {
let stem = path
.file_stem()
.and_then(|stem| stem.to_str())
.ok_or_else(|| AutomationStoreError::InvalidFilename {
path: path.to_path_buf(),
reason: "filename is not valid UTF-8".to_string(),
})?;
AutomationId::new(stem).map_err(|source| AutomationStoreError::InvalidFilename {
path: path.to_path_buf(),
reason: source.to_string(),
async fn insert_schedule_triggers(
transaction: &mut Transaction<'_, Sqlite>,
automation: &Automation,
) -> Result<(), AutomationStoreError> {
for trigger in automation.schedule_triggers() {
sqlx::query(
r"
INSERT INTO automation_triggers (automation_id, id, enabled, expression)
VALUES (?, ?, ?, ?)
",
)
.bind(automation.id.as_str())
.bind(trigger.id.as_str())
.bind(trigger.enabled)
.bind(&trigger.expression)
.execute(&mut **transaction)
.await?;
}
Ok(())
}
async fn current_revision(
transaction: &mut Transaction<'_, Sqlite>,
id: &AutomationId,
) -> Result<Option<AutomationRevision>, AutomationStoreError> {
let current = sqlx::query_scalar::<_, String>("SELECT revision FROM automations WHERE id = ?")
.bind(id.as_str())
.fetch_optional(&mut **transaction)
.await?;
current
.map(|revision| {
AutomationRevision::from_str(&revision).map_err(|source| {
AutomationStoreError::InvalidRevision {
id: id.clone(),
source,
}
})
})
.transpose()
}
async fn revision_mismatch_error(
transaction: &mut Transaction<'_, Sqlite>,
id: &AutomationId,
expected: &AutomationRevision,
) -> Result<AutomationStoreError, AutomationStoreError> {
let Some(actual) = current_revision(transaction, id).await? else {
return Err(AutomationStoreError::NotFound { id: id.clone() });
};
Ok(AutomationStoreError::StaleRevision {
id: id.clone(),
expected: expected.clone(),
actual,
})
}
fn is_toml_file(path: &Path) -> bool {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| extension == "toml")
}
async fn write_atomic(dir: &Path, path: &Path, bytes: &[u8]) -> Result<(), AutomationStoreError> {
fs::create_dir_all(dir)
.await
.map_err(|err| AutomationStoreError::io(dir, err))?;
let temp_path = temp_path_for(path);
let mut file = fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&temp_path)
.await
.map_err(|err| AutomationStoreError::io(&temp_path, err))?;
if let Err(err) = file.write_all(bytes).await {
cleanup_temp(&temp_path).await;
return Err(AutomationStoreError::io(&temp_path, err));
}
if let Err(err) = file.sync_all().await {
cleanup_temp(&temp_path).await;
return Err(AutomationStoreError::io(&temp_path, err));
}
drop(file);
if let Err(err) = fs::rename(&temp_path, path).await {
cleanup_temp(&temp_path).await;
return Err(AutomationStoreError::io(path, err));
}
Ok(())
}
async fn write_new(dir: &Path, path: &Path, bytes: &[u8]) -> Result<(), AutomationStoreError> {
fs::create_dir_all(dir)
.await
.map_err(|err| AutomationStoreError::io(dir, err))?;
let temp_path = temp_path_for(path);
let mut file = fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&temp_path)
.await
.map_err(|err| AutomationStoreError::io(&temp_path, err))?;
if let Err(err) = file.write_all(bytes).await {
cleanup_temp(&temp_path).await;
return Err(AutomationStoreError::io(&temp_path, err));
}
if let Err(err) = file.sync_all().await {
cleanup_temp(&temp_path).await;
return Err(AutomationStoreError::io(&temp_path, err));
}
drop(file);
if let Err(err) = fs::hard_link(&temp_path, path).await {
cleanup_temp(&temp_path).await;
return Err(AutomationStoreError::io(path, err));
}
cleanup_temp(&temp_path).await;
Ok(())
}
async fn cleanup_temp(path: &Path) {
let _ = fs::remove_file(path).await;
}
fn create_error_for(id: AutomationId, err: AutomationStoreError) -> AutomationStoreError {
match err {
AutomationStoreError::Io { source, .. } if source.kind() == ErrorKind::AlreadyExists => {
AutomationStoreError::AlreadyExists { id }
}
err => err,
}
}
fn temp_path_for(path: &Path) -> PathBuf {
let parent = path.parent().unwrap_or_else(|| Path::new("."));
let file_name = path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("automation.toml");
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |duration| duration.as_nanos());
parent.join(format!(".{file_name}.{}.{}.tmp", std::process::id(), now))
}
fn automation_path(dir: &Path, id: &AutomationId) -> PathBuf {
dir.join(format!("{id}.toml"))
}
#[cfg(test)]
mod tests {
use tokio::fs;
use crate::{
ApiTrigger, AutomationDraft, AutomationId, AutomationReplace, AutomationStore,
AutomationStoreError, AutomationTarget, AutomationTrigger, AutomationTriggerId,
ScheduleTrigger,
};
fn target() -> AutomationTarget {
AutomationTarget {
repository: "fabro-sh/fabro".to_string(),
ref_selector: "main".to_string(),
workflow: "release".to_string(),
}
}
fn draft(id: &str, name: &str) -> AutomationDraft {
AutomationDraft {
id: AutomationId::new(id).unwrap(),
name: name.to_string(),
description: None,
target: target(),
triggers: vec![
AutomationTrigger::Api(ApiTrigger {
id: AutomationTriggerId::new("manual").unwrap(),
enabled: true,
}),
AutomationTrigger::Schedule(ScheduleTrigger {
id: AutomationTriggerId::new("nightly").unwrap(),
enabled: true,
expression: "0 0 * * *".to_string(),
}),
],
}
}
fn replacement(name: &str) -> AutomationReplace {
AutomationReplace {
name: name.to_string(),
description: Some("updated".to_string()),
target: target(),
triggers: vec![AutomationTrigger::Api(ApiTrigger {
id: AutomationTriggerId::new("manual").unwrap(),
enabled: false,
})],
}
}
#[tokio::test]
async fn missing_directory_loads_empty_store() {
let dir = tempfile::tempdir().unwrap();
let store = AutomationStore::load(dir.path().join("automations")).unwrap();
assert!(store.list().await.is_empty());
}
#[tokio::test]
async fn load_ignores_non_toml_files_and_keeps_valid_automations() {
let dir = tempfile::tempdir().unwrap();
let automation_dir = dir.path().join("automations");
fs::create_dir_all(&automation_dir).await.unwrap();
fs::write(automation_dir.join("notes.txt"), "ignore")
.await
.unwrap();
fs::write(
automation_dir.join("valid.toml"),
r#"
name = "Valid"
[target]
repository = "fabro-sh/fabro"
ref = "main"
workflow = "release"
"#,
)
.await
.unwrap();
let store = AutomationStore::load(&automation_dir).unwrap();
let automations = store.list().await;
assert_eq!(automations.len(), 1);
assert_eq!(automations[0].id.as_str(), "valid");
assert_eq!(automations[0].name, "Valid");
}
#[tokio::test]
async fn load_fails_on_malformed_toml() {
let dir = tempfile::tempdir().unwrap();
let automation_dir = dir.path().join("automations");
fs::create_dir_all(&automation_dir).await.unwrap();
fs::write(automation_dir.join("broken.toml"), "not valid toml =")
.await
.unwrap();
let err = AutomationStore::load(&automation_dir).unwrap_err();
assert!(matches!(err, AutomationStoreError::Parse { .. }));
}
#[tokio::test]
async fn load_fails_on_invalid_filename_id() {
let dir = tempfile::tempdir().unwrap();
let automation_dir = dir.path().join("automations");
fs::create_dir_all(&automation_dir).await.unwrap();
fs::write(automation_dir.join("Bad Name.toml"), "name = \"Bad\"")
.await
.unwrap();
let err = AutomationStore::load(&automation_dir).unwrap_err();
assert!(matches!(err, AutomationStoreError::InvalidFilename { .. }));
}
#[tokio::test]
async fn create_replace_and_delete_round_trip_files_and_revisions() {
let dir = tempfile::tempdir().unwrap();
let automation_dir = dir.path().join("automations");
let store = AutomationStore::load(&automation_dir).unwrap();
let created = store.create(draft("nightly", "Nightly")).await.unwrap();
let path = automation_dir.join("nightly.toml");
let persisted = fs::read_to_string(&path).await.unwrap();
assert!(persisted.contains("name = \"Nightly\""));
assert!(!top_level_lines(&persisted).any(|line| line.starts_with("id = ")));
assert!(!top_level_lines(&persisted).any(|line| line.starts_with("revision = ")));
assert_eq!(
created.revision,
crate::AutomationRevision::from_bytes(persisted.as_bytes())
);
assert!(store.create(draft("nightly", "Duplicate")).await.is_err());
let stale = crate::AutomationRevision::from_bytes(b"stale");
assert!(
store
.replace(&created.id, &stale, replacement("Updated"))
.await
.is_err()
);
let replaced = store
.replace(&created.id, &created.revision, replacement("Updated"))
.await
.unwrap();
assert_ne!(replaced.revision, created.revision);
assert_eq!(
store.get(&created.id).await.unwrap().revision,
replaced.revision
);
store.delete(&created.id, &replaced.revision).await.unwrap();
assert!(store.get(&created.id).await.is_none());
assert!(!path.exists());
}
fn top_level_lines(toml: &str) -> impl Iterator<Item = &str> {
toml.lines().take_while(|line| !line.starts_with('['))
}
}

View file

@ -0,0 +1,369 @@
#![expect(
clippy::unwrap_used,
reason = "SQLite automation-store integration tests use panic-on-failure fixture setup"
)]
use std::path::Path;
use fabro_automation::{
ApiTrigger, AutomationDraft, AutomationId, AutomationReplace, AutomationStore,
AutomationStoreError, AutomationTarget, AutomationTrigger, AutomationTriggerId,
ScheduleTrigger,
};
use fabro_db::Database;
use tokio::fs;
async fn test_database() -> (tempfile::TempDir, Database) {
let dir = tempfile::tempdir().unwrap();
let database = Database::connect(dir.path().join("fabro.sqlite3"))
.await
.unwrap();
database.migrate().await.unwrap();
(dir, database)
}
fn target() -> AutomationTarget {
AutomationTarget {
repository: "fabro-sh/fabro".to_string(),
ref_selector: "main".to_string(),
workflow: "release".to_string(),
}
}
fn schedule(id: &str, expression: &str, enabled: bool) -> AutomationTrigger {
AutomationTrigger::Schedule(ScheduleTrigger {
id: AutomationTriggerId::new(id).unwrap(),
enabled,
expression: expression.to_string(),
})
}
fn draft(id: &str, api_enabled: bool) -> AutomationDraft {
AutomationDraft {
id: AutomationId::new(id).unwrap(),
name: "Nightly".to_string(),
description: Some("Runs every night".to_string()),
target: target(),
triggers: vec![
schedule("z-last", "0 2 * * *", false),
AutomationTrigger::Api(ApiTrigger {
id: AutomationTriggerId::new("custom-api-id").unwrap(),
enabled: api_enabled,
}),
schedule("a-first", "0 1 * * *", true),
],
}
}
fn replacement(name: &str, expression: &str) -> AutomationReplace {
AutomationReplace {
name: name.to_string(),
description: None,
target: target(),
triggers: vec![
schedule("nightly", expression, true),
AutomationTrigger::Api(ApiTrigger {
id: AutomationTriggerId::new("api").unwrap(),
enabled: true,
}),
],
}
}
fn trigger_ids(automation: &fabro_automation::Automation) -> Vec<&str> {
automation
.triggers
.iter()
.map(|trigger| trigger.id().as_str())
.collect()
}
#[tokio::test]
async fn crud_normalizes_api_and_schedule_order() {
let (_dir, database) = test_database().await;
let store = AutomationStore::new(database.clone_pool());
let created = store.create(draft("nightly", true)).await.unwrap();
assert_eq!(trigger_ids(&created), vec!["manual", "a-first", "z-last"]);
assert!(created.enabled_api_trigger().is_some());
let fetched = store.get(&created.id).await.unwrap().unwrap();
assert_eq!(fetched, created);
assert_eq!(store.list().await.unwrap(), vec![created.clone()]);
let replaced = store
.replace(
&created.id,
&created.revision,
replacement("Updated", "30 4 * * *"),
)
.await
.unwrap();
assert_ne!(replaced.revision, created.revision);
assert_eq!(replaced.name, "Updated");
store
.delete(&replaced.id, &replaced.revision)
.await
.unwrap();
assert!(store.get(&replaced.id).await.unwrap().is_none());
}
#[tokio::test]
async fn disabled_api_trigger_normalizes_to_absent() {
let (_dir, database) = test_database().await;
let store = AutomationStore::new(database.clone_pool());
let created = store.create(draft("nightly", false)).await.unwrap();
assert!(created.enabled_api_trigger().is_none());
assert_eq!(trigger_ids(&created), vec!["a-first", "z-last"]);
assert_eq!(store.get(&created.id).await.unwrap().unwrap(), created);
}
#[tokio::test]
async fn equivalent_trigger_orders_have_the_same_revision() {
let (_dir, database) = test_database().await;
let store = AutomationStore::new(database.clone_pool());
let first = store.create(draft("first", true)).await.unwrap();
let mut reordered = draft("second", true);
reordered.triggers.reverse();
let second = store.create(reordered).await.unwrap();
assert_eq!(first.revision, second.revision);
}
#[tokio::test]
async fn create_conflict_and_conditional_delete_errors_are_typed() {
let (_dir, database) = test_database().await;
let store = AutomationStore::new(database.clone_pool());
let created = store.create(draft("nightly", true)).await.unwrap();
let duplicate = store.create(draft("nightly", true)).await.unwrap_err();
assert!(matches!(
duplicate,
AutomationStoreError::AlreadyExists { .. }
));
let mut revision_source = draft("revision-source", true);
revision_source.name = "Different revision".to_string();
let stale_revision = store.create(revision_source).await.unwrap().revision;
let stale = store
.delete(&created.id, &stale_revision)
.await
.unwrap_err();
assert!(matches!(stale, AutomationStoreError::StaleRevision { .. }));
let missing = AutomationId::new("missing").unwrap();
let not_found = store.delete(&missing, &stale_revision).await.unwrap_err();
assert!(matches!(not_found, AutomationStoreError::NotFound { .. }));
}
#[tokio::test]
async fn independent_pools_observe_writes_and_revision_conflicts() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("fabro.sqlite3");
let first_database = Database::connect(&path).await.unwrap();
first_database.migrate().await.unwrap();
let second_database = Database::connect(&path).await.unwrap();
second_database.migrate().await.unwrap();
let first = AutomationStore::new(first_database.clone_pool());
let second = AutomationStore::new(second_database.clone_pool());
let created = first.create(draft("nightly", true)).await.unwrap();
assert_eq!(second.get(&created.id).await.unwrap().unwrap(), created);
let replaced = second
.replace(
&created.id,
&created.revision,
replacement("Winner", "0 5 * * *"),
)
.await
.unwrap();
let err = first
.replace(
&created.id,
&created.revision,
replacement("Loser", "0 6 * * *"),
)
.await
.unwrap_err();
assert!(matches!(err, AutomationStoreError::StaleRevision {
actual,
..
} if actual == replaced.revision));
}
#[tokio::test]
async fn failed_schedule_insert_rolls_back_parent_replace() {
let (_dir, database) = test_database().await;
let store = AutomationStore::new(database.clone_pool());
let created = store.create(draft("nightly", true)).await.unwrap();
sqlx::query(
r"
CREATE TRIGGER reject_blocked_schedule
BEFORE INSERT ON automation_triggers
WHEN NEW.id = 'blocked'
BEGIN
SELECT RAISE(ABORT, 'blocked schedule');
END
",
)
.execute(database.pool())
.await
.unwrap();
let replacement = AutomationReplace {
name: "Should roll back".to_string(),
description: None,
target: target(),
triggers: vec![schedule("blocked", "0 7 * * *", true)],
};
let err = store
.replace(&created.id, &created.revision, replacement)
.await
.unwrap_err();
assert!(matches!(err, AutomationStoreError::Db { .. }));
assert_eq!(store.get(&created.id).await.unwrap().unwrap(), created);
}
#[tokio::test]
async fn invalid_stored_schedule_is_rejected_on_read() {
let (_dir, database) = test_database().await;
let store = AutomationStore::new(database.clone_pool());
let created = store.create(draft("nightly", true)).await.unwrap();
sqlx::query("UPDATE automation_triggers SET expression = 'not cron' WHERE automation_id = ?")
.bind(created.id.as_str())
.execute(database.pool())
.await
.unwrap();
let err = store.get(&created.id).await.unwrap_err();
assert!(matches!(err, AutomationStoreError::StoredValidation { .. }));
}
#[tokio::test]
async fn legacy_import_is_transactional_and_sql_wins() {
let (dir, database) = test_database().await;
let store = AutomationStore::new(database.clone_pool());
store.create(draft("existing", true)).await.unwrap();
let source_dir = dir.path().join("automations");
fs::create_dir_all(&source_dir).await.unwrap();
write_legacy_automation(&source_dir, "existing", "Legacy existing").await;
write_legacy_automation(&source_dir, "imported", "Imported").await;
fs::write(source_dir.join("notes.txt"), "ignored")
.await
.unwrap();
let report = fabro_automation::import_legacy_directory_once(database.pool(), &source_dir)
.await
.unwrap()
.unwrap();
assert_eq!(report.imported_rows, 1);
assert_eq!(report.skipped_rows, 1);
assert_eq!(report.names, vec!["imported"]);
assert!(!source_dir.exists());
assert!(report.backup_path.exists());
assert_eq!(
store
.get(&AutomationId::new("existing").unwrap())
.await
.unwrap()
.unwrap()
.name,
"Nightly"
);
assert_eq!(
store
.get(&AutomationId::new("imported").unwrap())
.await
.unwrap()
.unwrap()
.name,
"Imported"
);
fs::create_dir_all(&source_dir).await.unwrap();
write_legacy_automation(&source_dir, "existing", "Legacy existing").await;
write_legacy_automation(&source_dir, "imported", "Imported again").await;
let retry = fabro_automation::import_legacy_directory_once(database.pool(), &source_dir)
.await
.unwrap()
.unwrap();
assert_eq!(retry.imported_rows, 0);
assert_eq!(retry.skipped_rows, 2);
assert!(retry.backup_path.exists());
assert_eq!(
store
.get(&AutomationId::new("imported").unwrap())
.await
.unwrap()
.unwrap()
.name,
"Imported"
);
assert!(
fabro_automation::import_legacy_directory_once(database.pool(), &source_dir)
.await
.unwrap()
.is_none()
);
}
#[tokio::test]
async fn invalid_legacy_file_leaves_directory_and_database_unchanged() {
let (dir, database) = test_database().await;
let source_dir = dir.path().join("automations");
fs::create_dir_all(&source_dir).await.unwrap();
write_legacy_automation(&source_dir, "valid", "Valid").await;
fs::write(source_dir.join("broken.toml"), "not valid toml =")
.await
.unwrap();
let err = fabro_automation::import_legacy_directory_once(database.pool(), &source_dir)
.await
.unwrap_err();
assert!(matches!(err, AutomationStoreError::Parse { .. }));
assert!(source_dir.exists());
assert!(
AutomationStore::new(database.clone_pool())
.list()
.await
.unwrap()
.is_empty()
);
}
async fn write_legacy_automation(dir: &Path, id: &str, name: &str) {
fs::write(
dir.join(format!("{id}.toml")),
format!(
r#"name = "{name}"
[target]
repository = "fabro-sh/fabro"
ref = "main"
workflow = "release"
[[triggers]]
id = "manual"
type = "api"
enabled = true
[[triggers]]
id = "nightly"
type = "schedule"
enabled = true
expression = "0 3 * * *"
"#
),
)
.await
.unwrap();
}

View file

@ -2174,10 +2174,7 @@ mod tests {
use super::*;
async fn load_secret_snapshot(storage: &Storage) -> Vault {
fabro_vault::SecretStore::open(storage.sqlite_path(), storage.secrets_path())
.await
.unwrap()
.snapshot()
fabro_vault::SecretStore::open_snapshot(storage.sqlite_path(), storage.secrets_path())
.await
.unwrap()
.into_vault()

View file

@ -278,17 +278,14 @@ async fn load_worker_vault(storage_dir: Option<&Path>) -> Result<Option<Arc<Asyn
};
let storage = Storage::new(storage_dir);
let vault = SecretStore::open(storage.sqlite_path(), storage.secrets_path())
let vault = SecretStore::open_snapshot(storage.sqlite_path(), storage.secrets_path())
.await
.with_context(|| {
format!(
"failed to open worker secret store from {}",
"failed to load worker secrets from {}",
storage.root().display()
)
})?
.snapshot()
.await
.context("loading worker secrets snapshot")?
.into_vault();
Ok(Some(Arc::new(AsyncRwLock::new(vault))))
}

View file

@ -292,10 +292,7 @@ async fn execute_daemon(
validate_startup_configuration(&resolved_settings)?;
let storage = Storage::new(storage_dir);
migrate_startup_vault(storage.secrets_path());
let startup_vault = SecretStore::open(storage.sqlite_path(), storage.secrets_path())
.await
.context("opening the Fabro secret store for startup validation")?
.snapshot()
let startup_vault = SecretStore::open_snapshot(storage.sqlite_path(), storage.secrets_path())
.await
.context("loading secrets for startup validation")?;
validate_startup(

View file

@ -14,10 +14,7 @@ use fabro_vault::{SecretType, Vault};
const INSTALL_COMMAND_TIMEOUT: Duration = Duration::from_secs(30);
async fn load_secret_snapshot(storage: &Storage) -> Vault {
fabro_vault::SecretStore::open(storage.sqlite_path(), storage.secrets_path())
.await
.expect("test secret store should open")
.snapshot()
fabro_vault::SecretStore::open_snapshot(storage.sqlite_path(), storage.secrets_path())
.await
.expect("test secret snapshot should load")
.into_vault()

View file

@ -0,0 +1,34 @@
CREATE TABLE automations (
id TEXT PRIMARY KEY NOT NULL,
revision TEXT NOT NULL,
name TEXT NOT NULL,
description TEXT,
api_enabled INTEGER NOT NULL,
target_repository TEXT NOT NULL,
target_ref TEXT NOT NULL,
target_workflow TEXT NOT NULL,
CHECK (length(id) BETWEEN 1 AND 63),
CHECK (substr(id, 1, 1) GLOB '[a-z0-9]'),
CHECK (id NOT GLOB '*[^a-z0-9-]*'),
CHECK (length(revision) = 64),
CHECK (revision NOT GLOB '*[^0-9a-f]*'),
CHECK (length(trim(name)) > 0),
CHECK (api_enabled IN (0, 1)),
CHECK (length(target_repository) BETWEEN 3 AND 140),
CHECK (length(target_ref) BETWEEN 1 AND 255),
CHECK (length(target_workflow) BETWEEN 1 AND 255)
);
CREATE TABLE automation_triggers (
automation_id TEXT NOT NULL,
id TEXT NOT NULL,
enabled INTEGER NOT NULL,
expression TEXT NOT NULL,
PRIMARY KEY (automation_id, id),
FOREIGN KEY (automation_id) REFERENCES automations(id) ON DELETE CASCADE,
CHECK (length(id) BETWEEN 1 AND 63),
CHECK (substr(id, 1, 1) GLOB '[a-z0-9]'),
CHECK (id NOT GLOB '*[^a-z0-9_-]*'),
CHECK (enabled IN (0, 1)),
CHECK (length(trim(expression)) > 0)
);

View file

@ -1,79 +0,0 @@
//! Helpers shared by the one-time imports that seed SQLite tables from
//! pre-SQLite on-disk stores.
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use chrono::{DateTime, Utc};
use tokio::fs;
/// Failure to move an imported legacy source aside, carrying the backup path
/// the caller needs for its own error variant.
#[derive(Debug)]
pub struct LegacyBackupError {
pub backup_path: PathBuf,
pub source: std::io::Error,
}
/// Move an imported legacy file or directory aside to
/// `<name>.imported-<timestamp>.bak` next to the original. `fallback_name` is
/// used when the source path has no final component.
pub async fn rename_to_legacy_backup(
source: &Path,
fallback_name: &str,
) -> Result<PathBuf, LegacyBackupError> {
let backup_path = legacy_backup_path(source, fallback_name, Utc::now());
fs::rename(source, &backup_path)
.await
.map_err(|source| LegacyBackupError {
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}
fn legacy_backup_path(source: &Path, fallback_name: &str, imported_at: DateTime<Utc>) -> PathBuf {
let timestamp = imported_at.format("%Y%m%dT%H%M%S%fZ");
let mut file_name = source
.file_name()
.map_or_else(|| OsString::from(fallback_name), OsString::from);
file_name.push(format!(".imported-{timestamp}.bak"));
source.with_file_name(file_name)
}
/// True when `path` has a `.toml` extension (legacy per-item store files).
pub fn is_toml_file(path: &Path) -> bool {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| extension == "toml")
}
#[cfg(test)]
mod tests {
use chrono::TimeZone as _;
use super::*;
#[test]
fn backup_path_appends_timestamped_suffix() {
let imported_at = Utc.with_ymd_and_hms(2026, 7, 11, 1, 2, 3).unwrap();
let backup =
legacy_backup_path(Path::new("/data/secrets.json"), "secrets.json", imported_at);
assert_eq!(
backup,
Path::new("/data/secrets.json.imported-20260711T010203000000000Z.bak")
);
}
#[test]
fn backup_path_uses_fallback_when_source_has_no_file_name() {
let imported_at = Utc.with_ymd_and_hms(2026, 7, 11, 1, 2, 3).unwrap();
let backup = legacy_backup_path(Path::new("/"), "mcps", imported_at);
assert!(
backup
.file_name()
.and_then(|name| name.to_str())
.is_some_and(|name| name.starts_with("mcps.imported-"))
);
}
}

View file

@ -1,4 +1,5 @@
use std::collections::HashSet;
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use std::time::Duration;
@ -11,15 +12,8 @@ use tokio::fs;
use tokio::task::spawn_blocking;
use tracing::info;
pub mod legacy;
pub type DbPool = sqlx::SqlitePool;
/// Parse an RFC 3339 TEXT column value into a UTC timestamp.
pub fn parse_rfc3339_utc(value: &str) -> Result<DateTime<Utc>, chrono::ParseError> {
DateTime::parse_from_rfc3339(value).map(|timestamp| timestamp.with_timezone(&Utc))
}
static MIGRATOR: Migrator = sqlx::migrate!("./migrations");
#[derive(Clone)]
@ -161,6 +155,36 @@ impl Database {
}
}
/// Parse an RFC 3339 timestamp column value into UTC.
pub fn parse_rfc3339(value: &str) -> Result<DateTime<Utc>, chrono::ParseError> {
DateTime::parse_from_rfc3339(value).map(|timestamp| timestamp.with_timezone(&Utc))
}
/// Result of a one-time import of a legacy file or directory into SQLite.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ImportReport {
pub source_path: PathBuf,
pub backup_path: PathBuf,
pub imported_rows: usize,
pub skipped_rows: usize,
pub names: Vec<String>,
}
/// Backup destination for a legacy file or directory after a one-time import
/// into SQLite. `default_name` is used when `source_path` has no file name.
pub fn legacy_backup_path(
source_path: &Path,
default_name: &str,
imported_at: DateTime<Utc>,
) -> PathBuf {
let timestamp = imported_at.format("%Y%m%dT%H%M%S%fZ");
let mut file_name = source_path
.file_name()
.map_or_else(|| OsString::from(default_name), OsString::from);
file_name.push(format!(".imported-{timestamp}.bak"));
source_path.with_file_name(file_name)
}
/// Rollback artifact written by [`Database::migrate`] before applying new
/// migrations: the database path with `.pre-migration.bak` appended.
pub fn pre_migration_snapshot_path(database_path: &Path) -> PathBuf {

View file

@ -56,6 +56,16 @@ async fn connect_creates_parent_directory_and_migrate_is_idempotent() -> anyhow:
.await?;
assert_eq!(mcp_servers_table_count, 1);
for table in ["automations", "automation_triggers"] {
let count: i64 = sqlx::query_scalar(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = ?",
)
.bind(table)
.fetch_one(database.pool())
.await?;
assert_eq!(count, 1, "{table} table should exist");
}
let legacy_import_table_count: i64 = sqlx::query_scalar(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'table' AND name = 'legacy_imports'",
)
@ -222,6 +232,105 @@ async fn insert_mcp_server(
Ok(())
}
#[tokio::test]
async fn automations_schema_enforces_aggregate_constraints() -> anyhow::Result<()> {
let dir = tempfile::tempdir()?;
let database = fabro_db::Database::connect(dir.path().join("fabro.sqlite3")).await?;
database.migrate().await?;
insert_minimal_automation(database.pool(), "valid", 1).await?;
sqlx::query(
"INSERT INTO automation_triggers (automation_id, id, enabled, expression) \
VALUES (?, ?, ?, ?)",
)
.bind("valid")
.bind("nightly")
.bind(true)
.bind("0 3 * * *")
.execute(database.pool())
.await?;
assert!(
insert_minimal_automation(database.pool(), "Bad", 1)
.await
.is_err()
);
assert!(
insert_minimal_automation(database.pool(), "bad-bool", 2_i64)
.await
.is_err()
);
assert!(
sqlx::query(
"INSERT INTO automation_triggers (automation_id, id, enabled, expression) \
VALUES (?, ?, ?, ?)",
)
.bind("valid")
.bind("Bad!")
.bind(true)
.bind("0 4 * * *")
.execute(database.pool())
.await
.is_err()
);
assert!(
sqlx::query(
"INSERT INTO automation_triggers (automation_id, id, enabled, expression) \
VALUES (?, ?, ?, ?)",
)
.bind("missing")
.bind("nightly")
.bind(true)
.bind("0 4 * * *")
.execute(database.pool())
.await
.is_err()
);
sqlx::query("DELETE FROM automations WHERE id = ?")
.bind("valid")
.execute(database.pool())
.await?;
let trigger_count: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM automation_triggers WHERE automation_id = ?")
.bind("valid")
.fetch_one(database.pool())
.await?;
assert_eq!(trigger_count, 0);
Ok(())
}
async fn insert_minimal_automation(
pool: &fabro_db::DbPool,
id: &str,
api_enabled: i64,
) -> Result<(), sqlx::Error> {
sqlx::query(
r"
INSERT INTO automations (
id,
revision,
name,
api_enabled,
target_repository,
target_ref,
target_workflow
) VALUES (?, ?, ?, ?, ?, ?, ?)
",
)
.bind(id)
.bind("a".repeat(64))
.bind("Automation")
.bind(api_enabled)
.bind("fabro-sh/fabro")
.bind("main")
.bind("release")
.execute(pool)
.await?;
Ok(())
}
#[tokio::test]
async fn environments_schema_rejects_invalid_rows() -> anyhow::Result<()> {
let dir = tempfile::tempdir()?;

View file

@ -3,11 +3,12 @@ use std::path::{Path, PathBuf};
use std::str::FromStr;
use std::sync::Arc;
use chrono::Utc;
use fabro_config::{
EnvironmentDockerfileLayer, EnvironmentImageLayer, EnvironmentLayer, EnvironmentLifecycleLayer,
EnvironmentNetworkLayer, EnvironmentResourcesLayer, MergeMap, StickyMap,
};
use fabro_db::{DbPool, legacy};
use fabro_db::DbPool;
use fabro_types::settings::run::{DockerfileSource, EnvironmentProvider, EnvironmentSettings};
use fabro_types::settings::{Duration, InterpString, Size};
use serde::de::DeserializeOwned;
@ -704,7 +705,7 @@ async fn legacy_environment_paths(
.file_type()
.await
.map_err(|source| EnvironmentStoreError::io(&path, source))?;
if file_type.is_file() && legacy::is_toml_file(&path) {
if file_type.is_file() && is_toml_file(&path) {
paths.push(LegacyEnvironmentPath {
id: id_from_path(&path)?,
path,
@ -754,9 +755,11 @@ async fn existing_environment_ids(
async fn rename_imported_legacy_directory(
source_dir: &Path,
) -> Result<PathBuf, EnvironmentStoreError> {
legacy::rename_to_legacy_backup(source_dir, "environments")
let backup_path = fabro_db::legacy_backup_path(source_dir, "environments", Utc::now());
fs::rename(source_dir, &backup_path)
.await
.map_err(|err| EnvironmentStoreError::io(&err.backup_path, err.source))
.map_err(|source| EnvironmentStoreError::io(&backup_path, source))?;
Ok(backup_path)
}
fn id_from_path(path: &Path) -> Result<EnvironmentId, EnvironmentStoreError> {
@ -773,6 +776,12 @@ fn id_from_path(path: &Path) -> Result<EnvironmentId, EnvironmentStoreError> {
})
}
fn is_toml_file(path: &Path) -> bool {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| extension == "toml")
}
fn encode_json<T: serde::Serialize>(
field: &'static str,
value: &T,

View file

@ -723,10 +723,7 @@ mod tests {
use fabro_vault::{SecretType as VaultSecretType, Vault};
async fn load_secret_snapshot(storage: &Storage) -> Vault {
SecretStore::open(storage.sqlite_path(), storage.secrets_path())
.await
.unwrap()
.snapshot()
SecretStore::open_snapshot(storage.sqlite_path(), storage.secrets_path())
.await
.unwrap()
.into_vault()

View file

@ -10,4 +10,5 @@ mod model;
mod store;
pub use error::McpServerStoreError;
pub use store::{ImportReport, McpServerStore, import_legacy_directory_once};
pub use fabro_db::ImportReport;
pub use store::{McpServerStore, import_legacy_directory_once};

View file

@ -4,7 +4,8 @@ use std::path::{Path, PathBuf};
use std::str::FromStr as _;
use std::sync::RwLock;
use fabro_db::{DbPool, legacy};
use chrono::Utc;
use fabro_db::{DbPool, ImportReport};
use fabro_types::settings::run::{McpHttpProtocol, McpServerSettings, McpTransport};
use fabro_types::{
McpServerDefinition, McpServerDraft, McpServerId, McpServerReplace, McpServerRevision,
@ -35,15 +36,6 @@ pub struct McpServerStore {
defs: RwLock<HashMap<McpServerId, McpServerDefinition>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ImportReport {
pub source_path: PathBuf,
pub backup_path: PathBuf,
pub imported_rows: usize,
pub skipped_rows: usize,
pub mcp_server_ids: Vec<String>,
}
#[derive(Debug, Clone, Copy, Display, EnumString, IntoStaticStr)]
#[strum(serialize_all = "snake_case")]
enum TransportType {
@ -642,14 +634,14 @@ pub async fn import_legacy_directory_once(
backup_path,
imported_rows: imported_ids.len(),
skipped_rows,
mcp_server_ids: imported_ids,
names: imported_ids,
};
info!(
source_path = %report.source_path.display(),
backup_path = %report.backup_path.display(),
imported_rows = report.imported_rows,
skipped_rows = report.skipped_rows,
mcp_server_ids = ?report.mcp_server_ids,
mcp_server_ids = ?report.names,
"Imported legacy MCP server directory into SQLite"
);
Ok(Some(report))
@ -675,7 +667,7 @@ async fn legacy_definition_paths(
.file_type()
.await
.map_err(|source| McpServerStoreError::io(&path, source))?;
if file_type.is_file() && legacy::is_toml_file(&path) {
if file_type.is_file() && is_toml_file(&path) {
paths.push((id_from_path(&path)?, path));
}
}
@ -711,14 +703,22 @@ fn id_from_path(path: &Path) -> Result<McpServerId, McpServerStoreError> {
})
}
fn is_toml_file(path: &Path) -> bool {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| extension == "toml")
}
async fn rename_imported_legacy_directory(
source_dir: &Path,
) -> Result<PathBuf, McpServerStoreError> {
legacy::rename_to_legacy_backup(source_dir, "mcps")
let backup_path = fabro_db::legacy_backup_path(source_dir, "mcps", Utc::now());
fs::rename(source_dir, &backup_path)
.await
.map_err(|err| McpServerStoreError::LegacyBackup {
.map_err(|source| McpServerStoreError::LegacyBackup {
source_path: source_dir.to_path_buf(),
backup_path: err.backup_path,
source: err.source,
})
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}

View file

@ -233,7 +233,7 @@ async fn imports_legacy_directory_once_without_overwriting_sql() {
assert_eq!(report.imported_rows, 1);
assert_eq!(report.skipped_rows, 1);
assert_eq!(report.mcp_server_ids, vec!["new"]);
assert_eq!(report.names, vec!["new"]);
assert!(!source.exists());
assert!(report.backup_path.join("notes.txt").exists());
let reloaded = McpServerStore::load(database.clone_pool()).await.unwrap();

View file

@ -113,6 +113,7 @@ tower = "0.5"
http-body-util = "0.1"
httpmock = "0.8"
serde_yaml = "0.9"
sqlx.workspace = true
tracing-subscriber.workspace = true
tokio-util.workspace = true
tokio-tungstenite.workspace = true

View file

@ -2502,12 +2502,10 @@ mod tests {
assert!(!server_env.contains_key(EnvVars::GITHUB_APP_CLIENT_SECRET));
assert!(!server_env.contains_key(EnvVars::GITHUB_APP_WEBHOOK_SECRET));
let vault = fabro_vault::SecretStore::open(storage.sqlite_path(), storage.secrets_path())
.await
.unwrap()
.snapshot()
.await
.unwrap();
let vault =
fabro_vault::SecretStore::open_snapshot(storage.sqlite_path(), storage.secrets_path())
.await
.unwrap();
assert_eq!(vault.get(EnvVars::GITHUB_TOKEN), None);
assert_eq!(
vault.get(EnvVars::GITHUB_APP_CLIENT_SECRET),

View file

@ -34,7 +34,7 @@ use crate::canonical_origin::resolve_canonical_origin;
use crate::github_webhooks::{TailscaleFunnelManager, WEBHOOK_ROUTE, WEBHOOK_SECRET_ENV};
use crate::interp::process_env_var;
use crate::server::{
AppState, AppStateConfig, ResolvedAppStateSettings, RouterOptions, build_app_state,
self, AppState, AppStateConfig, ResolvedAppStateSettings, RouterOptions, build_app_state,
build_router_with_options, reconcile_incomplete_runs_on_startup, shutdown_active_workers,
spawn_automation_scheduler, spawn_scheduler,
};
@ -735,6 +735,15 @@ where
legacy_environment_dir.display()
)
})?;
let legacy_automation_dir = server::automation_dir_for_active_config(&active_config_path);
fabro_automation::import_legacy_directory_once(database.pool(), &legacy_automation_dir)
.await
.with_context(|| {
format!(
"importing legacy automations directory {}",
legacy_automation_dir.display()
)
})?;
let db_pool = database.clone_pool();
let max_concurrent_runs = resolved_server_settings.scheduler.max_concurrent_runs;
// In `--watch-web` mode the build watcher will populate `dist/` shortly

View file

@ -2309,7 +2309,7 @@ fn build_sandbox_provider_registry(
SandboxProviderRegistry::new(providers)
}
fn automation_dir_for_active_config(active_config_path: &std::path::Path) -> PathBuf {
pub(crate) fn automation_dir_for_active_config(active_config_path: &std::path::Path) -> PathBuf {
active_config_path
.parent()
.unwrap_or_else(|| std::path::Path::new("."))
@ -2369,12 +2369,7 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
automation_materializer_override,
} = config;
let automation_dir = automation_dir_for_active_config(&active_config_path);
let automation_store = Arc::new(
AutomationStore::load(automation_dir)
.map_err(anyhow::Error::new)
.context("load automations")?,
);
let automation_store = Arc::new(AutomationStore::new(db_pool.clone()));
let local_provider_enabled = resolved_settings
.server_settings
.server

View file

@ -180,7 +180,18 @@ pub(crate) fn spawn_automation_scheduler(state: Arc<AppState>) {
break;
}
let automations = state.automation_store().list().await;
let automations = match state.automation_store().list().await {
Ok(automations) => automations,
Err(err) => {
error!(error = ?err, "Failed to load automations for scheduler");
tokio::select! {
() = shutdown.cancelled() => break,
() = state.automation_scheduler_notified() => {},
() = sleep(AUTOMATION_SCHEDULER_MAX_SLEEP) => {},
}
continue;
}
};
let now = Utc::now();
for due in planner.tick(&automations, now) {
let state = Arc::clone(&state);
@ -304,7 +315,11 @@ fn run_due_schedules_once<'a>(
now: DateTime<Utc>,
) -> std::pin::Pin<Box<dyn std::future::Future<Output = ()> + Send + 'a>> {
Box::pin(async move {
let automations = state.automation_store().list().await;
let automations = state
.automation_store()
.list()
.await
.expect("test automations should load");
for trigger in planner.tick(&automations, now) {
Box::pin(fire_scheduled_automation_run(
Arc::clone(&state),

View file

@ -46,17 +46,20 @@ pub(super) fn routes() -> Router<Arc<AppState>> {
)
}
async fn list_automations(_auth: RequiredUser, State(state): State<Arc<AppState>>) -> Response {
let data = state.automation_store().list().await;
async fn list_automations(
_auth: RequiredUser,
State(state): State<Arc<AppState>>,
) -> Result<Response, ApiError> {
let data = state.automation_store().list().await?;
let total = data.len();
(
Ok((
StatusCode::OK,
Json(AutomationListResponse {
data,
meta: AutomationListMeta { total },
}),
)
.into_response()
.into_response())
}
async fn list_automation_runs(
@ -69,8 +72,12 @@ async fn list_automation_runs(
Ok(id) => id,
Err(err) => return err.into_response(),
};
if state.automation_store().get(&id).await.is_none() {
return ApiError::not_found(format!("automation not found: {id}")).into_response();
match state.automation_store().exists(&id).await {
Ok(true) => {}
Ok(false) => {
return ApiError::not_found(format!("automation not found: {id}")).into_response();
}
Err(err) => return ApiError::from(err).into_response(),
}
let entries = match state
@ -126,8 +133,12 @@ async fn create_automation_run(
Ok(id) => id,
Err(err) => return err.into_response(),
};
let Some(automation) = state.automation_store().get(&id).await else {
return ApiError::not_found(format!("automation not found: {id}")).into_response();
let automation = match state.automation_store().get(&id).await {
Ok(Some(automation)) => automation,
Ok(None) => {
return ApiError::not_found(format!("automation not found: {id}")).into_response();
}
Err(err) => return ApiError::from(err).into_response(),
};
let Some(api_trigger) = automation.enabled_api_trigger() else {
return ApiError::with_code(
@ -210,7 +221,7 @@ async fn get_automation(
Path(id): Path<String>,
) -> Result<Response, ApiError> {
let id = parse_path_id(id)?;
match state.automation_store().get(&id).await {
match state.automation_store().get(&id).await? {
Some(automation) => Ok(automation_with_etag_response(StatusCode::OK, automation)),
None => Err(ApiError::not_found(format!("automation not found: {id}"))),
}
@ -273,18 +284,13 @@ impl From<AutomationStoreError> for ApiError {
AutomationStoreError::Validation { source } => {
Self::new(StatusCode::UNPROCESSABLE_ENTITY, source.to_string())
}
// The handlers parse `If-Match` before reaching the store, so a
// missing-revision error from the store would indicate an internal
// bug rather than a client problem.
AutomationStoreError::MissingRevision { .. }
| AutomationStoreError::InvalidFilename { .. }
| AutomationStoreError::Parse { .. }
| AutomationStoreError::InvalidUtf8 { .. }
| AutomationStoreError::Serialize { .. }
| AutomationStoreError::Io { .. } => Self::new(
StatusCode::INTERNAL_SERVER_ERROR,
"automation store operation failed",
),
err => {
tracing::error!(error = ?err, "Automation store operation failed");
Self::new(
StatusCode::INTERNAL_SERVER_ERROR,
"automation store operation failed",
)
}
}
}
}

View file

@ -252,6 +252,10 @@ impl TestAppStateBuilder {
&vault_path,
self.default_environment_provider,
)?;
import_test_legacy_automations(
db_pool.clone(),
server::automation_dir_for_active_config(&active_config_path),
)?;
let preloaded_vault = test_secret_snapshot(db_pool.clone())?;
build_app_state(AppStateConfig {
resolved_settings: resolved_runtime_settings_for_tests(
@ -581,6 +585,26 @@ fn test_db_pool(
.expect("test database setup thread should not panic")
}
#[expect(
clippy::disallowed_methods,
reason = "sync test builders may run inside async tests; a short-lived OS thread avoids nested Tokio runtimes"
)]
fn import_test_legacy_automations(pool: DbPool, source_dir: PathBuf) -> anyhow::Result<()> {
std::thread::spawn(move || {
let runtime = TokioRuntimeBuilder::new_current_thread()
.enable_all()
.build()?;
runtime
.block_on(fabro_automation::import_legacy_directory_once(
&pool, source_dir,
))
.map(|_| ())
.map_err(anyhow::Error::new)
})
.join()
.expect("test automation import thread should not panic")
}
pub fn test_app_state_with_store_and_runtime_settings(
server_settings: ServerSettings,
manifest_run_defaults: RunLayer,

View file

@ -2,11 +2,13 @@ use std::path::{Path, PathBuf};
use axum::body::Body;
use axum::http::{Method, Request, StatusCode, header};
use fabro_config::Storage;
use fabro_server::server::build_router;
use fabro_server::test_support::{
TestAppStateBuilder, TestAutomationRunMaterializer, build_test_router, test_auth_mode,
};
use serde_json::{Value, json};
use sqlx::Row as _;
use tower::ServiceExt;
use crate::helpers::{
@ -62,17 +64,20 @@ fn replacement_body(name: &str) -> Value {
fn automation_app() -> (axum::Router, tempfile::TempDir, PathBuf) {
let temp_dir = tempfile::tempdir().expect("automation test tempdir should be created");
let active_config_path = temp_dir.path().join("settings.toml");
let automation_dir = temp_dir.path().join("automations");
let vault_path = temp_dir.path().join("secrets.json");
let sqlite_path = Storage::new(temp_dir.path()).sqlite_path();
let state = TestAppStateBuilder::new()
.active_config_path(active_config_path)
.vault_path(vault_path)
.build();
(build_test_router(state), temp_dir, automation_dir)
(build_test_router(state), temp_dir, sqlite_path)
}
fn automation_app_with_fake_materializer() -> (axum::Router, tempfile::TempDir, PathBuf) {
let temp_dir = tempfile::tempdir().expect("automation test tempdir should be created");
let active_config_path = temp_dir.path().join("settings.toml");
let automation_dir = temp_dir.path().join("automations");
let vault_path = temp_dir.path().join("secrets.json");
let sqlite_path = Storage::new(temp_dir.path()).sqlite_path();
let materialized_manifest: fabro_api::types::RunManifest =
serde_json::from_value(minimal_manifest_json(MINIMAL_DOT))
.expect("minimal run manifest fixture should deserialize");
@ -80,12 +85,13 @@ fn automation_app_with_fake_materializer() -> (axum::Router, tempfile::TempDir,
serde_json::to_vec(&materialized_manifest).expect("minimal run manifest should serialize");
let state = TestAppStateBuilder::new()
.active_config_path(active_config_path)
.vault_path(vault_path)
.automation_materializer(TestAutomationRunMaterializer::succeed(
materialized_manifest,
submitted_manifest_bytes,
))
.build();
(build_test_router(state), temp_dir, automation_dir)
(build_test_router(state), temp_dir, sqlite_path)
}
fn json_request(method: Method, path: &str, body: &Value) -> Request<Body> {
@ -197,35 +203,36 @@ fn assert_schedule_trigger(body: &Value, expression: &str, enabled: bool) {
);
}
async fn persisted_automation_toml(automation_dir: &Path, id: &str) -> toml::Value {
let persisted = tokio::fs::read_to_string(automation_dir.join(format!("{id}.toml")))
async fn persisted_automation(sqlite_path: &Path, id: &str) -> Option<Value> {
let database = fabro_db::Database::connect(sqlite_path)
.await
.expect("persisted automation TOML should be readable");
toml::from_str(&persisted).expect("persisted automation TOML should parse")
}
fn assert_persisted_schedule_trigger(body: &toml::Value, expression: &str, enabled: bool) {
let triggers = body
.get("triggers")
.and_then(toml::Value::as_array)
.expect("persisted automation TOML should include triggers");
let trigger = triggers
.iter()
.find(|trigger| trigger.get("id").and_then(toml::Value::as_str) == Some("nightly"))
.expect("persisted automation TOML should include nightly trigger");
assert_eq!(
trigger.get("type").and_then(toml::Value::as_str),
Some("schedule")
);
assert_eq!(
trigger.get("expression").and_then(toml::Value::as_str),
Some(expression)
);
assert_eq!(
trigger.get("enabled").and_then(toml::Value::as_bool),
Some(enabled)
);
.expect("automation test database should open");
let parent = sqlx::query("SELECT api_enabled FROM automations WHERE id = ?")
.bind(id)
.fetch_optional(database.pool())
.await
.expect("persisted automation should query")?;
let triggers = sqlx::query(
"SELECT id, enabled, expression FROM automation_triggers \
WHERE automation_id = ? ORDER BY id",
)
.bind(id)
.fetch_all(database.pool())
.await
.expect("persisted automation triggers should query")
.into_iter()
.map(|row| {
json!({
"id": row.get::<String, _>("id"),
"enabled": row.get::<bool, _>("enabled"),
"expression": row.get::<String, _>("expression")
})
})
.collect::<Vec<_>>();
Some(json!({
"api_enabled": parent.get::<bool, _>("api_enabled"),
"triggers": triggers
}))
}
#[tokio::test]
@ -250,19 +257,60 @@ async fn empty_automation_list_returns_total_zero() {
}
#[tokio::test]
async fn create_automation_persists_sibling_toml_file() {
let (app, _temp_dir, automation_dir) = automation_app();
async fn create_automation_persists_sql_aggregate() {
let (app, _temp_dir, sqlite_path) = automation_app();
let body = create_automation(&app, "nightly", "Nightly").await;
assert_eq!(body["id"], "nightly");
assert_eq!(body["name"], "Nightly");
assert!(automation_dir.join("nightly.toml").exists());
assert_eq!(
persisted_automation(&sqlite_path, "nightly").await,
Some(json!({
"api_enabled": true,
"triggers": [{
"id": "nightly",
"enabled": true,
"expression": "0 3 * * *"
}]
}))
);
}
#[tokio::test]
async fn schedule_trigger_round_trips_through_create_list_get_and_toml() {
let (app, _temp_dir, automation_dir) = automation_app();
async fn automation_persists_across_app_rebuild() {
let temp_dir = tempfile::tempdir().expect("automation test tempdir should be created");
let active_config_path = temp_dir.path().join("settings.toml");
let vault_path = temp_dir.path().join("secrets.json");
let state = TestAppStateBuilder::new()
.active_config_path(active_config_path.clone())
.vault_path(vault_path.clone())
.build();
let app = build_test_router(state);
let created = create_automation(&app, "nightly", "Nightly").await;
drop(app);
let state = TestAppStateBuilder::new()
.active_config_path(active_config_path)
.vault_path(vault_path)
.build();
let response = build_test_router(state)
.oneshot(empty_request(Method::GET, "/automations/nightly"))
.await
.expect("get persisted automation should respond");
let retrieved = response_json(
response,
StatusCode::OK,
"GET /api/v1/automations/nightly after app rebuild",
)
.await;
assert_eq!(retrieved, created);
}
#[tokio::test]
async fn schedule_trigger_round_trips_through_create_list_get_and_sql() {
let (app, _temp_dir, sqlite_path) = automation_app();
let created = create_automation(&app, "nightly", "Nightly").await;
assert_schedule_trigger(&created, "0 3 * * *", true);
@ -283,11 +331,43 @@ async fn schedule_trigger_round_trips_through_create_list_get_and_toml() {
let list = response_json(response, StatusCode::OK, "GET /api/v1/automations").await;
assert_schedule_trigger(&list["data"][0], "0 3 * * *", true);
let persisted = persisted_automation_toml(&automation_dir, "nightly").await;
assert_persisted_schedule_trigger(&persisted, "0 3 * * *", true);
assert!(persisted.get("id").is_none());
assert!(persisted.get("revision").is_none());
assert!(persisted.get("enabled").is_none());
assert_eq!(
persisted_automation(&sqlite_path, "nightly")
.await
.expect("automation should be persisted")["triggers"][0],
json!({
"id": "nightly",
"enabled": true,
"expression": "0 3 * * *"
})
);
}
#[tokio::test]
async fn invalid_stored_automation_returns_internal_server_error() {
let (app, _temp_dir, sqlite_path) = automation_app();
create_automation(&app, "nightly", "Nightly").await;
let database = fabro_db::Database::connect(sqlite_path)
.await
.expect("automation test database should open");
sqlx::query(
"UPDATE automation_triggers SET expression = 'not cron' WHERE automation_id = 'nightly'",
)
.execute(database.pool())
.await
.expect("stored schedule should be corrupted for the test");
let response = app
.oneshot(empty_request(Method::GET, "/automations/nightly"))
.await
.expect("get invalid stored automation should respond");
response_status(
response,
StatusCode::INTERNAL_SERVER_ERROR,
"GET /api/v1/automations/nightly with invalid stored schedule",
)
.await;
}
#[tokio::test]
@ -386,7 +466,7 @@ async fn replace_automation_accepts_unquoted_if_match_and_returns_new_etag() {
#[tokio::test]
async fn replace_automation_round_trips_schedule_trigger() {
let (app, _temp_dir, automation_dir) = automation_app();
let (app, _temp_dir, sqlite_path) = automation_app();
let created = create_automation(&app, "nightly", "Nightly").await;
let revision = revision_from(&created);
let mut replacement = replacement_body("Rescheduled");
@ -416,8 +496,16 @@ async fn replace_automation_round_trips_schedule_trigger() {
let body = response_json(response, StatusCode::OK, "PUT /api/v1/automations/nightly").await;
assert_schedule_trigger(&body, "30 4 * * *", false);
let persisted = persisted_automation_toml(&automation_dir, "nightly").await;
assert_persisted_schedule_trigger(&persisted, "30 4 * * *", false);
assert_eq!(
persisted_automation(&sqlite_path, "nightly")
.await
.expect("automation should be persisted")["triggers"][0],
json!({
"id": "nightly",
"enabled": false,
"expression": "30 4 * * *"
})
);
}
#[tokio::test]
@ -495,8 +583,8 @@ async fn replace_and_delete_automation_require_if_match() {
}
#[tokio::test]
async fn delete_automation_removes_file_and_resource() {
let (app, _temp_dir, automation_dir) = automation_app();
async fn delete_automation_removes_sql_aggregate() {
let (app, _temp_dir, sqlite_path) = automation_app();
let created = create_automation(&app, "nightly", "Nightly").await;
let revision = revision_from(&created);
@ -517,7 +605,11 @@ async fn delete_automation_removes_file_and_resource() {
)
.await;
assert!(!automation_dir.join("nightly.toml").exists());
assert!(
persisted_automation(&sqlite_path, "nightly")
.await
.is_none()
);
let response = app
.oneshot(empty_request(Method::GET, "/automations/nightly"))
.await
@ -670,6 +762,69 @@ async fn automation_store_malformed_persisted_toml_fails_startup() {
assert!(result.is_err());
}
#[tokio::test]
async fn legacy_automation_is_imported_and_directory_is_backed_up() {
let temp_dir = tempfile::tempdir().expect("automation test tempdir should be created");
let automation_dir = temp_dir.path().join("automations");
tokio::fs::create_dir_all(&automation_dir)
.await
.expect("automation dir should be created");
tokio::fs::write(
automation_dir.join("nightly.toml"),
r#"name = "Legacy nightly"
[target]
repository = "fabro-sh/fabro"
ref = "main"
workflow = "release"
[[triggers]]
type = "api"
id = "manual"
enabled = true
"#,
)
.await
.expect("legacy automation fixture should be written");
let state = TestAppStateBuilder::new()
.active_config_path(temp_dir.path().join("settings.toml"))
.vault_path(temp_dir.path().join("secrets.json"))
.build();
let response = build_test_router(state)
.oneshot(empty_request(Method::GET, "/automations/nightly"))
.await
.expect("get imported automation should respond");
let body = response_json(
response,
StatusCode::OK,
"GET /api/v1/automations/nightly after legacy import",
)
.await;
assert_eq!(body["name"], "Legacy nightly");
assert!(!automation_dir.exists());
let mut entries = tokio::fs::read_dir(temp_dir.path())
.await
.expect("storage directory should be readable");
let mut backup_exists = false;
while let Some(entry) = entries
.next_entry()
.await
.expect("storage directory entry should be readable")
{
if entry
.file_name()
.to_string_lossy()
.starts_with("automations.imported-")
{
backup_exists = true;
break;
}
}
assert!(backup_exists);
}
#[tokio::test]
async fn automations_routes_require_authenticated_user() {
let temp_dir = tempfile::tempdir().expect("automation test tempdir should be created");

View file

@ -37,10 +37,7 @@ fn spa_fixture_root() -> PathBuf {
async fn load_secret_snapshot(storage_dir: &std::path::Path) -> Vault {
let storage = Storage::new(storage_dir);
fabro_vault::SecretStore::open(storage.sqlite_path(), storage.secrets_path())
.await
.expect("test secret store should open")
.snapshot()
fabro_vault::SecretStore::open_snapshot(storage.sqlite_path(), storage.secrets_path())
.await
.expect("test secret snapshot should load")
.into_vault()

View file

@ -2,7 +2,7 @@ use std::collections::HashMap;
use std::path::{Path, PathBuf};
use chrono::{DateTime, Utc};
use fabro_db::{DbPool, legacy};
use fabro_db::DbPool;
use fabro_types::{Variable, is_env_style_name};
use sqlx::Row as _;
use sqlx::sqlite::SqliteRow;
@ -312,13 +312,15 @@ fn parse_legacy_entries(
}
async fn rename_imported_legacy_file(source_path: &Path) -> Result<PathBuf, Error> {
legacy::rename_to_legacy_backup(source_path, "variables.json")
let backup_path = fabro_db::legacy_backup_path(source_path, "variables.json", Utc::now());
fs::rename(source_path, &backup_path)
.await
.map_err(|err| Error::LegacyBackup {
.map_err(|source| Error::LegacyBackup {
source_path: source_path.to_path_buf(),
backup_path: err.backup_path,
source: err.source,
})
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}
fn row_count(count: usize) -> Result<i64, Error> {
@ -342,7 +344,7 @@ fn variable_from_row(row: &SqliteRow) -> Result<Variable, Error> {
}
fn parse_timestamp(name: &str, column: &'static str, value: &str) -> Result<DateTime<Utc>, Error> {
fabro_db::parse_rfc3339_utc(value).map_err(|source| Error::Timestamp {
fabro_db::parse_rfc3339(value).map_err(|source| Error::Timestamp {
name: name.to_string(),
column,
source,

View file

@ -10,12 +10,12 @@ use std::path::{Component, Path, PathBuf};
use std::{fmt, io};
use chrono::{DateTime, Utc};
pub use fabro_db::ImportReport;
use fabro_static::EnvVars;
pub use fabro_types::SecretType;
use fabro_types::{SecretMetadata, is_env_style_name};
pub use store::{
ImportReport, SecretSnapshot, SecretStore, SecretStoreError, SecretStoreWrite,
import_legacy_json_once,
SecretSnapshot, SecretStore, SecretStoreError, SecretStoreWrite, import_legacy_json_once,
};
#[derive(Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]

View file

@ -3,7 +3,7 @@ use std::path::{Path, PathBuf};
use std::str::FromStr as _;
use chrono::{DateTime, Utc};
use fabro_db::{Database, DbPool, legacy};
use fabro_db::{Database, DbPool, ImportReport};
use fabro_types::{OAuthCredential, SecretMetadata, SecretType};
use sqlx::sqlite::SqliteRow;
use sqlx::{Row as _, Sqlite, Transaction};
@ -87,15 +87,6 @@ pub enum SecretStoreError {
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ImportReport {
pub source_path: PathBuf,
pub backup_path: PathBuf,
pub imported_rows: usize,
pub skipped_rows: usize,
pub secret_names: Vec<String>,
}
#[derive(Clone, PartialEq, Eq)]
pub struct SecretStoreWrite {
pub name: String,
@ -124,6 +115,10 @@ pub struct SecretStore {
pub struct SecretSnapshot(Vault);
impl SecretSnapshot {
/// Converts the snapshot into a path-less in-memory [`Vault`]. Mutations to
/// the returned vault are never persisted; callers that need durable writes
/// must go through [`SecretStore`] (see `SqlVaultCredentialSource`'s
/// before/after diffing for the OAuth refresh write-back).
#[must_use]
pub fn into_vault(self) -> Vault {
self.0
@ -168,6 +163,19 @@ impl SecretStore {
Ok(Self::new(database.clone_pool()))
}
/// Like [`SecretStore::open`], but returns a point-in-time snapshot of the
/// store instead of the store itself.
pub async fn open_snapshot(
sqlite_path: impl AsRef<Path>,
legacy_secrets_path: impl AsRef<Path>,
) -> anyhow::Result<SecretSnapshot> {
let snapshot = Self::open(sqlite_path, legacy_secrets_path)
.await?
.snapshot()
.await?;
Ok(snapshot)
}
pub async fn get(&self, name: &str) -> Result<Option<SecretEntry>, SecretStoreError> {
let row = sqlx::query(
"SELECT name, secret_type, value, description, revision, created_at, updated_at \
@ -176,7 +184,10 @@ impl SecretStore {
.bind(name)
.fetch_optional(&self.pool)
.await?;
row.as_ref().map(entry_from_row).transpose()
row.as_ref()
.map(entry_from_row)
.transpose()
.map(|entry| entry.map(|(_, entry)| entry))
}
pub async fn list(&self) -> Result<Vec<SecretMetadata>, SecretStoreError> {
@ -284,7 +295,7 @@ impl SecretStore {
.await?;
if let Some(row) = row {
return entry_from_row(&row);
return entry_from_row(&row).map(|(_, entry)| entry);
}
match self.get(name).await? {
Some(entry) => Err(SecretStoreError::StaleRevision {
@ -305,8 +316,8 @@ impl SecretStore {
.await?;
let entries = rows
.iter()
.map(|row| Ok((row.try_get::<String, _>("name")?, entry_from_row(row)?)))
.collect::<Result<HashMap<_, _>, SecretStoreError>>()?;
.map(entry_from_row)
.collect::<Result<HashMap<_, _>, _>>()?;
Ok(SecretSnapshot(Vault::from_entries(entries)))
}
}
@ -402,14 +413,14 @@ pub async fn import_legacy_json_once(
backup_path,
imported_rows: imported_names.len(),
skipped_rows,
secret_names: imported_names,
names: imported_names,
};
info!(
source_path = %report.source_path.display(),
backup_path = %report.backup_path.display(),
imported_rows = report.imported_rows,
skipped_rows = report.skipped_rows,
secret_names = ?report.secret_names,
secret_names = ?report.names,
removal_deadline = IMPORT_REMOVAL_DEADLINE,
"Imported legacy secrets JSON into SQLite"
);
@ -440,7 +451,7 @@ async fn insert_legacy_entry(
Ok(result.rows_affected() == 1)
}
fn entry_from_row(row: &SqliteRow) -> Result<SecretEntry, SecretStoreError> {
fn entry_from_row(row: &SqliteRow) -> Result<(String, SecretEntry), SecretStoreError> {
let metadata = metadata_from_row(row)?;
let name = metadata.name;
let value = row.try_get::<String, _>("value")?;
@ -454,14 +465,15 @@ fn entry_from_row(row: &SqliteRow) -> Result<SecretEntry, SecretStoreError> {
if revision <= 0 {
return Err(SecretStoreError::StoredRevision { name, revision });
}
Ok(SecretEntry {
let entry = SecretEntry {
value,
secret_type: metadata.secret_type,
description: metadata.description,
created_at: metadata.created_at,
updated_at: metadata.updated_at,
revision,
})
};
Ok((name, entry))
}
fn metadata_from_row(row: &SqliteRow) -> Result<SecretMetadata, SecretStoreError> {
@ -495,7 +507,7 @@ fn parse_timestamp(
column: &'static str,
value: &str,
) -> Result<DateTime<Utc>, SecretStoreError> {
fabro_db::parse_rfc3339_utc(value).map_err(|source| SecretStoreError::Timestamp {
fabro_db::parse_rfc3339(value).map_err(|source| SecretStoreError::Timestamp {
name: name.to_string(),
column,
source,
@ -529,11 +541,13 @@ fn validate_stored_name(name: &str, secret_type: SecretType) -> Result<(), Secre
}
async fn rename_imported_legacy_file(source_path: &Path) -> Result<PathBuf, SecretStoreError> {
legacy::rename_to_legacy_backup(source_path, "secrets.json")
let backup_path = fabro_db::legacy_backup_path(source_path, "secrets.json", Utc::now());
fs::rename(source_path, &backup_path)
.await
.map_err(|err| SecretStoreError::LegacyBackup {
.map_err(|source| SecretStoreError::LegacyBackup {
source_path: source_path.to_path_buf(),
backup_path: err.backup_path,
source: err.source,
})
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}