diff --git a/Cargo.lock b/Cargo.lock index 2b0b894cf..2ac91d62f 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -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", diff --git a/docs/internal/parallel-strategy.md b/docs/internal/parallel-strategy.md new file mode 100644 index 000000000..351764ade --- /dev/null +++ b/docs/internal/parallel-strategy.md @@ -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///pass/`), 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.`, 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.` 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: ` 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.`) 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/@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.` 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. diff --git a/docs/plans/2026-07-11-automations-sqlite-schema.md b/docs/plans/2026-07-11-automations-sqlite-schema.md new file mode 100644 index 000000000..c23280197 --- /dev/null +++ b/docs/plans/2026-07-11-automations-sqlite-schema.md @@ -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-.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. diff --git a/docs/public/execution/automations.mdx b/docs/public/execution/automations.mdx index b2b314339..e6d28b6ea 100644 --- a/docs/public/execution/automations.mdx +++ b/docs/public/execution/automations.mdx @@ -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" diff --git a/lib/crates/fabro-automation/Cargo.toml b/lib/crates/fabro-automation/Cargo.toml index 6ed2bd595..6284862e7 100644 --- a/lib/crates/fabro-automation/Cargo.toml +++ b/lib/crates/fabro-automation/Cargo.toml @@ -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"] } diff --git a/lib/crates/fabro-automation/migrations/2026071101_file_definitions_to_sqlite.rs b/lib/crates/fabro-automation/migrations/2026071101_file_definitions_to_sqlite.rs new file mode 100644 index 000000000..2621a712b --- /dev/null +++ b/lib/crates/fabro-automation/migrations/2026071101_file_definitions_to_sqlite.rs @@ -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, +) -> Result, 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>, 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 { + 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 { + 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) +} diff --git a/lib/crates/fabro-automation/src/error.rs b/lib/crates/fabro-automation/src/error.rs index cf5ffa595..8ea767365 100644 --- a/lib/crates/fabro-automation/src/error.rs +++ b/lib/crates/fabro-automation/src/error.rs @@ -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", } } } diff --git a/lib/crates/fabro-automation/src/lib.rs b/lib/crates/fabro-automation/src/lib.rs index f5f453100..7498e2878 100644 --- a/lib/crates/fabro-automation/src/lib.rs +++ b/lib/crates/fabro-automation/src/lib.rs @@ -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, diff --git a/lib/crates/fabro-automation/src/migrations.rs b/lib/crates/fabro-automation/src/migrations.rs new file mode 100644 index 000000000..824c47179 --- /dev/null +++ b/lib/crates/fabro-automation/src/migrations.rs @@ -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; diff --git a/lib/crates/fabro-automation/src/model.rs b/lib/crates/fabro-automation/src/model.rs index f00be4fdf..90be56d8b 100644 --- a/lib/crates/fabro-automation/src/model.rs +++ b/lib/crates/fabro-automation/src/model.rs @@ -21,6 +21,8 @@ static SCHEDULE_CRON_PARSER: LazyLock = 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), 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 { + 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 { + 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 { - 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 { + 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::>(); + 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 { diff --git a/lib/crates/fabro-automation/src/store.rs b/lib/crates/fabro-automation/src/store.rs index ccd6ea7a4..775580e69 100644 --- a/lib/crates/fabro-automation/src/store.rs +++ b/lib/crates/fabro-automation/src/store.rs @@ -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>, + 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) -> Result { - 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 { - let automations = self.automations.read().await; - let mut values = automations.values().cloned().collect::>(); - values.sort_by(|left, right| left.id.cmp(&right.id)); - values + pub async fn list(&self) -> Result, 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 { - self.automations.read().await.get(id).cloned() + pub async fn get(&self, id: &AutomationId) -> Result, 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 { + 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 { 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 { - 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 ¤t.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 ¤t.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, 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, + api_enabled: bool, + target: AutomationTarget, + schedule_triggers: Vec, +} - 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 { + let id_value = row.try_get::("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::("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::, _>("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::, _>("trigger_enabled")? + .ok_or_else(|| AutomationStoreError::StoredTriggerShape { + id: self.id.clone(), + })?, + expression: row + .try_get::, _>("trigger_expression")? + .ok_or_else(|| AutomationStoreError::StoredTriggerShape { + id: self.id.clone(), + })?, + }); + Ok(()) + } + + fn finish(self) -> Result { + // `from_stored` canonicalizes trigger order and the manual API trigger. + let mut triggers = self + .schedule_triggers + .into_iter() + .map(AutomationTrigger::Schedule) + .collect::>(); + 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, AutomationStoreError> { + let mut automations = Vec::new(); + let mut current: Option = None; + + for row in rows { + let row_id = row.try_get::("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 { - 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 { + 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 { - 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, 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 { + 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 { - toml.lines().take_while(|line| !line.starts_with('[')) - } -} diff --git a/lib/crates/fabro-automation/tests/store.rs b/lib/crates/fabro-automation/tests/store.rs new file mode 100644 index 000000000..3f419895b --- /dev/null +++ b/lib/crates/fabro-automation/tests/store.rs @@ -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(); +} diff --git a/lib/crates/fabro-cli/src/commands/install.rs b/lib/crates/fabro-cli/src/commands/install.rs index 7ea8885c0..e12cf04c2 100644 --- a/lib/crates/fabro-cli/src/commands/install.rs +++ b/lib/crates/fabro-cli/src/commands/install.rs @@ -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() diff --git a/lib/crates/fabro-cli/src/commands/run/runner.rs b/lib/crates/fabro-cli/src/commands/run/runner.rs index 90283c353..4e6b72b8b 100644 --- a/lib/crates/fabro-cli/src/commands/run/runner.rs +++ b/lib/crates/fabro-cli/src/commands/run/runner.rs @@ -278,17 +278,14 @@ async fn load_worker_vault(storage_dir: Option<&Path>) -> Result 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() diff --git a/lib/crates/fabro-db/migrations/2026071103_automations.sql b/lib/crates/fabro-db/migrations/2026071103_automations.sql new file mode 100644 index 000000000..64b638a91 --- /dev/null +++ b/lib/crates/fabro-db/migrations/2026071103_automations.sql @@ -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) +); diff --git a/lib/crates/fabro-db/src/legacy.rs b/lib/crates/fabro-db/src/legacy.rs deleted file mode 100644 index 74065be40..000000000 --- a/lib/crates/fabro-db/src/legacy.rs +++ /dev/null @@ -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 -/// `.imported-.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 { - 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) -> 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-")) - ); - } -} diff --git a/lib/crates/fabro-db/src/lib.rs b/lib/crates/fabro-db/src/lib.rs index ba930523e..16111ff16 100644 --- a/lib/crates/fabro-db/src/lib.rs +++ b/lib/crates/fabro-db/src/lib.rs @@ -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, 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, 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, +} + +/// 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, +) -> 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 { diff --git a/lib/crates/fabro-db/tests/sqlite.rs b/lib/crates/fabro-db/tests/sqlite.rs index 60c371ffa..384dbcec2 100644 --- a/lib/crates/fabro-db/tests/sqlite.rs +++ b/lib/crates/fabro-db/tests/sqlite.rs @@ -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()?; diff --git a/lib/crates/fabro-environment/src/store.rs b/lib/crates/fabro-environment/src/store.rs index 798231a28..441906fb3 100644 --- a/lib/crates/fabro-environment/src/store.rs +++ b/lib/crates/fabro-environment/src/store.rs @@ -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 { - 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 { @@ -773,6 +776,12 @@ fn id_from_path(path: &Path) -> Result { }) } +fn is_toml_file(path: &Path) -> bool { + path.extension() + .and_then(|extension| extension.to_str()) + .is_some_and(|extension| extension == "toml") +} + fn encode_json( field: &'static str, value: &T, diff --git a/lib/crates/fabro-install/src/lib.rs b/lib/crates/fabro-install/src/lib.rs index 0d465368f..bfcfe9ab4 100644 --- a/lib/crates/fabro-install/src/lib.rs +++ b/lib/crates/fabro-install/src/lib.rs @@ -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() diff --git a/lib/crates/fabro-mcp-store/src/lib.rs b/lib/crates/fabro-mcp-store/src/lib.rs index 793c84ae4..872275454 100644 --- a/lib/crates/fabro-mcp-store/src/lib.rs +++ b/lib/crates/fabro-mcp-store/src/lib.rs @@ -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}; diff --git a/lib/crates/fabro-mcp-store/src/store.rs b/lib/crates/fabro-mcp-store/src/store.rs index 321ef88d1..fcf3f2697 100644 --- a/lib/crates/fabro-mcp-store/src/store.rs +++ b/lib/crates/fabro-mcp-store/src/store.rs @@ -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>, } -#[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, -} - #[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 { }) } +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 { - 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) } diff --git a/lib/crates/fabro-mcp-store/tests/store.rs b/lib/crates/fabro-mcp-store/tests/store.rs index f39cd1ea9..368907c7d 100644 --- a/lib/crates/fabro-mcp-store/tests/store.rs +++ b/lib/crates/fabro-mcp-store/tests/store.rs @@ -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(); diff --git a/lib/crates/fabro-server/Cargo.toml b/lib/crates/fabro-server/Cargo.toml index 3d54fb216..b951b339f 100644 --- a/lib/crates/fabro-server/Cargo.toml +++ b/lib/crates/fabro-server/Cargo.toml @@ -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 diff --git a/lib/crates/fabro-server/src/install.rs b/lib/crates/fabro-server/src/install.rs index a47b76557..841f405b5 100644 --- a/lib/crates/fabro-server/src/install.rs +++ b/lib/crates/fabro-server/src/install.rs @@ -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), diff --git a/lib/crates/fabro-server/src/serve.rs b/lib/crates/fabro-server/src/serve.rs index 642fa9435..6325e04a3 100644 --- a/lib/crates/fabro-server/src/serve.rs +++ b/lib/crates/fabro-server/src/serve.rs @@ -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 diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 7ed6514af..d31506ab6 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -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) { 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, ) -> std::pin::Pin + 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), diff --git a/lib/crates/fabro-server/src/server/handler/automations.rs b/lib/crates/fabro-server/src/server/handler/automations.rs index b3cf3423e..39c432ea6 100644 --- a/lib/crates/fabro-server/src/server/handler/automations.rs +++ b/lib/crates/fabro-server/src/server/handler/automations.rs @@ -46,17 +46,20 @@ pub(super) fn routes() -> Router> { ) } -async fn list_automations(_auth: RequiredUser, State(state): State>) -> Response { - let data = state.automation_store().list().await; +async fn list_automations( + _auth: RequiredUser, + State(state): State>, +) -> Result { + 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, ) -> Result { 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 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", + ) + } } } } diff --git a/lib/crates/fabro-server/src/test_support.rs b/lib/crates/fabro-server/src/test_support.rs index 31f8aad5f..4394c02a9 100644 --- a/lib/crates/fabro-server/src/test_support.rs +++ b/lib/crates/fabro-server/src/test_support.rs @@ -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, diff --git a/lib/crates/fabro-server/tests/it/api/automations.rs b/lib/crates/fabro-server/tests/it/api/automations.rs index 11ad21a28..4a2b96de0 100644 --- a/lib/crates/fabro-server/tests/it/api/automations.rs +++ b/lib/crates/fabro-server/tests/it/api/automations.rs @@ -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 { @@ -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 { + 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::("id"), + "enabled": row.get::("enabled"), + "expression": row.get::("expression") + }) + }) + .collect::>(); + Some(json!({ + "api_enabled": parent.get::("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"); diff --git a/lib/crates/fabro-server/tests/it/api/install.rs b/lib/crates/fabro-server/tests/it/api/install.rs index 937d12047..8dde87fdd 100644 --- a/lib/crates/fabro-server/tests/it/api/install.rs +++ b/lib/crates/fabro-server/tests/it/api/install.rs @@ -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() diff --git a/lib/crates/fabro-variable/src/lib.rs b/lib/crates/fabro-variable/src/lib.rs index 1148f62cc..974ef6a36 100644 --- a/lib/crates/fabro-variable/src/lib.rs +++ b/lib/crates/fabro-variable/src/lib.rs @@ -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 { - 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 { @@ -342,7 +344,7 @@ fn variable_from_row(row: &SqliteRow) -> Result { } fn parse_timestamp(name: &str, column: &'static str, value: &str) -> Result, 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, diff --git a/lib/crates/fabro-vault/src/lib.rs b/lib/crates/fabro-vault/src/lib.rs index 144ff7d23..6b9618759 100644 --- a/lib/crates/fabro-vault/src/lib.rs +++ b/lib/crates/fabro-vault/src/lib.rs @@ -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)] diff --git a/lib/crates/fabro-vault/src/store.rs b/lib/crates/fabro-vault/src/store.rs index 8bd7df938..68f3ae81e 100644 --- a/lib/crates/fabro-vault/src/store.rs +++ b/lib/crates/fabro-vault/src/store.rs @@ -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, -} - #[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, + legacy_secrets_path: impl AsRef, + ) -> anyhow::Result { + let snapshot = Self::open(sqlite_path, legacy_secrets_path) + .await? + .snapshot() + .await?; + Ok(snapshot) + } + pub async fn get(&self, name: &str) -> Result, 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, 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::("name")?, entry_from_row(row)?))) - .collect::, SecretStoreError>>()?; + .map(entry_from_row) + .collect::, _>>()?; 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 { +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::("value")?; @@ -454,14 +465,15 @@ fn entry_from_row(row: &SqliteRow) -> Result { 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 { @@ -495,7 +507,7 @@ fn parse_timestamp( column: &'static str, value: &str, ) -> Result, 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 { - 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) }