From 4b279e82718067d2af86e155e8bc9835ed34c51d Mon Sep 17 00:00:00 2001 From: sami Date: Thu, 10 Sep 2026 15:15:13 +0400 Subject: [PATCH] =?UTF-8?q?daemon:=20store/=20=E2=80=94=20SQLite=20WAL=20s?= =?UTF-8?q?chema=20+=20forward-only=20migrator=20(build=20step=203)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The daemon's persistent state. SQLite in WAL mode, foreign keys on, 5 s busy timeout so a writer waits rather than SQLITE_BUSY under the RPC loop. - store/sqlite — RAII Db/Stmt over the C API; errors returned as DbResult (std::expected), never thrown — the RPC loop must not unwind. transaction() helper: BEGIN / fn / COMMIT, ROLLBACK on error. - store/migrations/0001_initial.sql — the eight tables from the brief: settings, categories, queues, tasks, segments, rules, history, pairings. Notable choices: * tasks columns project onto proto TaskSummary with no computation; requested vs effective segments/buffer split per ADR 0010/0012; pause_reason column per ADR 0013. * segments end_byte is NOT constrained >= 0 so a whole-file zero-length download is one row with end_byte = -1 (ADR 0010 B3a). * pairings stores only token_sha256 — the plaintext token is returned once from session.pair and never persisted (CLAUDE.md §4). * indices on tasks(state), (category_id), (queue_id, queue_position), (created_at), (completed_at) for the "1000 tasks, download.list under 50 ms" DoD. * six built-in categories + a Main queue seeded. - store/migrations — runs every embedded migration past PRAGMA user_version, each in its own transaction, forward-only. SQL files are embedded at build time by cmake/embed_migrations.cmake. Test veloxd.store_migrations (ASan+UBSan and TSan clean): fresh DB -> head, all tables present, seed rows, FK cascade (segment orphan rejected, task delete cascades), the end_byte=-1 zero-length case, idempotent re-run, and forward-only from every released user_version. Also: daemon/docs/proto-requests-m1.md — P1 marked landed on lane/proto as 1.4.0 (HandlerError/HandlerResult), to be adopted in rpc/ once that merges to main; P2 resolved. Not linked into the running daemon yet — the store is wired to the dispatcher when download.add/list/get get real bodies, next. Co-Authored-By: Claude Sonnet 5 Claude-Session: https://claude.ai/code/session_01Upd9WhG9oppieig5nRDLig --- daemon/CMakeLists.txt | 31 +++- daemon/cmake/embed_migrations.cmake | 52 ++++++ daemon/docs/proto-requests-m1.md | 13 ++ daemon/src/store/migrations.cpp | 44 +++++ daemon/src/store/migrations.hpp | 38 ++++ daemon/src/store/migrations/.gitkeep | 0 daemon/src/store/migrations/0001_initial.sql | 176 +++++++++++++++++++ daemon/src/store/sqlite.cpp | 146 +++++++++++++++ daemon/src/store/sqlite.hpp | 109 ++++++++++++ daemon/tests/CMakeLists.txt | 5 + daemon/tests/store_migrations_test.cpp | 118 +++++++++++++ 11 files changed, 731 insertions(+), 1 deletion(-) create mode 100644 daemon/cmake/embed_migrations.cmake create mode 100644 daemon/src/store/migrations.cpp create mode 100644 daemon/src/store/migrations.hpp delete mode 100644 daemon/src/store/migrations/.gitkeep create mode 100644 daemon/src/store/migrations/0001_initial.sql create mode 100644 daemon/src/store/sqlite.cpp create mode 100644 daemon/src/store/sqlite.hpp create mode 100644 daemon/tests/store_migrations_test.cpp diff --git a/daemon/CMakeLists.txt b/daemon/CMakeLists.txt index 6191528..d0e442b 100644 --- a/daemon/CMakeLists.txt +++ b/daemon/CMakeLists.txt @@ -11,6 +11,35 @@ if(NOT TARGET nlohmann_json::nlohmann_json) find_package(nlohmann_json 3.11 REQUIRED) endif() find_package(Threads REQUIRED) +find_package(SQLite3 REQUIRED) + +# --- generated: migrations_embedded.hpp from src/store/migrations/*.sql --------------- +set(_mig_dir ${CMAKE_CURRENT_SOURCE_DIR}/src/store/migrations) +set(_mig_hdr ${CMAKE_CURRENT_BINARY_DIR}/generated/migrations_embedded.hpp) +file(GLOB _mig_srcs ${_mig_dir}/*.sql) +add_custom_command( + OUTPUT ${_mig_hdr} + COMMAND ${CMAKE_COMMAND} -DMIG_DIR=${_mig_dir} -DOUT=${_mig_hdr} + -P ${CMAKE_CURRENT_SOURCE_DIR}/cmake/embed_migrations.cmake + DEPENDS ${_mig_srcs} ${CMAKE_CURRENT_SOURCE_DIR}/cmake/embed_migrations.cmake + COMMENT "Embedding SQL migrations" + VERBATIM) +add_custom_target(veloxd_migrations_hdr DEPENDS ${_mig_hdr}) + +# --- veloxd_store — SQLite store + migrations ----------------------------------------- +add_library(veloxd_store STATIC + src/store/sqlite.cpp + src/store/migrations.cpp + ${_mig_hdr} +) +add_library(velox::daemon_store ALIAS veloxd_store) +target_include_directories(veloxd_store + PUBLIC ${CMAKE_CURRENT_SOURCE_DIR}/src + PRIVATE ${CMAKE_CURRENT_BINARY_DIR}/generated +) +target_compile_features(veloxd_store PUBLIC cxx_std_23) +target_compile_options(veloxd_store PRIVATE -Wall -Wextra -Wpedantic -Werror) +target_link_libraries(veloxd_store PUBLIC SQLite::SQLite3) # --- veloxd_rpc — the server library ---------------------------------------------------- add_library(veloxd_rpc STATIC @@ -34,7 +63,7 @@ target_link_libraries(veloxd_rpc add_executable(veloxd src/main.cpp) target_compile_features(veloxd PRIVATE cxx_std_23) target_compile_options(veloxd PRIVATE -Wall -Wextra -Wpedantic -Werror) -target_link_libraries(veloxd PRIVATE veloxd_rpc) +target_link_libraries(veloxd PRIVATE veloxd_rpc veloxd_store) if(VELOX_BUILD_TESTS AND EXISTS ${CMAKE_CURRENT_SOURCE_DIR}/tests/CMakeLists.txt) add_subdirectory(tests) diff --git a/daemon/cmake/embed_migrations.cmake b/daemon/cmake/embed_migrations.cmake new file mode 100644 index 0000000..e018260 --- /dev/null +++ b/daemon/cmake/embed_migrations.cmake @@ -0,0 +1,52 @@ +# Generates migrations_embedded.hpp from daemon/src/store/migrations/*.sql. +# +# cmake -DMIG_DIR= -DOUT= -P embed_migrations.cmake +# +# Each NNNN_name.sql becomes a Migration{ version = NNNN, name = "NNNN_name", sql = R"..." }. +# The delimiter for the raw string literal is chosen to not collide with the file body. + +file(GLOB _sql_files "${MIG_DIR}/*.sql") +list(SORT _sql_files) + +set(_entries "") +foreach(_f ${_sql_files}) + get_filename_component(_stem "${_f}" NAME_WE) # 0001_initial + string(REGEX MATCH "^([0-9]+)_" _m "${_stem}") + if(NOT _m) + message(FATAL_ERROR "migration file '${_f}' does not start with NNNN_") + endif() + string(REGEX REPLACE "^0*([0-9]+)_.*$" "\\1" _ver "${_stem}") + + file(READ "${_f}" _body) + # Pick a raw-string delimiter guaranteed absent from the body. + set(_delim "MIGSQL") + while(_body MATCHES "\\)${_delim}\"") + set(_delim "${_delim}X") + endwhile() + + string(APPEND _entries + " Migration{ ${_ver}, \"${_stem}\", R\"${_delim}(\n${_body}\n)${_delim}\" },\n") +endforeach() + +list(LENGTH _sql_files _count) + +set(_out "// GENERATED by embed_migrations.cmake — do not edit. Source: src/store/migrations/*.sql +#pragma once +#include +#include \"store/migrations.hpp\" + +namespace velox::daemon::store { + +inline constexpr std::array kEmbeddedMigrations = {{ +${_entries}}}; + +} // namespace velox::daemon::store +") + +if(EXISTS "${OUT}") + file(READ "${OUT}" _existing) + if(_existing STREQUAL "${_out}") + return() # unchanged — do not rewrite, keeps the build stable + endif() +endif() +file(WRITE "${OUT}" "${_out}") diff --git a/daemon/docs/proto-requests-m1.md b/daemon/docs/proto-requests-m1.md index 0833155..383510e 100644 --- a/daemon/docs/proto-requests-m1.md +++ b/daemon/docs/proto-requests-m1.md @@ -8,6 +8,19 @@ it needs a version note and a regen, not a silent change. --- +## Status + +- **P1 — landed** on `lane/proto` as `contracts/` **1.4.0** (commit `5e3e215`), as the + `HandlerError` / `HandlerResult` sketch below. Wire is byte-identical; C++-binding + bump only. `rpc/` adopts it (the predicted `Result` → `HandlerResult` swap on the + `on_*` overrides) **once `lane/proto` merges to `main`** — not against the unmerged + branch. `uds_roundtrip`'s `-32603`-collapse guard flips to `-32010` in the same change. +- **P2 — resolved.** `session.hello.version-mismatch`'s `data.expected` is now `$any`; + the error-fixture compare is on `code` only, so `rpc/` echoes `kProtocolVersion` there. + PROTO's writeup: `contracts/proto-answers-daemon-m1.md`. + +--- + ## P1. The generated `Dispatcher` has no error channel below `-32603` — **blocking a conformant server** `velox::proto::Dispatcher`'s 39 methods each return `Result` = diff --git a/daemon/src/store/migrations.cpp b/daemon/src/store/migrations.cpp new file mode 100644 index 0000000..1732811 --- /dev/null +++ b/daemon/src/store/migrations.cpp @@ -0,0 +1,44 @@ +#include "store/migrations.hpp" + +#include + +// Generated at build time from store/migrations/*.sql by embed_migrations.cmake. +#include "migrations_embedded.hpp" + +namespace velox::daemon::store { + +std::span embedded_migrations() { + return {kEmbeddedMigrations.data(), kEmbeddedMigrations.size()}; +} + +DbResult migrate_to_head(Db& db) { + const auto all = embedded_migrations(); + + MigrationOutcome out; + out.from_version = db.user_version(); + out.to_version = out.from_version; + if (out.from_version < 0) { + return std::unexpected(DbError{0, "could not read PRAGMA user_version"}); + } + + for (const auto& m : all) { + if (m.version <= out.from_version) continue; + + // Each migration is one transaction: a failure half-way leaves user_version and + // the schema exactly where they were. + auto r = db.transaction([&]() -> DbResult { + if (auto e = db.exec(m.sql); !e) return e; + return db.set_user_version(m.version); + }); + if (!r) { + return std::unexpected(DbError{r.error().code, "migration " + std::string(m.name) + + " failed: " + r.error().message}); + } + out.to_version = m.version; + ++out.applied; + } + + return out; +} + +} // namespace velox::daemon::store diff --git a/daemon/src/store/migrations.hpp b/daemon/src/store/migrations.hpp new file mode 100644 index 0000000..5d9e7c6 --- /dev/null +++ b/daemon/src/store/migrations.hpp @@ -0,0 +1,38 @@ +#pragma once + +// The schema migrator. Numbered SQL files in store/migrations/ are embedded at build time +// (see embed_migrations.cmake). On startup the daemon calls migrate_to_head(db): every +// migration whose version exceeds PRAGMA user_version is applied in order, each in its own +// transaction, and user_version is advanced to match. +// +// Forward-only: a released migration is immutable. The forward-only test in +// daemon/tests replays from every prior released user_version to head. + +#include +#include +#include + +#include "store/sqlite.hpp" + +namespace velox::daemon::store { + +struct Migration { + std::int64_t version; // 1, 2, 3, ... ; matches the NNNN_ prefix + std::string_view name; // e.g. "0001_initial" + std::string_view sql; // the file body +}; + +// The embedded set, sorted by version ascending. Defined in the generated header. +std::span embedded_migrations(); + +struct MigrationOutcome { + std::int64_t from_version = 0; + std::int64_t to_version = 0; + int applied = 0; +}; + +// Apply every migration newer than db.user_version(). A no-op (applied == 0) when the DB +// is already at or beyond the highest embedded version. +DbResult migrate_to_head(Db& db); + +} // namespace velox::daemon::store diff --git a/daemon/src/store/migrations/.gitkeep b/daemon/src/store/migrations/.gitkeep deleted file mode 100644 index e69de29..0000000 diff --git a/daemon/src/store/migrations/0001_initial.sql b/daemon/src/store/migrations/0001_initial.sql new file mode 100644 index 0000000..61f5eb2 --- /dev/null +++ b/daemon/src/store/migrations/0001_initial.sql @@ -0,0 +1,176 @@ +-- Migration 0001 — initial schema. +-- +-- Applied when PRAGMA user_version < 1. The migrator wraps this file in one transaction +-- and sets user_version = 1 on success. Forward-only: never edit a released migration, +-- add 0002_*.sql instead (AGENT-DAEMON.md build step 3). +-- +-- Conventions: +-- * ids are lowercase UUID text, except the built-in rows below. +-- * timestamps are RFC 3339 UTC strings ("2026-09-10T14:55:02Z") — same on the wire, +-- so projection to TaskSummary is a copy. +-- * JSON-valued columns hold a TEXT document; SQLite's json1 validates on read where +-- it matters. Marked "-- json" below. +-- * credentials NEVER live here (CLAUDE.md §4) — the Secret Service holds those. + +-- --- settings : the whole config bag, one row per SettingKey -------------------------- +CREATE TABLE settings ( + key TEXT PRIMARY KEY, + value TEXT NOT NULL -- json: the value as it appears in the Settings schema +) WITHOUT ROWID; + +-- --- categories : folder + extension routing ---------------------------------------- +CREATE TABLE categories ( + category_id TEXT PRIMARY KEY, + name TEXT NOT NULL, + save_dir TEXT NOT NULL, + extensions TEXT NOT NULL DEFAULT '[]', -- json array of lowercase extensions, no dot + builtin INTEGER NOT NULL DEFAULT 0 -- 1 = cannot be deleted (category.remove -> -32602) +) WITHOUT ROWID; + +INSERT INTO categories (category_id, name, save_dir, extensions, builtin) VALUES + ('general', 'General', '~/Downloads', '[]', 1), + ('programs', 'Programs', '~/Downloads/Programs', '["exe","msi","deb","rpm","dmg","appimage","iso","zip","tar","gz","xz","7z"]', 1), + ('video', 'Video', '~/Downloads/Video', '["mp4","mkv","webm","avi","mov","flv","m4v","ts"]', 1), + ('audio', 'Audio', '~/Downloads/Audio', '["mp3","flac","aac","ogg","opus","wav","m4a"]', 1), + ('documents','Documents', '~/Downloads/Documents', '["pdf","doc","docx","xls","xlsx","ppt","pptx","odt","epub"]', 1), + ('images', 'Images', '~/Downloads/Images', '["jpg","jpeg","png","gif","webp","svg","bmp","tiff"]', 1); + +-- --- queues : ordered runs with their own concurrency cap --------------------------- +CREATE TABLE queues ( + queue_id TEXT PRIMARY KEY, + name TEXT NOT NULL, + state TEXT NOT NULL DEFAULT 'stopped' -- 'running' | 'stopped' + CHECK (state IN ('running','stopped')), + max_concurrent INTEGER NOT NULL DEFAULT 2 CHECK (max_concurrent BETWEEN 1 AND 32), + on_complete TEXT NOT NULL DEFAULT 'nothing' -- 'nothing'|'exit'|'shutdown'|'hangup' + CHECK (on_complete IN ('nothing','exit','shutdown','hangup')), + schedule TEXT -- json Schedule, or NULL for manual +) WITHOUT ROWID; + +INSERT INTO queues (queue_id, name, state, max_concurrent) VALUES + ('main', 'Main Queue', 'stopped', 4); + +-- --- tasks : the download list ----------------------------------------------------- +-- Column set is chosen so a row projects onto proto TaskSummary with no computation +-- beyond reading segments/history for the detail view. +CREATE TABLE tasks ( + task_id TEXT PRIMARY KEY, + url TEXT NOT NULL, -- as supplied + effective_url TEXT, -- after redirects; NULL until first probe + filename TEXT NOT NULL DEFAULT '', + save_dir TEXT NOT NULL, -- absolute, canonicalized, inside an allowed root + category_id TEXT REFERENCES categories(category_id) ON DELETE SET NULL, + queue_id TEXT REFERENCES queues(queue_id) ON DELETE SET NULL, + queue_position INTEGER, -- NULL unless queued; run order within the queue + + state TEXT NOT NULL DEFAULT 'new' + CHECK (state IN ('new','probing','queued','connecting','downloading','paused', + 'retry_wait','assembling','verifying','complete','failed','cancelled')), + -- ADR 0013: why a paused task is paused. NULL unless state='paused'. 'auto' means CORE + -- entered it (auth_required/server_file_changed/disk_full); the code is in error_code. + pause_reason TEXT CHECK (pause_reason IN + ('user','schedule','queue_stopped','admission_reconcile','auto')), + + size_bytes INTEGER, -- NULL when the server reported no length + downloaded_bytes INTEGER NOT NULL DEFAULT 0, + resumable INTEGER NOT NULL DEFAULT 0, + + -- Requested vs effective, per ADR 0010 / ADR 0012. Requested values come from the + -- DownloadSpec; effective values are written by the engine as it runs. + req_segments INTEGER, -- DownloadSpec.segments (NULL = use setting) + eff_segments INTEGER NOT NULL DEFAULT 0,-- TaskSummary.segments (in use right now) + req_buffer_bytes INTEGER, + eff_buffer_bytes INTEGER, -- TaskDetail.effectiveBufferBytes + + start_mode TEXT NOT NULL DEFAULT 'auto' + CHECK (start_mode IN ('auto','now','queue','manual')), + description TEXT, + + -- Validators, kept for If-Range resume revalidation (docs/04 §5). + etag TEXT, + last_modified TEXT, + content_type TEXT, + + checksum_algo TEXT CHECK (checksum_algo IN ('md5','sha1','sha256','sha512')), + checksum_value TEXT, + + -- proto TaskError, flattened. Set on failed / retry_wait, and on an auto-pause. + error_code TEXT, -- TaskErrorCode string + error_message TEXT, + error_http_status INTEGER, + error_retryable INTEGER, + error_attempt INTEGER, + error_next_retry_at TEXT, + + created_at TEXT NOT NULL, + last_try_at TEXT, + completed_at TEXT +) STRICT; + +-- download.list filters/sorts in the daemon (brief: never materialize 100k rows for 40). +-- These cover the common filter columns and both default sorts. +CREATE INDEX idx_tasks_state ON tasks(state); +CREATE INDEX idx_tasks_category ON tasks(category_id); +CREATE INDEX idx_tasks_queue_order ON tasks(queue_id, queue_position); +CREATE INDEX idx_tasks_created ON tasks(created_at); +CREATE INDEX idx_tasks_completed ON tasks(completed_at); + +-- --- segments : per-connection byte ranges for one task --------------------------- +-- Inclusive ranges [start_byte, end_byte], matching HTTP Range and ADR 0010. A whole-file +-- zero-length download is one row with end_byte = start_byte - 1 = -1 (ADR 0010 B3a), so +-- end_byte is not constrained to >= 0. +CREATE TABLE segments ( + task_id TEXT NOT NULL REFERENCES tasks(task_id) ON DELETE CASCADE, + idx INTEGER NOT NULL, + start_byte INTEGER NOT NULL, + end_byte INTEGER NOT NULL, + completed_bytes INTEGER NOT NULL DEFAULT 0, + state TEXT NOT NULL DEFAULT 'connecting' + CHECK (state IN ('connecting','downloading','stalled','complete','failed')), + PRIMARY KEY (task_id, idx) +) STRICT, WITHOUT ROWID; + +-- --- rules : the routing / capture rules engine table --------------------------- +CREATE TABLE rules ( + rule_id TEXT PRIMARY KEY, + priority INTEGER NOT NULL, -- lower runs first; rules.list returns priority order + enabled INTEGER NOT NULL DEFAULT 1, + match TEXT NOT NULL, -- json: the match clause (host/ext/size/mime/...) + action TEXT NOT NULL -- json: capture decision + category + queue + start mode +) WITHOUT ROWID; + +CREATE INDEX idx_rules_priority ON rules(priority); + +-- --- history : completed and removed tasks, for the History view -------------------- +-- A task leaving the list (complete, or removed by the user) drops a snapshot here so the +-- main tasks table stays the size of the active list. +CREATE TABLE history ( + history_id INTEGER PRIMARY KEY AUTOINCREMENT, + task_id TEXT NOT NULL, + filename TEXT NOT NULL, + url TEXT NOT NULL, + save_dir TEXT NOT NULL, + size_bytes INTEGER, + final_state TEXT NOT NULL, -- 'complete' | 'cancelled' | 'failed' + category_id TEXT, + finished_at TEXT NOT NULL, + snapshot TEXT NOT NULL -- json: the full TaskSummary at the time it left +); + +CREATE INDEX idx_history_finished ON history(finished_at); +CREATE INDEX idx_history_task ON history(task_id); + +-- --- pairings : WebSocket transport tokens, HASHED (docs/05 §4, CLAUDE.md §4) -------- +-- The plaintext token is returned to the extension exactly once, from session.pair, and +-- never stored. token_sha256 is the lookup key on every subsequent connect. +CREATE TABLE pairings ( + pairing_id TEXT PRIMARY KEY, + token_sha256 TEXT NOT NULL UNIQUE, -- hex SHA-256 of the 256-bit token + origin TEXT NOT NULL, -- moz-extension://, verified on the WS upgrade + label TEXT NOT NULL DEFAULT '', -- human-readable, shown in Options -> Unpair + created_at TEXT NOT NULL, + last_seen_at TEXT, + revoked_at TEXT -- non-NULL once unpaired; kept for the audit trail +) WITHOUT ROWID; + +CREATE INDEX idx_pairings_origin ON pairings(origin); diff --git a/daemon/src/store/sqlite.cpp b/daemon/src/store/sqlite.cpp new file mode 100644 index 0000000..136d60e --- /dev/null +++ b/daemon/src/store/sqlite.cpp @@ -0,0 +1,146 @@ +#include "store/sqlite.hpp" + +#include + +#include + +namespace velox::daemon::store { + +// --- Db ------------------------------------------------------------------------------ + +Db::~Db() { + if (db_ != nullptr) sqlite3_close(db_); +} + +Db::Db(Db&& o) noexcept : db_(std::exchange(o.db_, nullptr)) {} + +Db& Db::operator=(Db&& o) noexcept { + if (this != &o) { + if (db_ != nullptr) sqlite3_close(db_); + db_ = std::exchange(o.db_, nullptr); + } + return *this; +} + +DbError Db::last_error() const { + return DbError{sqlite3_extended_errcode(db_), sqlite3_errmsg(db_)}; +} + +DbResult Db::open(const std::string& path) { + sqlite3* handle = nullptr; + const int rc = sqlite3_open_v2( + path.c_str(), &handle, SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE | SQLITE_OPEN_NOMUTEX, + nullptr); + if (rc != SQLITE_OK) { + DbError e{rc, handle != nullptr ? sqlite3_errmsg(handle) : "sqlite3_open_v2 failed"}; + if (handle != nullptr) sqlite3_close(handle); + return std::unexpected(std::move(e)); + } + + Db db(handle); + // WAL for crash-safe concurrent readers (docs/01 §1). busy_timeout so a writer waits + // rather than returning SQLITE_BUSY under the RPC loop. foreign_keys is per-connection. + for (const char* pragma : {"PRAGMA journal_mode=WAL", "PRAGMA synchronous=NORMAL", + "PRAGMA foreign_keys=ON", "PRAGMA busy_timeout=5000"}) { + if (auto r = db.exec(pragma); !r) return std::unexpected(r.error()); + } + return db; +} + +DbResult Db::exec(std::string_view sql) { + char* err = nullptr; + const int rc = sqlite3_exec(db_, std::string(sql).c_str(), nullptr, nullptr, &err); + if (rc != SQLITE_OK) { + DbError e{rc, err != nullptr ? err : sqlite3_errmsg(db_)}; + sqlite3_free(err); + return std::unexpected(std::move(e)); + } + return {}; +} + +DbResult Db::prepare(std::string_view sql) { + sqlite3_stmt* s = nullptr; + const int rc = + sqlite3_prepare_v2(db_, sql.data(), static_cast(sql.size()), &s, nullptr); + if (rc != SQLITE_OK) return std::unexpected(last_error()); + return Stmt(db_, s); +} + +std::int64_t Db::user_version() { + auto st = prepare("PRAGMA user_version"); + if (!st) return -1; + auto row = st->step(); + if (!row || !*row) return -1; + return st->column_int(0); +} + +DbResult Db::set_user_version(std::int64_t v) { + // PRAGMA does not accept a bound parameter; the value is our own integer. + return exec("PRAGMA user_version=" + std::to_string(v)); +} + +// --- Stmt ---------------------------------------------------------------------------- + +Stmt::~Stmt() { + if (stmt_ != nullptr) sqlite3_finalize(stmt_); +} + +Stmt::Stmt(Stmt&& o) noexcept + : db_(std::exchange(o.db_, nullptr)), stmt_(std::exchange(o.stmt_, nullptr)) {} + +Stmt& Stmt::operator=(Stmt&& o) noexcept { + if (this != &o) { + if (stmt_ != nullptr) sqlite3_finalize(stmt_); + db_ = std::exchange(o.db_, nullptr); + stmt_ = std::exchange(o.stmt_, nullptr); + } + return *this; +} + +DbError Stmt::last_error() const { + return DbError{sqlite3_extended_errcode(db_), sqlite3_errmsg(db_)}; +} + +DbResult Stmt::bind(int i, std::int64_t v) { + if (sqlite3_bind_int64(stmt_, i, v) != SQLITE_OK) return std::unexpected(last_error()); + return {}; +} + +DbResult Stmt::bind(int i, std::string_view v) { + if (sqlite3_bind_text(stmt_, i, v.data(), static_cast(v.size()), SQLITE_TRANSIENT) != + SQLITE_OK) + return std::unexpected(last_error()); + return {}; +} + +DbResult Stmt::bind_null(int i) { + if (sqlite3_bind_null(stmt_, i) != SQLITE_OK) return std::unexpected(last_error()); + return {}; +} + +DbResult Stmt::step() { + const int rc = sqlite3_step(stmt_); + if (rc == SQLITE_ROW) return true; + if (rc == SQLITE_DONE) return false; + return std::unexpected(last_error()); +} + +DbResult Stmt::reset() { + if (sqlite3_reset(stmt_) != SQLITE_OK) return std::unexpected(last_error()); + return {}; +} + +std::int64_t Stmt::column_int(int i) const { return sqlite3_column_int64(stmt_, i); } + +std::string Stmt::column_text(int i) const { + const auto* p = sqlite3_column_text(stmt_, i); + if (p == nullptr) return {}; + return std::string(reinterpret_cast(p), + static_cast(sqlite3_column_bytes(stmt_, i))); +} + +bool Stmt::column_is_null(int i) const { + return sqlite3_column_type(stmt_, i) == SQLITE_NULL; +} + +} // namespace velox::daemon::store diff --git a/daemon/src/store/sqlite.hpp b/daemon/src/store/sqlite.hpp new file mode 100644 index 0000000..c79899d --- /dev/null +++ b/daemon/src/store/sqlite.hpp @@ -0,0 +1,109 @@ +#pragma once + +// A thin RAII wrapper over the SQLite C API — just enough for the store: open in WAL mode, +// run statements, prepare/bind/step. Errors are returned, never thrown (the RPC loop must +// not unwind through an exception). No ORM, no query builder. + +#include +#include +#include +#include +#include + +struct sqlite3; +struct sqlite3_stmt; + +namespace velox::daemon::store { + +struct DbError { + int code = 0; // SQLite result code + std::string message; + + std::string to_string() const { + return message + " (sqlite " + std::to_string(code) + ")"; + } +}; + +template +using DbResult = std::expected; + +class Stmt; + +class Db { +public: + Db() = default; + ~Db(); + Db(Db&&) noexcept; + Db& operator=(Db&&) noexcept; + Db(const Db&) = delete; + Db& operator=(const Db&) = delete; + + // Open (creating if absent) at `path`, set WAL, busy timeout, and foreign_keys=ON. + // ":memory:" is accepted for tests. + static DbResult open(const std::string& path); + + // Run one or more statements with no result rows (DDL, PRAGMA, INSERT without + // returning). Uses sqlite3_exec, so it accepts a multi-statement script. + DbResult exec(std::string_view sql); + + DbResult prepare(std::string_view sql); + + // Convenience: run `fn` between BEGIN and COMMIT; ROLLBACK and propagate on error. + template + DbResult transaction(Fn&& fn) { + if (auto r = exec("BEGIN"); !r) return r; + auto r = std::forward(fn)(); + if (!r) { + exec("ROLLBACK"); // best effort; original error wins + return r; + } + return exec("COMMIT"); + } + + std::int64_t user_version(); + DbResult set_user_version(std::int64_t v); + + sqlite3* raw() const noexcept { return db_; } + explicit operator bool() const noexcept { return db_ != nullptr; } + +private: + explicit Db(sqlite3* db) : db_(db) {} + DbError last_error() const; + + sqlite3* db_ = nullptr; +}; + +// A prepared statement. bind_* are 1-indexed. step() returns true while rows remain. +class Stmt { +public: + Stmt() = default; + ~Stmt(); + Stmt(Stmt&&) noexcept; + Stmt& operator=(Stmt&&) noexcept; + Stmt(const Stmt&) = delete; + Stmt& operator=(const Stmt&) = delete; + + DbResult bind(int i, std::int64_t v); + DbResult bind(int i, std::string_view v); + DbResult bind_null(int i); + + // true: a row is available; false: done. Any error code is surfaced via error(). + DbResult step(); + DbResult reset(); + + std::int64_t column_int(int i) const; + std::string column_text(int i) const; + bool column_is_null(int i) const; + + sqlite3_stmt* raw() const noexcept { return stmt_; } + +private: + friend class Db; + explicit Stmt(sqlite3* db, sqlite3_stmt* s) : db_(db), stmt_(s) {} + DbError last_error() const; + + sqlite3* db_ = nullptr; // borrowed, for error messages + sqlite3_stmt* stmt_ = nullptr; +}; + +} // namespace velox::daemon::store diff --git a/daemon/tests/CMakeLists.txt b/daemon/tests/CMakeLists.txt index 04983f7..fa8ee78 100644 --- a/daemon/tests/CMakeLists.txt +++ b/daemon/tests/CMakeLists.txt @@ -12,3 +12,8 @@ target_link_libraries(veloxd_uds_roundtrip_test PRIVATE veloxd_rpc) target_compile_options(veloxd_uds_roundtrip_test PRIVATE -Wall -Wextra -Wpedantic -Werror) add_test(NAME veloxd.uds_roundtrip COMMAND veloxd_uds_roundtrip_test) set_tests_properties(veloxd.uds_roundtrip PROPERTIES TIMEOUT 30) + +add_executable(veloxd_store_migrations_test store_migrations_test.cpp) +target_link_libraries(veloxd_store_migrations_test PRIVATE veloxd_store) +target_compile_options(veloxd_store_migrations_test PRIVATE -Wall -Wextra -Wpedantic -Werror) +add_test(NAME veloxd.store_migrations COMMAND veloxd_store_migrations_test) diff --git a/daemon/tests/store_migrations_test.cpp b/daemon/tests/store_migrations_test.cpp new file mode 100644 index 0000000..50971aa --- /dev/null +++ b/daemon/tests/store_migrations_test.cpp @@ -0,0 +1,118 @@ +// The migrator: fresh DB -> head, idempotent re-run, and forward-only from every released +// user_version (M1 DoD: "a forward-only test from every released schema version"). + +#include + +#include "check.hpp" +#include "store/migrations.hpp" +#include "store/sqlite.hpp" + +using namespace velox::daemon::store; + +namespace { + +std::int64_t head_version() { + std::int64_t v = 0; + for (const auto& m : embedded_migrations()) v = std::max(v, m.version); + return v; +} + +bool table_exists(Db& db, const char* name) { + auto st = db.prepare("SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1"); + if (!st) return false; + if (!st->bind(1, std::string_view(name))) return false; + auto row = st->step(); + return row && *row; +} + +std::int64_t count(Db& db, const char* sql) { + auto st = db.prepare(sql); + if (!st) return -1; + auto row = st->step(); + if (!row || !*row) return -1; + return st->column_int(0); +} + +} // namespace + +void run() { + const std::int64_t head = head_version(); + CHECK(head >= 1); + + // --- fresh in-memory DB migrates cleanly to head -------------------------------- + { + auto db = Db::open(":memory:"); + CHECK(db.has_value()); + if (!db) return; + CHECK_EQ(db->user_version(), 0); + + auto out = migrate_to_head(*db); + CHECK(out.has_value()); + if (out) { + CHECK_EQ(out->from_version, 0); + CHECK_EQ(out->to_version, head); + CHECK_EQ(static_cast(out->applied), head); + } + CHECK_EQ(db->user_version(), head); + + for (const char* t : {"settings", "categories", "queues", "tasks", "segments", + "rules", "history", "pairings"}) { + CHECK(table_exists(*db, t)); + } + // Seed rows the initial migration inserts. + CHECK_EQ(count(*db, "SELECT count(*) FROM categories WHERE builtin=1"), 6); + CHECK_EQ(count(*db, "SELECT count(*) FROM queues"), 1); + + // FK + cascade wired: a segment for a missing task is rejected; deleting a task + // takes its segments with it. + CHECK(db->exec("INSERT INTO tasks(task_id,url,save_dir,created_at) " + "VALUES('t1','http://x','/tmp','2026-09-10T00:00:00Z')") + .has_value()); + CHECK(db->exec("INSERT INTO segments(task_id,idx,start_byte,end_byte) " + "VALUES('t1',0,0,99)") + .has_value()); + CHECK(!db->exec("INSERT INTO segments(task_id,idx,start_byte,end_byte) " + "VALUES('nope',0,0,99)") + .has_value()); + CHECK(db->exec("DELETE FROM tasks WHERE task_id='t1'").has_value()); + CHECK_EQ(count(*db, "SELECT count(*) FROM segments"), 0); + + // A whole-file zero-length download: one segment, end_byte = -1 (ADR 0010 B3a). + CHECK(db->exec("INSERT INTO tasks(task_id,url,save_dir,created_at,size_bytes) " + "VALUES('z','http://x','/tmp','2026-09-10T00:00:00Z',0)") + .has_value()); + CHECK(db->exec("INSERT INTO segments(task_id,idx,start_byte,end_byte) " + "VALUES('z',0,0,-1)") + .has_value()); + } + + // --- re-running the migrator on an at-head DB is a no-op ------------------------ + { + auto db = Db::open(":memory:"); + CHECK(db.has_value()); + (void)migrate_to_head(*db); + auto again = migrate_to_head(*db); + CHECK(again.has_value()); + if (again) { + CHECK_EQ(again->applied, 0); + CHECK_EQ(again->to_version, head); + } + } + + // --- forward-only: from every released version [0 .. head-1], reach head -------- + for (std::int64_t start = 0; start < head; ++start) { + auto db = Db::open(":memory:"); + CHECK(db.has_value()); + if (!db) continue; + CHECK(db->set_user_version(start).has_value()); + auto out = migrate_to_head(*db); + CHECK(out.has_value()); + if (out) { + CHECK_EQ(out->from_version, start); + CHECK_EQ(out->to_version, head); + } + CHECK_EQ(db->user_version(), head); + } +} + +TEST_MAIN()