From 9dccf99c50d331cf65f503a0db96b81faea195e0 Mon Sep 17 00:00:00 2001 From: Hare Date: Tue, 11 Aug 2026 04:03:54 +0900 Subject: [PATCH] ticket: own typed schema migrations --- crates/ticket/src/lib.rs | 194 +--- crates/ticket/src/sqlite_schema.rs | 1077 ++++++++++++++++++++++ crates/workspace-server/src/authority.rs | 5 +- crates/workspace-server/src/server.rs | 9 +- crates/workspace-server/src/store.rs | 40 + crates/yoi/src/ticket_cli.rs | 4 +- 6 files changed, 1171 insertions(+), 158 deletions(-) create mode 100644 crates/ticket/src/sqlite_schema.rs diff --git a/crates/ticket/src/lib.rs b/crates/ticket/src/lib.rs index e613a841..29f6f244 100644 --- a/crates/ticket/src/lib.rs +++ b/crates/ticket/src/lib.rs @@ -20,8 +20,13 @@ use serde_yaml::{Mapping as YamlMapping, Value as YamlValue}; use thiserror::Error; pub mod config; +mod sqlite_schema; pub mod tool; +pub use sqlite_schema::{ + LATEST_SQLITE_TICKET_SCHEMA_VERSION, migrate_sqlite_ticket_schema, verify_sqlite_ticket_schema, +}; + const REQUIRED_FIELDS: [&str; 4] = ["title", "state", "created_at", "updated_at"]; const MAX_STATE_CHANGE_REASON_BYTES: usize = 1024; const MAX_INTAKE_SUMMARY_BODY_BYTES: usize = 16 * 1024; @@ -2301,7 +2306,7 @@ impl fmt::Debug for SqliteTicketBackend { } impl SqliteTicketBackend { - pub fn new(db_path: impl Into, workspace_id: impl Into) -> Self { + fn configured(db_path: impl Into, workspace_id: impl Into) -> Self { Self { db_path: db_path.into(), workspace_id: workspace_id.into(), @@ -2311,6 +2316,27 @@ impl SqliteTicketBackend { } } + /// Opens a standalone Ticket backend, applying all Ticket-owned migrations once. + pub fn open(db_path: impl Into, workspace_id: impl Into) -> Result { + let backend = Self::configured(db_path, workspace_id); + let connection = backend.connect()?; + migrate_sqlite_ticket_schema(&connection)?; + Ok(backend) + } + + /// Connects to a database whose Ticket schema was composed by its startup owner. + /// + /// This performs verification only and never creates or alters schema objects. + pub fn open_verified( + db_path: impl Into, + workspace_id: impl Into, + ) -> Result { + let backend = Self::configured(db_path, workspace_id); + let connection = backend.connect()?; + verify_sqlite_ticket_schema(&connection)?; + Ok(backend) + } + pub fn with_event_attributes(mut self, attributes: BTreeMap) -> Self { self.event_attributes = attributes; self @@ -2338,7 +2364,6 @@ impl SqliteTicketBackend { pub fn import_from_local_backend(&self, local: &LocalTicketBackend) -> Result<()> { let conn = self.open_connection()?; - self.ensure_schema(&conn)?; conn.execute_batch("BEGIN IMMEDIATE").map_err(sqlite_err)?; let result = (|| { for summary in local.list(TicketListQuery::all())? { @@ -2351,7 +2376,7 @@ impl SqliteTicketBackend { finish_sqlite_transaction(&conn, result) } - fn open_connection(&self) -> Result { + fn connect(&self) -> Result { if let Some(parent) = self.db_path.parent() { fs::create_dir_all(parent).map_err(|error| io_err(parent, error))?; } @@ -2361,115 +2386,20 @@ impl SqliteTicketBackend { Ok(conn) } - fn ensure_schema(&self, conn: &Connection) -> Result<()> { - conn.execute_batch(r#" -CREATE TABLE IF NOT EXISTS typed_tickets ( - workspace_id TEXT NOT NULL, - ticket_id TEXT NOT NULL, - slug TEXT NOT NULL, - title TEXT NOT NULL, - status TEXT NOT NULL, - kind TEXT NOT NULL, - priority TEXT NOT NULL, - body TEXT NOT NULL, - created_at TEXT, - updated_at TEXT, - assignee TEXT, - readiness TEXT, - workflow_state TEXT NOT NULL, - workflow_state_explicit INTEGER NOT NULL, - queued_by TEXT, - queued_at TEXT, - resolution TEXT, - repository_id TEXT, - ref_selector TEXT, - PRIMARY KEY (workspace_id, ticket_id) -); -CREATE TABLE IF NOT EXISTS typed_ticket_labels ( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, ordinal INTEGER NOT NULL, label TEXT NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, ordinal), - FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE -); -CREATE TABLE IF NOT EXISTS typed_ticket_risk_flags ( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, ordinal INTEGER NOT NULL, risk_flag TEXT NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, ordinal), - FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE -); -CREATE TABLE IF NOT EXISTS typed_ticket_raw_frontmatter ( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, key TEXT NOT NULL, value TEXT NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, key), - FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE -); -CREATE TABLE IF NOT EXISTS typed_ticket_events ( - workspace_id TEXT NOT NULL, - ticket_id TEXT NOT NULL, - event_index INTEGER NOT NULL, - kind TEXT NOT NULL, - author TEXT, - at TEXT, - status TEXT, - from_state TEXT, - to_state TEXT, - reason TEXT, - state_field TEXT, - heading TEXT, - body TEXT NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, event_index), - FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE -); -CREATE TABLE IF NOT EXISTS typed_ticket_event_references ( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, event_index INTEGER NOT NULL, ordinal INTEGER NOT NULL, kind TEXT NOT NULL, target TEXT NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, event_index, ordinal), - FOREIGN KEY (workspace_id, ticket_id, event_index) REFERENCES typed_ticket_events(workspace_id, ticket_id, event_index) ON DELETE CASCADE -); -CREATE TABLE IF NOT EXISTS typed_ticket_event_attributes ( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, event_index INTEGER NOT NULL, key TEXT NOT NULL, value TEXT NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, event_index, key), - FOREIGN KEY (workspace_id, ticket_id, event_index) REFERENCES typed_ticket_events(workspace_id, ticket_id, event_index) ON DELETE CASCADE -); -CREATE TABLE IF NOT EXISTS typed_ticket_relations ( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, kind TEXT NOT NULL, target TEXT NOT NULL, note TEXT, author TEXT NOT NULL, at TEXT NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, kind, target), - FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE -); -CREATE TABLE IF NOT EXISTS typed_ticket_orchestration_plans ( - workspace_id TEXT NOT NULL, - ticket_id TEXT NOT NULL, - record_id TEXT NOT NULL, - kind TEXT NOT NULL, - related_ticket TEXT, - note TEXT, - accepted_summary TEXT, - accepted_branch TEXT, - accepted_worktree TEXT, - accepted_role_plan TEXT, - author TEXT NOT NULL, - at TEXT NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, record_id), - FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE -); -CREATE TABLE IF NOT EXISTS typed_ticket_artifacts ( - workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, relative_path TEXT NOT NULL, content BLOB NOT NULL, - PRIMARY KEY (workspace_id, ticket_id, relative_path), - FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE -); -"#) - .map_err(sqlite_err)?; - ensure_sqlite_ticket_column(conn, "repository_id", "TEXT")?; - ensure_sqlite_ticket_column(conn, "ref_selector", "TEXT")?; - Ok(()) + fn open_connection(&self) -> Result { + let connection = self.connect()?; + verify_sqlite_ticket_schema(&connection)?; + Ok(connection) } fn with_write(&self, op: impl FnOnce(&Connection) -> Result) -> Result { let conn = self.open_connection()?; - self.ensure_schema(&conn)?; conn.execute_batch("BEGIN IMMEDIATE").map_err(sqlite_err)?; finish_sqlite_transaction(&conn, op(&conn)) } fn with_read(&self, op: impl FnOnce(&Connection) -> Result) -> Result { let conn = self.open_connection()?; - self.ensure_schema(&conn)?; op(&conn) } @@ -2915,30 +2845,6 @@ CREATE TABLE IF NOT EXISTS typed_ticket_artifacts ( } } -fn ensure_sqlite_ticket_column( - conn: &rusqlite::Connection, - name: &str, - sql_type: &str, -) -> Result<()> { - let mut statement = conn - .prepare("PRAGMA table_info(typed_tickets)") - .map_err(sqlite_err)?; - let columns = statement - .query_map([], |row| row.get::<_, String>(1)) - .map_err(sqlite_err)?; - for column in columns { - if column.map_err(sqlite_err)? == name { - return Ok(()); - } - } - conn.execute( - format!("ALTER TABLE typed_tickets ADD COLUMN {name} {sql_type}").as_str(), - [], - ) - .map_err(sqlite_err)?; - Ok(()) -} - fn finish_sqlite_transaction(conn: &Connection, result: Result) -> Result { match result { Ok(output) => { @@ -6406,28 +6312,11 @@ state: planning #[test] fn sqlite_backend_persists_and_edits_ticket_target() { let tmp = TempDir::new().unwrap(); - let backend = SqliteTicketBackend::new(tmp.path().join("workspace.db"), "workspace-test"); + let backend = + SqliteTicketBackend::open(tmp.path().join("workspace.db"), "workspace-test").unwrap(); assert_ticket_target_edit_semantics(&backend); } - #[test] - fn sqlite_ticket_target_columns_are_added_to_existing_table() { - let tmp = TempDir::new().unwrap(); - let conn = rusqlite::Connection::open(tmp.path().join("workspace.db")).unwrap(); - conn.execute_batch("CREATE TABLE typed_tickets (ticket_id TEXT PRIMARY KEY);") - .unwrap(); - ensure_sqlite_ticket_column(&conn, "repository_id", "TEXT").unwrap(); - ensure_sqlite_ticket_column(&conn, "ref_selector", "TEXT").unwrap(); - let mut statement = conn.prepare("PRAGMA table_info(typed_tickets)").unwrap(); - let columns = statement - .query_map([], |row| row.get::<_, String>(1)) - .unwrap() - .collect::, _>>() - .unwrap(); - assert!(columns.iter().any(|column| column == "repository_id")); - assert!(columns.iter().any(|column| column == "ref_selector")); - } - #[test] fn local_backend_edit_item_supports_partial_body_replacement() { let tmp = TempDir::new().unwrap(); @@ -6438,7 +6327,8 @@ state: planning #[test] fn sqlite_backend_edit_item_supports_partial_body_replacement() { let tmp = TempDir::new().unwrap(); - let backend = SqliteTicketBackend::new(tmp.path().join("workspace.db"), "workspace-test"); + let backend = + SqliteTicketBackend::open(tmp.path().join("workspace.db"), "workspace-test").unwrap(); assert_partial_body_replacement_semantics(&backend); } @@ -6446,7 +6336,7 @@ state: planning fn sqlite_mutation_hook_failure_rolls_back_ticket_event() { let tmp = TempDir::new().unwrap(); let db_path = tmp.path().join("workspace.db"); - let backend = SqliteTicketBackend::new(&db_path, "workspace-test"); + let backend = SqliteTicketBackend::open(&db_path, "workspace-test").unwrap(); let created = backend.create(NewTicket::new("Atomic mutation")).unwrap(); let before = backend .show(TicketIdOrSlug::Id(created.id.clone())) @@ -6480,7 +6370,8 @@ state: planning #[test] fn sqlite_backend_persists_core_ticket_operations() { let tmp = TempDir::new().unwrap(); - let backend = SqliteTicketBackend::new(tmp.path().join("workspace.db"), "workspace-test"); + let backend = + SqliteTicketBackend::open(tmp.path().join("workspace.db"), "workspace-test").unwrap(); let created = backend.create(NewTicket::new("SQLite Ticket")).unwrap(); backend .add_event( @@ -6501,7 +6392,9 @@ state: planning ) .unwrap(); - let reopened = SqliteTicketBackend::new(tmp.path().join("workspace.db"), "workspace-test"); + let reopened = + SqliteTicketBackend::open_verified(tmp.path().join("workspace.db"), "workspace-test") + .unwrap(); let list = reopened.list(TicketListQuery::all()).unwrap(); assert_eq!(list.len(), 1); assert_eq!(list[0].id, created.id); @@ -6531,7 +6424,8 @@ state: planning let tmp = TempDir::new().unwrap(); let local = backend(&tmp); let created = local.create(NewTicket::new("Legacy Ticket")).unwrap(); - let db = SqliteTicketBackend::new(tmp.path().join("workspace.db"), "workspace-test"); + let db = + SqliteTicketBackend::open(tmp.path().join("workspace.db"), "workspace-test").unwrap(); db.import_from_local_backend(&local).unwrap(); let ticket = db.show(TicketIdOrSlug::Id(created.id.clone())).unwrap(); diff --git a/crates/ticket/src/sqlite_schema.rs b/crates/ticket/src/sqlite_schema.rs new file mode 100644 index 00000000..87536232 --- /dev/null +++ b/crates/ticket/src/sqlite_schema.rs @@ -0,0 +1,1077 @@ +use std::collections::{BTreeMap, BTreeSet}; +use std::time::Duration; + +use rusqlite::{Connection, OptionalExtension, params}; + +use crate::{Result, TicketError, sqlite_err}; + +const MIGRATION_TABLE: &str = "ticket_schema_migrations"; +const MAX_SCHEMA_DIAGNOSTICS: usize = 32; +pub const LATEST_SQLITE_TICKET_SCHEMA_VERSION: i64 = 2; + +#[derive(Clone, Copy)] +struct Migration { + version: i64, + name: &'static str, + apply: fn(&Connection) -> Result<()>, +} + +const MIGRATIONS: &[Migration] = &[ + Migration { + version: 1, + name: "create_typed_ticket_tables", + apply: create_typed_ticket_tables, + }, + Migration { + version: 2, + name: "add_ticket_repository_target", + apply: add_ticket_repository_target, + }, +]; + +#[derive(Clone, Copy)] +struct ExpectedColumn { + name: &'static str, + data_type: &'static str, + not_null: bool, + primary_key_position: i64, +} + +#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] +struct ExpectedForeignKey { + from: &'static str, + target_table: &'static str, + to: &'static str, + on_delete: &'static str, +} + +const OWNED_TABLES: &[&str] = &[ + "typed_tickets", + "typed_ticket_labels", + "typed_ticket_risk_flags", + "typed_ticket_raw_frontmatter", + "typed_ticket_events", + "typed_ticket_event_references", + "typed_ticket_event_attributes", + "typed_ticket_relations", + "typed_ticket_orchestration_plans", + "typed_ticket_artifacts", +]; + +const MIGRATION_COLUMNS: &[ExpectedColumn] = &[ + column("version", "INTEGER", false, 1), + column("name", "TEXT", true, 0), + column("applied_at", "TEXT", true, 0), +]; + +const TICKET_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("slug", "TEXT", true, 0), + column("title", "TEXT", true, 0), + column("status", "TEXT", true, 0), + column("kind", "TEXT", true, 0), + column("priority", "TEXT", true, 0), + column("body", "TEXT", true, 0), + column("created_at", "TEXT", false, 0), + column("updated_at", "TEXT", false, 0), + column("assignee", "TEXT", false, 0), + column("readiness", "TEXT", false, 0), + column("workflow_state", "TEXT", true, 0), + column("workflow_state_explicit", "INTEGER", true, 0), + column("queued_by", "TEXT", false, 0), + column("queued_at", "TEXT", false, 0), + column("resolution", "TEXT", false, 0), + column("repository_id", "TEXT", false, 0), + column("ref_selector", "TEXT", false, 0), +]; + +const LABEL_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("ordinal", "INTEGER", true, 3), + column("label", "TEXT", true, 0), +]; + +const RISK_FLAG_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("ordinal", "INTEGER", true, 3), + column("risk_flag", "TEXT", true, 0), +]; + +const RAW_FRONTMATTER_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("key", "TEXT", true, 3), + column("value", "TEXT", true, 0), +]; + +const EVENT_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("event_index", "INTEGER", true, 3), + column("kind", "TEXT", true, 0), + column("author", "TEXT", false, 0), + column("at", "TEXT", false, 0), + column("status", "TEXT", false, 0), + column("from_state", "TEXT", false, 0), + column("to_state", "TEXT", false, 0), + column("reason", "TEXT", false, 0), + column("state_field", "TEXT", false, 0), + column("heading", "TEXT", false, 0), + column("body", "TEXT", true, 0), +]; + +const EVENT_REFERENCE_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("event_index", "INTEGER", true, 3), + column("ordinal", "INTEGER", true, 4), + column("kind", "TEXT", true, 0), + column("target", "TEXT", true, 0), +]; + +const EVENT_ATTRIBUTE_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("event_index", "INTEGER", true, 3), + column("key", "TEXT", true, 4), + column("value", "TEXT", true, 0), +]; + +const RELATION_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("kind", "TEXT", true, 3), + column("target", "TEXT", true, 4), + column("note", "TEXT", false, 0), + column("author", "TEXT", true, 0), + column("at", "TEXT", true, 0), +]; + +const ORCHESTRATION_PLAN_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("record_id", "TEXT", true, 3), + column("kind", "TEXT", true, 0), + column("related_ticket", "TEXT", false, 0), + column("note", "TEXT", false, 0), + column("accepted_summary", "TEXT", false, 0), + column("accepted_branch", "TEXT", false, 0), + column("accepted_worktree", "TEXT", false, 0), + column("accepted_role_plan", "TEXT", false, 0), + column("author", "TEXT", true, 0), + column("at", "TEXT", true, 0), +]; + +const ARTIFACT_COLUMNS: &[ExpectedColumn] = &[ + column("workspace_id", "TEXT", true, 1), + column("ticket_id", "TEXT", true, 2), + column("relative_path", "TEXT", true, 3), + column("content", "BLOB", true, 0), +]; + +const TICKET_FOREIGN_KEYS: &[ExpectedForeignKey] = &[]; +const CHILD_FOREIGN_KEYS: &[ExpectedForeignKey] = &[ + ExpectedForeignKey { + from: "workspace_id", + target_table: "typed_tickets", + to: "workspace_id", + on_delete: "CASCADE", + }, + ExpectedForeignKey { + from: "ticket_id", + target_table: "typed_tickets", + to: "ticket_id", + on_delete: "CASCADE", + }, +]; +const EVENT_CHILD_FOREIGN_KEYS: &[ExpectedForeignKey] = &[ + ExpectedForeignKey { + from: "workspace_id", + target_table: "typed_ticket_events", + to: "workspace_id", + on_delete: "CASCADE", + }, + ExpectedForeignKey { + from: "ticket_id", + target_table: "typed_ticket_events", + to: "ticket_id", + on_delete: "CASCADE", + }, + ExpectedForeignKey { + from: "event_index", + target_table: "typed_ticket_events", + to: "event_index", + on_delete: "CASCADE", + }, +]; + +const fn column( + name: &'static str, + data_type: &'static str, + not_null: bool, + primary_key_position: i64, +) -> ExpectedColumn { + ExpectedColumn { + name, + data_type, + not_null, + primary_key_position, + } +} + +/// Applies the Ticket crate's SQLite migrations and verifies the resulting schema. +/// +/// This is a startup/standalone-open operation. Normal Ticket request handling must +/// use [`verify_sqlite_ticket_schema`] instead, so request paths never acquire DDL +/// authority. +pub fn migrate_sqlite_ticket_schema(connection: &Connection) -> Result<()> { + connection + .busy_timeout(Duration::from_secs(5)) + .map_err(sqlite_err)?; + connection + .execute_batch("BEGIN IMMEDIATE") + .map_err(sqlite_err)?; + + let result = (|| { + connection + .execute_batch( + "CREATE TABLE IF NOT EXISTS ticket_schema_migrations ( + version INTEGER PRIMARY KEY, + name TEXT NOT NULL, + applied_at TEXT NOT NULL + );", + ) + .map_err(sqlite_err)?; + verify_table(connection, MIGRATION_TABLE, MIGRATION_COLUMNS, &[], false)?; + + let applied = load_applied_migrations(connection)?; + validate_applied_migrations(&applied)?; + + for migration in MIGRATIONS { + if applied.contains_key(&migration.version) { + continue; + } + (migration.apply)(connection)?; + connection + .execute( + "INSERT INTO ticket_schema_migrations (version, name, applied_at) + VALUES (?1, ?2, ?3)", + params![ + migration.version, + migration.name, + chrono::Utc::now().to_rfc3339() + ], + ) + .map_err(sqlite_err)?; + } + + verify_sqlite_ticket_schema(connection) + })(); + + match result { + Ok(()) => connection.execute_batch("COMMIT").map_err(sqlite_err), + Err(error) => { + let _ = connection.execute_batch("ROLLBACK"); + Err(error) + } + } +} + +/// Verifies the current Ticket-owned SQLite schema without executing DDL. +pub fn verify_sqlite_ticket_schema(connection: &Connection) -> Result<()> { + let mut diagnostics = Vec::new(); + + collect_table_diagnostics( + connection, + MIGRATION_TABLE, + MIGRATION_COLUMNS, + &[], + &mut diagnostics, + ); + + match load_applied_migrations(connection) { + Ok(applied) => { + if let Err(error) = validate_applied_migrations(&applied) { + push_diagnostic(&mut diagnostics, error.to_string()); + } else if applied.len() != MIGRATIONS.len() { + push_diagnostic( + &mut diagnostics, + format!( + "Ticket schema is not current: found {} migration(s), expected {}", + applied.len(), + MIGRATIONS.len() + ), + ); + } + } + Err(error) => push_diagnostic(&mut diagnostics, error.to_string()), + } + + for (table, columns, foreign_keys) in [ + ("typed_tickets", TICKET_COLUMNS, TICKET_FOREIGN_KEYS), + ("typed_ticket_labels", LABEL_COLUMNS, CHILD_FOREIGN_KEYS), + ( + "typed_ticket_risk_flags", + RISK_FLAG_COLUMNS, + CHILD_FOREIGN_KEYS, + ), + ( + "typed_ticket_raw_frontmatter", + RAW_FRONTMATTER_COLUMNS, + CHILD_FOREIGN_KEYS, + ), + ("typed_ticket_events", EVENT_COLUMNS, CHILD_FOREIGN_KEYS), + ( + "typed_ticket_event_references", + EVENT_REFERENCE_COLUMNS, + EVENT_CHILD_FOREIGN_KEYS, + ), + ( + "typed_ticket_event_attributes", + EVENT_ATTRIBUTE_COLUMNS, + EVENT_CHILD_FOREIGN_KEYS, + ), + ( + "typed_ticket_relations", + RELATION_COLUMNS, + CHILD_FOREIGN_KEYS, + ), + ( + "typed_ticket_orchestration_plans", + ORCHESTRATION_PLAN_COLUMNS, + CHILD_FOREIGN_KEYS, + ), + ( + "typed_ticket_artifacts", + ARTIFACT_COLUMNS, + CHILD_FOREIGN_KEYS, + ), + ] { + collect_table_diagnostics(connection, table, columns, foreign_keys, &mut diagnostics); + } + + for table in OWNED_TABLES { + collect_foreign_key_check_diagnostics(connection, table, &mut diagnostics); + } + + if diagnostics.is_empty() { + Ok(()) + } else { + let was_truncated = diagnostics.len() > MAX_SCHEMA_DIAGNOSTICS; + diagnostics.truncate(MAX_SCHEMA_DIAGNOSTICS); + let mut message = format!( + "Ticket SQLite schema verification failed: {}", + diagnostics.join("; ") + ); + if was_truncated { + message.push_str("; additional diagnostics omitted"); + } + Err(TicketError::Sqlite(message)) + } +} + +fn create_typed_ticket_tables(connection: &Connection) -> Result<()> { + connection + .execute_batch( + r#" +CREATE TABLE IF NOT EXISTS typed_tickets ( + workspace_id TEXT NOT NULL, + ticket_id TEXT NOT NULL, + slug TEXT NOT NULL, + title TEXT NOT NULL, + status TEXT NOT NULL, + kind TEXT NOT NULL, + priority TEXT NOT NULL, + body TEXT NOT NULL, + created_at TEXT, + updated_at TEXT, + assignee TEXT, + readiness TEXT, + workflow_state TEXT NOT NULL, + workflow_state_explicit INTEGER NOT NULL, + queued_by TEXT, + queued_at TEXT, + resolution TEXT, + PRIMARY KEY (workspace_id, ticket_id) +); +CREATE TABLE IF NOT EXISTS typed_ticket_labels ( + workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, ordinal INTEGER NOT NULL, label TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, ordinal), + FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE +); +CREATE TABLE IF NOT EXISTS typed_ticket_risk_flags ( + workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, ordinal INTEGER NOT NULL, risk_flag TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, ordinal), + FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE +); +CREATE TABLE IF NOT EXISTS typed_ticket_raw_frontmatter ( + workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, key TEXT NOT NULL, value TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, key), + FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE +); +CREATE TABLE IF NOT EXISTS typed_ticket_events ( + workspace_id TEXT NOT NULL, + ticket_id TEXT NOT NULL, + event_index INTEGER NOT NULL, + kind TEXT NOT NULL, + author TEXT, + at TEXT, + status TEXT, + from_state TEXT, + to_state TEXT, + reason TEXT, + state_field TEXT, + heading TEXT, + body TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, event_index), + FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE +); +CREATE TABLE IF NOT EXISTS typed_ticket_event_references ( + workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, event_index INTEGER NOT NULL, ordinal INTEGER NOT NULL, kind TEXT NOT NULL, target TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, event_index, ordinal), + FOREIGN KEY (workspace_id, ticket_id, event_index) REFERENCES typed_ticket_events(workspace_id, ticket_id, event_index) ON DELETE CASCADE +); +CREATE TABLE IF NOT EXISTS typed_ticket_event_attributes ( + workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, event_index INTEGER NOT NULL, key TEXT NOT NULL, value TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, event_index, key), + FOREIGN KEY (workspace_id, ticket_id, event_index) REFERENCES typed_ticket_events(workspace_id, ticket_id, event_index) ON DELETE CASCADE +); +CREATE TABLE IF NOT EXISTS typed_ticket_relations ( + workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, kind TEXT NOT NULL, target TEXT NOT NULL, note TEXT, author TEXT NOT NULL, at TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, kind, target), + FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE +); +CREATE TABLE IF NOT EXISTS typed_ticket_orchestration_plans ( + workspace_id TEXT NOT NULL, + ticket_id TEXT NOT NULL, + record_id TEXT NOT NULL, + kind TEXT NOT NULL, + related_ticket TEXT, + note TEXT, + accepted_summary TEXT, + accepted_branch TEXT, + accepted_worktree TEXT, + accepted_role_plan TEXT, + author TEXT NOT NULL, + at TEXT NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, record_id), + FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE +); +CREATE TABLE IF NOT EXISTS typed_ticket_artifacts ( + workspace_id TEXT NOT NULL, ticket_id TEXT NOT NULL, relative_path TEXT NOT NULL, content BLOB NOT NULL, + PRIMARY KEY (workspace_id, ticket_id, relative_path), + FOREIGN KEY (workspace_id, ticket_id) REFERENCES typed_tickets(workspace_id, ticket_id) ON DELETE CASCADE +); +"#, + ) + .map_err(sqlite_err) +} + +fn add_ticket_repository_target(connection: &Connection) -> Result<()> { + add_column_if_missing(connection, "typed_tickets", "repository_id", "TEXT")?; + add_column_if_missing(connection, "typed_tickets", "ref_selector", "TEXT") +} + +fn add_column_if_missing( + connection: &Connection, + table: &str, + column: &str, + declaration: &str, +) -> Result<()> { + let columns = load_columns(connection, table)?; + if columns.iter().any(|found| found.name == column) { + return Ok(()); + } + connection + .execute_batch(&format!( + "ALTER TABLE {table} ADD COLUMN {column} {declaration}" + )) + .map_err(sqlite_err) +} + +fn load_applied_migrations(connection: &Connection) -> Result> { + let mut statement = connection + .prepare("SELECT version, name FROM ticket_schema_migrations ORDER BY version") + .map_err(sqlite_err)?; + let rows = statement + .query_map([], |row| { + Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)) + }) + .map_err(sqlite_err)?; + let mut applied = BTreeMap::new(); + for row in rows { + let (version, name) = row.map_err(sqlite_err)?; + if applied.insert(version, name).is_some() { + return Err(TicketError::Sqlite(format!( + "duplicate Ticket schema migration version {version}" + ))); + } + } + Ok(applied) +} + +fn validate_applied_migrations(applied: &BTreeMap) -> Result<()> { + for (&version, name) in applied { + let Some(expected) = MIGRATIONS + .iter() + .find(|migration| migration.version == version) + else { + return Err(TicketError::Sqlite(format!( + "unsupported Ticket schema migration version {version}; latest supported version is {LATEST_SQLITE_TICKET_SCHEMA_VERSION}" + ))); + }; + if name != expected.name { + return Err(TicketError::Sqlite(format!( + "Ticket schema migration {version} is named {name:?}, expected {:?}", + expected.name + ))); + } + } + for migration in MIGRATIONS { + if applied.keys().any(|version| *version > migration.version) + && !applied.contains_key(&migration.version) + { + return Err(TicketError::Sqlite(format!( + "Ticket schema migration history has a gap at version {}", + migration.version + ))); + } + } + Ok(()) +} + +#[derive(Debug)] +struct ActualColumn { + name: String, + data_type: String, + not_null: bool, + primary_key_position: i64, +} + +fn load_columns(connection: &Connection, table: &str) -> Result> { + let exists = connection + .query_row( + "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1", + params![table], + |_| Ok(()), + ) + .optional() + .map_err(sqlite_err)? + .is_some(); + if !exists { + return Err(TicketError::Sqlite(format!( + "required Ticket schema table {table:?} is missing" + ))); + } + + let mut statement = connection + .prepare(&format!("PRAGMA table_info({table})")) + .map_err(sqlite_err)?; + let rows = statement + .query_map([], |row| { + Ok(ActualColumn { + name: row.get(1)?, + data_type: row.get::<_, String>(2)?.to_ascii_uppercase(), + not_null: row.get::<_, i64>(3)? != 0, + primary_key_position: row.get(5)?, + }) + }) + .map_err(sqlite_err)?; + rows.collect::, _>>() + .map_err(sqlite_err) +} + +fn verify_table( + connection: &Connection, + table: &str, + expected_columns: &[ExpectedColumn], + expected_foreign_keys: &[ExpectedForeignKey], + check_foreign_keys: bool, +) -> Result<()> { + let mut diagnostics = Vec::new(); + collect_column_diagnostics(connection, table, expected_columns, &mut diagnostics); + if check_foreign_keys { + collect_foreign_key_diagnostics(connection, table, expected_foreign_keys, &mut diagnostics); + } + if diagnostics.is_empty() { + Ok(()) + } else { + Err(TicketError::Sqlite(diagnostics.join("; "))) + } +} + +fn collect_table_diagnostics( + connection: &Connection, + table: &str, + expected_columns: &[ExpectedColumn], + expected_foreign_keys: &[ExpectedForeignKey], + diagnostics: &mut Vec, +) { + collect_column_diagnostics(connection, table, expected_columns, diagnostics); + if diagnostics.len() < MAX_SCHEMA_DIAGNOSTICS { + collect_foreign_key_diagnostics(connection, table, expected_foreign_keys, diagnostics); + } +} + +fn collect_column_diagnostics( + connection: &Connection, + table: &str, + expected: &[ExpectedColumn], + diagnostics: &mut Vec, +) { + let actual = match load_columns(connection, table) { + Ok(actual) => actual, + Err(error) => { + push_diagnostic(diagnostics, error.to_string()); + return; + } + }; + + let expected_names = expected + .iter() + .map(|column| column.name) + .collect::>(); + let actual_names = actual + .iter() + .map(|column| column.name.as_str()) + .collect::>(); + for missing in expected_names.difference(&actual_names) { + push_diagnostic( + diagnostics, + format!("table {table:?} is missing column {missing:?}"), + ); + } + for unexpected in actual_names.difference(&expected_names) { + push_diagnostic( + diagnostics, + format!("table {table:?} has unexpected column {unexpected:?}"), + ); + } + + for expected in expected { + let Some(actual) = actual.iter().find(|column| column.name == expected.name) else { + continue; + }; + if actual.data_type != expected.data_type { + push_diagnostic( + diagnostics, + format!( + "table {table:?} column {:?} has type {:?}, expected {:?}", + expected.name, actual.data_type, expected.data_type + ), + ); + } + if actual.not_null != expected.not_null { + push_diagnostic( + diagnostics, + format!( + "table {table:?} column {:?} NOT NULL is {}, expected {}", + expected.name, actual.not_null, expected.not_null + ), + ); + } + if actual.primary_key_position != expected.primary_key_position { + push_diagnostic( + diagnostics, + format!( + "table {table:?} column {:?} primary-key position is {}, expected {}", + expected.name, actual.primary_key_position, expected.primary_key_position + ), + ); + } + } +} + +fn collect_foreign_key_diagnostics( + connection: &Connection, + table: &str, + expected: &[ExpectedForeignKey], + diagnostics: &mut Vec, +) { + let mut statement = match connection.prepare(&format!("PRAGMA foreign_key_list({table})")) { + Ok(statement) => statement, + Err(error) => { + push_diagnostic(diagnostics, sqlite_err(error).to_string()); + return; + } + }; + let rows = match statement.query_map([], |row| { + Ok(( + row.get::<_, String>(3)?, + row.get::<_, String>(2)?, + row.get::<_, String>(4)?, + row.get::<_, String>(6)?.to_ascii_uppercase(), + )) + }) { + Ok(rows) => rows, + Err(error) => { + push_diagnostic(diagnostics, sqlite_err(error).to_string()); + return; + } + }; + let mut actual = BTreeSet::new(); + for row in rows { + match row { + Ok(foreign_key) => { + actual.insert(foreign_key); + } + Err(error) => push_diagnostic(diagnostics, sqlite_err(error).to_string()), + } + } + let expected = expected + .iter() + .map(|foreign_key| { + ( + foreign_key.from.to_string(), + foreign_key.target_table.to_string(), + foreign_key.to.to_string(), + foreign_key.on_delete.to_string(), + ) + }) + .collect::>(); + for missing in expected.difference(&actual) { + push_diagnostic( + diagnostics, + format!("table {table:?} is missing foreign key {missing:?}"), + ); + } + for unexpected in actual.difference(&expected) { + push_diagnostic( + diagnostics, + format!("table {table:?} has unexpected foreign key {unexpected:?}"), + ); + } +} + +fn collect_foreign_key_check_diagnostics( + connection: &Connection, + table: &str, + diagnostics: &mut Vec, +) { + let mut statement = match connection.prepare(&format!("PRAGMA foreign_key_check({table})")) { + Ok(statement) => statement, + Err(error) => { + push_diagnostic(diagnostics, sqlite_err(error).to_string()); + return; + } + }; + let rows = match statement.query_map([], |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, Option>(1)?, + row.get::<_, String>(2)?, + row.get::<_, i64>(3)?, + )) + }) { + Ok(rows) => rows, + Err(error) => { + push_diagnostic(diagnostics, sqlite_err(error).to_string()); + return; + } + }; + for row in rows { + match row { + Ok((table, row_id, parent, foreign_key_id)) => push_diagnostic( + diagnostics, + format!( + "foreign-key violation in table {table:?} row {row_id:?} referencing {parent:?} (foreign key {foreign_key_id})" + ), + ), + Err(error) => push_diagnostic(diagnostics, sqlite_err(error).to_string()), + } + } +} + +fn push_diagnostic(diagnostics: &mut Vec, diagnostic: String) { + if diagnostics.len() <= MAX_SCHEMA_DIAGNOSTICS { + diagnostics.push(diagnostic); + } +} + +#[cfg(test)] +mod tests { + use std::sync::{Arc, Barrier}; + use std::thread; + + use rusqlite::Connection; + use tempfile::tempdir; + + use super::*; + use crate::SqliteTicketBackend; + + #[test] + fn migrates_fresh_database_to_current_ticket_schema() { + let connection = Connection::open_in_memory().unwrap(); + migrate_sqlite_ticket_schema(&connection).unwrap(); + verify_sqlite_ticket_schema(&connection).unwrap(); + + let versions = load_applied_migrations(&connection).unwrap(); + assert_eq!(versions.len(), 2); + assert_eq!( + versions.get(&LATEST_SQLITE_TICKET_SCHEMA_VERSION), + Some(&"add_ticket_repository_target".to_string()) + ); + } + + #[test] + fn adopts_existing_current_schema_without_losing_data() { + let connection = Connection::open_in_memory().unwrap(); + create_typed_ticket_tables(&connection).unwrap(); + add_ticket_repository_target(&connection).unwrap(); + connection + .execute( + "INSERT INTO typed_tickets ( + workspace_id, ticket_id, slug, title, status, kind, priority, body, + workflow_state, workflow_state_explicit, repository_id, ref_selector + ) VALUES ('workspace-1', 'ticket-1', 'ticket-1', 'kept', 'open', + 'task', 'medium', 'body', 'ready', 1, 'main', 'develop')", + [], + ) + .unwrap(); + connection + .execute_batch( + "INSERT INTO typed_ticket_events ( + workspace_id, ticket_id, event_index, kind, author, at, heading, body + ) VALUES ( + 'workspace-1', 'ticket-1', 0, 'comment', 'hare', + '2026-08-10T00:00:00Z', 'Evidence', 'event kept' + ); + INSERT INTO typed_ticket_event_references ( + workspace_id, ticket_id, event_index, ordinal, kind, target + ) VALUES ('workspace-1', 'ticket-1', 0, 0, 'commit', 'abc123'); + INSERT INTO typed_ticket_relations ( + workspace_id, ticket_id, kind, target, note, author, at + ) VALUES ( + 'workspace-1', 'ticket-1', 'related', 'ticket-2', 'relation kept', + 'hare', '2026-08-10T00:00:00Z' + ); + INSERT INTO typed_ticket_orchestration_plans ( + workspace_id, ticket_id, record_id, kind, note, author, at + ) VALUES ( + 'workspace-1', 'ticket-1', 'plan-1', 'waiting_capacity_note', + 'plan kept', 'hare', '2026-08-10T00:00:00Z' + ); + INSERT INTO typed_ticket_artifacts ( + workspace_id, ticket_id, relative_path, content + ) VALUES ('workspace-1', 'ticket-1', 'evidence.txt', X'6b657074');", + ) + .unwrap(); + + migrate_sqlite_ticket_schema(&connection).unwrap(); + + let row = connection + .query_row( + "SELECT title, repository_id, ref_selector FROM typed_tickets", + [], + |row| { + Ok(( + row.get::<_, String>(0)?, + row.get::<_, String>(1)?, + row.get::<_, String>(2)?, + )) + }, + ) + .unwrap(); + assert_eq!(row, ("kept".into(), "main".into(), "develop".into())); + let preserved = connection + .query_row( + "SELECT + (SELECT COUNT(*) FROM typed_ticket_events), + (SELECT COUNT(*) FROM typed_ticket_event_references), + (SELECT COUNT(*) FROM typed_ticket_relations), + (SELECT COUNT(*) FROM typed_ticket_orchestration_plans), + (SELECT COUNT(*) FROM typed_ticket_artifacts)", + [], + |row| { + Ok(( + row.get::<_, i64>(0)?, + row.get::<_, i64>(1)?, + row.get::<_, i64>(2)?, + row.get::<_, i64>(3)?, + row.get::<_, i64>(4)?, + )) + }, + ) + .unwrap(); + assert_eq!(preserved, (1, 1, 1, 1, 1)); + } + + #[test] + fn upgrades_legacy_schema_without_repository_target_columns() { + let connection = Connection::open_in_memory().unwrap(); + create_typed_ticket_tables(&connection).unwrap(); + connection + .execute( + "INSERT INTO typed_tickets ( + workspace_id, ticket_id, slug, title, status, kind, priority, body, + workflow_state, workflow_state_explicit + ) VALUES ('workspace-1', 'ticket-1', 'ticket-1', 'legacy', 'open', + 'task', 'medium', 'body', 'ready', 1)", + [], + ) + .unwrap(); + connection + .execute_batch( + "CREATE TABLE ticket_schema_migrations ( + version INTEGER PRIMARY KEY, + name TEXT NOT NULL, + applied_at TEXT NOT NULL + ); + INSERT INTO ticket_schema_migrations (version, name, applied_at) + VALUES (1, 'create_typed_ticket_tables', '2026-08-10T00:00:00Z');", + ) + .unwrap(); + + migrate_sqlite_ticket_schema(&connection).unwrap(); + verify_sqlite_ticket_schema(&connection).unwrap(); + + let columns = load_columns(&connection, "typed_tickets").unwrap(); + assert!(columns.iter().any(|column| column.name == "repository_id")); + assert!(columns.iter().any(|column| column.name == "ref_selector")); + let title = connection + .query_row("SELECT title FROM typed_tickets", [], |row| { + row.get::<_, String>(0) + }) + .unwrap(); + assert_eq!(title, "legacy"); + } + + #[test] + fn rejects_unknown_future_migration_history_without_changing_schema() { + let connection = Connection::open_in_memory().unwrap(); + migrate_sqlite_ticket_schema(&connection).unwrap(); + connection + .execute( + "INSERT INTO ticket_schema_migrations (version, name, applied_at) + VALUES (99, 'future', '2026-08-10T00:00:00Z')", + [], + ) + .unwrap(); + + let error = migrate_sqlite_ticket_schema(&connection).unwrap_err(); + assert!( + error + .to_string() + .contains("unsupported Ticket schema migration version 99") + ); + assert_eq!(load_applied_migrations(&connection).unwrap().len(), 3); + } + + #[test] + fn verified_backend_open_fails_on_drift_without_repairing_request_schema() { + let directory = tempdir().unwrap(); + let database = directory.path().join("tickets.db"); + SqliteTicketBackend::open(&database, "workspace-1").unwrap(); + let connection = Connection::open(&database).unwrap(); + connection + .execute_batch("DROP TABLE typed_ticket_artifacts") + .unwrap(); + + let error = match SqliteTicketBackend::open_verified(&database, "workspace-1") { + Ok(_) => panic!("drifted schema unexpectedly passed request-path verification"), + Err(error) => error, + }; + assert!(error.to_string().contains("typed_ticket_artifacts")); + let still_missing = connection + .query_row( + "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'typed_ticket_artifacts'", + [], + |_| Ok(()), + ) + .optional() + .unwrap() + .is_none(); + assert!(still_missing); + } + + #[test] + fn verification_does_not_claim_unrelated_foreign_key_authority() { + let connection = Connection::open_in_memory().unwrap(); + migrate_sqlite_ticket_schema(&connection).unwrap(); + connection + .pragma_update(None, "foreign_keys", "OFF") + .unwrap(); + connection + .execute_batch( + "CREATE TABLE unrelated_parent (id TEXT PRIMARY KEY); + CREATE TABLE unrelated_child ( + id TEXT PRIMARY KEY, + parent_id TEXT NOT NULL REFERENCES unrelated_parent(id) + ); + INSERT INTO unrelated_child (id, parent_id) VALUES ('child', 'missing');", + ) + .unwrap(); + + verify_sqlite_ticket_schema(&connection).unwrap(); + } + + #[test] + fn migration_rejects_constraint_drift_and_rolls_back_version_adoption() { + let connection = Connection::open_in_memory().unwrap(); + connection + .execute_batch( + "CREATE TABLE typed_tickets ( + workspace_id TEXT NOT NULL, + ticket_id TEXT NOT NULL, + slug TEXT NOT NULL, + title TEXT NOT NULL, + status TEXT NOT NULL, + kind TEXT NOT NULL, + priority TEXT NOT NULL, + body TEXT NOT NULL, + created_at TEXT, + updated_at TEXT, + assignee TEXT, + readiness TEXT, + workflow_state TEXT NOT NULL, + workflow_state_explicit INTEGER NOT NULL, + queued_by TEXT, + queued_at TEXT, + resolution TEXT, + PRIMARY KEY (ticket_id, workspace_id) + );", + ) + .unwrap(); + + let error = migrate_sqlite_ticket_schema(&connection).unwrap_err(); + assert!(error.to_string().contains("primary-key position")); + let migration_table_exists = connection + .query_row( + "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'ticket_schema_migrations'", + [], + |_| Ok(()), + ) + .optional() + .unwrap() + .is_some(); + assert!(!migration_table_exists); + } + + #[test] + fn concurrent_migrators_converge_on_one_version_history() { + let directory = tempdir().unwrap(); + let database = directory.path().join("tickets.db"); + let barrier = Arc::new(Barrier::new(3)); + let mut joins = Vec::new(); + for _ in 0..2 { + let database = database.clone(); + let barrier = barrier.clone(); + joins.push(thread::spawn(move || { + let connection = Connection::open(database).unwrap(); + barrier.wait(); + migrate_sqlite_ticket_schema(&connection) + })); + } + barrier.wait(); + for join in joins { + join.join().unwrap().unwrap(); + } + + let connection = Connection::open(database).unwrap(); + verify_sqlite_ticket_schema(&connection).unwrap(); + assert_eq!(load_applied_migrations(&connection).unwrap().len(), 2); + } +} diff --git a/crates/workspace-server/src/authority.rs b/crates/workspace-server/src/authority.rs index 528b955e..e982bb67 100644 --- a/crates/workspace-server/src/authority.rs +++ b/crates/workspace-server/src/authority.rs @@ -131,7 +131,7 @@ impl SqliteWorkspaceAuthority { Ok(Self { workspace_id: workspace_id.clone(), store: SqliteWorkspaceStore::open(&database_path)?, - ticket_backend: SqliteTicketBackend::new(database_path, workspace_id), + ticket_backend: SqliteTicketBackend::open_verified(database_path, workspace_id)?, }) } @@ -725,7 +725,8 @@ mod tests { let dir = tempfile::tempdir().unwrap(); write_ticket(dir.path(), "00000000001J2", "Read bridge", "ready"); let db_path = dir.path().join("workspace.db"); - SqliteTicketBackend::new(&db_path, "workspace-test") + SqliteTicketBackend::open(&db_path, "workspace-test") + .unwrap() .import_from_local_backend(&ticket::LocalTicketBackend::new( dir.path().join(".yoi/tickets"), )) diff --git a/crates/workspace-server/src/server.rs b/crates/workspace-server/src/server.rs index 6652821e..fdc60a38 100644 --- a/crates/workspace-server/src/server.rs +++ b/crates/workspace-server/src/server.rs @@ -2383,10 +2383,10 @@ struct BrowserCloseTicketRequest { fn browser_ticket_backend(api: &WorkspaceApi) -> Result { let config = ticket::config::TicketConfig::load_workspace(&api.config.workspace_root) .map_err(|error| Error::Config(format!("load Ticket workspace settings: {error}")))?; - Ok(SqliteTicketBackend::new( + Ok(SqliteTicketBackend::open_verified( api.config.database_path.clone(), api.config.workspace_id.clone(), - ) + )? .with_record_language(config.ticket_record_language())) } @@ -2553,10 +2553,11 @@ async fn execute_worker_ticket_rest_operation( validate_workspace_scope(api, workspace_id)?; let config = ticket::config::TicketConfig::load_workspace(&api.config.workspace_root) .map_err(|error| Error::Config(format!("load Ticket workspace settings: {error}")))?; - let mut backend = SqliteTicketBackend::new( + let mut backend = SqliteTicketBackend::open_verified( api.config.database_path.clone(), api.config.workspace_id.clone(), ) + .map_err(Error::from)? .with_record_language(config.ticket_record_language()); let operation_kind = ticket_mutation_operation_kind(&operation); let is_mutation = operation_kind != "read"; @@ -15225,7 +15226,7 @@ mod tests { ) { use ticket::TicketBackend as _; - let backend = ticket::SqliteTicketBackend::new(database_path, workspace_id); + let backend = ticket::SqliteTicketBackend::open(database_path, workspace_id).unwrap(); let mut input = ticket::NewTicket::new(title); input.workflow_state = Some(state); backend.create(input).unwrap(); diff --git a/crates/workspace-server/src/store.rs b/crates/workspace-server/src/store.rs index c180414f..0b3617a6 100644 --- a/crates/workspace-server/src/store.rs +++ b/crates/workspace-server/src/store.rs @@ -754,6 +754,7 @@ impl SqliteWorkspaceStore { pub fn from_connection(conn: Connection) -> Result { configure_sqlite(&conn)?; apply_migrations(&conn)?; + ticket::migrate_sqlite_ticket_schema(&conn)?; Ok(Self { conn: Arc::new(Mutex::new(conn)), }) @@ -4526,6 +4527,45 @@ mod tests { use super::*; use std::collections::BTreeSet; + #[test] + fn startup_composes_ticket_migrations_when_control_plane_is_current() { + let conn = Connection::open_in_memory().unwrap(); + configure_sqlite(&conn).unwrap(); + apply_migrations(&conn).unwrap(); + assert!(!table_exists(&conn, "ticket_schema_migrations").unwrap()); + + let store = SqliteWorkspaceStore::from_connection(conn).unwrap(); + store + .with_conn(|conn| { + ticket::verify_sqlite_ticket_schema(conn)?; + let latest = conn.query_row( + "SELECT MAX(version) FROM ticket_schema_migrations", + [], + |row| row.get::<_, i64>(0), + )?; + assert_eq!(latest, ticket::LATEST_SQLITE_TICKET_SCHEMA_VERSION); + Ok(()) + }) + .unwrap(); + } + + #[test] + fn startup_fails_closed_when_current_ticket_schema_has_drifted() { + let conn = Connection::open_in_memory().unwrap(); + configure_sqlite(&conn).unwrap(); + apply_migrations(&conn).unwrap(); + ticket::migrate_sqlite_ticket_schema(&conn).unwrap(); + conn.execute_batch("DROP TABLE typed_ticket_artifacts") + .unwrap(); + + let result = SqliteWorkspaceStore::from_connection(conn); + let error = match result { + Ok(_) => panic!("schema drift unexpectedly passed startup verification"), + Err(error) => error, + }; + assert!(error.to_string().contains("typed_ticket_artifacts")); + } + #[test] fn removes_unused_control_plane_ticket_tables() { let conn = Connection::open_in_memory().unwrap(); diff --git a/crates/yoi/src/ticket_cli.rs b/crates/yoi/src/ticket_cli.rs index 487f9e4f..c51ca8a8 100644 --- a/crates/yoi/src/ticket_cli.rs +++ b/crates/yoi/src/ticket_cli.rs @@ -370,7 +370,7 @@ fn backend_for_workspace(workspace: &Path) -> Result, Tic let workspace_id = workspace_id_for_workspace(workspace)?; let db_path = server_database_path(workspace)?; Ok(Box::new( - SqliteTicketBackend::new(db_path, workspace_id) + SqliteTicketBackend::open(db_path, workspace_id)? .with_record_language(config.ticket_record_language()), )) } @@ -381,7 +381,7 @@ fn import_local(workspace: &Path) -> Result { .with_record_language(config.ticket_record_language()); let workspace_id = workspace_id_for_workspace(workspace)?; let db_path = server_database_path(workspace)?; - let sqlite = SqliteTicketBackend::new(db_path.clone(), workspace_id) + let sqlite = SqliteTicketBackend::open(db_path.clone(), workspace_id)? .with_record_language(config.ticket_record_language()); sqlite.import_from_local_backend(&local)?; Ok(success(format!(