From 40a3d164b2d71fb6cfe2de01aa57ecac18e53a26 Mon Sep 17 00:00:00 2001 From: Rod Boev Date: Thu, 3 Sep 2026 22:50:50 -0400 Subject: [PATCH 1/4] Derive current-schema column expectations from the applied schema --- docs/internal/storage-design.md | 4 + internal/store/upgrade.go | 127 ++++++++++++++++++++++++-------- internal/store/upgrade_test.go | 96 +++++++++++++++++++++++- 3 files changed, 194 insertions(+), 33 deletions(-) diff --git a/docs/internal/storage-design.md b/docs/internal/storage-design.md index f16eaf11..91134716 100644 --- a/docs/internal/storage-design.md +++ b/docs/internal/storage-design.md @@ -486,6 +486,10 @@ Store startup runs the embedded idempotent schema in one immediate transaction and ensures the root exists. This safely creates missing compatible tables and indexes, but it is not a general migration system. +The embedded `schema.sql` is the sole authority for current-layout columns. +Startup derives the six guarded table layouts by applying it to an isolated +temporary database. Released layouts remain pinned by their versioned adapters. + Metadata-v1 identity changes are vertical changes to the live store, ingest, reachability, and backup/restore paths. A parallel schema or codec that production code does not consume has no authority; shared metadata helpers diff --git a/internal/store/upgrade.go b/internal/store/upgrade.go index 7ba9df7c..817430dc 100644 --- a/internal/store/upgrade.go +++ b/internal/store/upgrade.go @@ -10,6 +10,7 @@ import ( "path/filepath" "slices" "strings" + "sync" "go.kenn.io/kit/pack" @@ -70,6 +71,21 @@ var ( removeInvalidUpgradeStage = removeUpgradeFileSet ) +type currentSchemaColumnsCacheKey struct { + driverName string + schemaSQL string +} + +var currentSchemaColumnsCache = struct { + sync.Mutex + + values map[currentSchemaColumnsCacheKey]map[string][]string +}{values: make(map[currentSchemaColumnsCacheKey]map[string][]string)} + +var currentSchemaTables = [...]string{ + "blobs", "blob_packs", "vault_metadata", "blob_stores", "blob_locations", "blob_pack_entries", +} + // prepareReleasedSchemaUpgrade recognizes only storage layouts that shipped in // a public release. Older layouts rebuild through the same deterministic JSONL // authority used by backup and restore; released schemas are never mutated in @@ -103,7 +119,7 @@ func prepareReleasedSchemaUpgrade(path string, driver docsqlite.Driver) error { if err != nil { return fmt.Errorf("inspecting database schema with %s: %w", driver.Name(), err) } - kind, classifyErr := classifyDatabaseSchema(db) + kind, classifyErr := classifyDatabaseSchema(driver, db) closeErr := db.Close() if classifyErr != nil || closeErr != nil { return errors.Join(classifyErr, closeErr) @@ -146,7 +162,7 @@ func validateReleasedStorageSchemas() error { return nil } -func classifyDatabaseSchema(db *sql.DB) (databaseSchema, error) { +func classifyDatabaseSchema(driver docsqlite.Driver, db *sql.DB) (databaseSchema, error) { blobs, err := tableColumns(db, "blobs") if err != nil { return databaseSchema{}, err @@ -175,7 +191,7 @@ func classifyDatabaseSchema(db *sql.DB) (databaseSchema, error) { ) } if version == currentStorageSchemaVersion { - if err := validateCurrentSchemaColumns(db, blobs, packs, vaultMetadata); err != nil { + if err := validateCurrentSchemaColumns(driver, db, blobs, packs, vaultMetadata); err != nil { return databaseSchema{}, err } return databaseSchema{version: version, current: true}, nil @@ -214,45 +230,27 @@ func classifyDatabaseSchema(db *sql.DB) (databaseSchema, error) { } func validateCurrentSchemaColumns( - db *sql.DB, + driver docsqlite.Driver, db *sql.DB, blobs, packs, vaultMetadata []string, ) error { - currentBlobColumns := []string{ - metadataCreatedAtField, "hash", metadataSizeField, - } - currentPackColumns := []string{ - metadataCreatedAtField, "entry_count", "live_entries", "live_raw_bytes", "live_stored_bytes", - "max_live_raw_len", "max_live_stored_len", "pack_id", "scan_hash", "store_id", "stored_bytes", + wantTables, err := canonicalCurrentSchemaColumns(driver) + if err != nil { + return fmt.Errorf("deriving current schema columns: %w", err) } - currentVaultColumns := []string{"schema_version", "singleton", "vault_uid"} - if !slices.Equal(blobs, currentBlobColumns) || !slices.Equal(packs, currentPackColumns) || - !slices.Equal(vaultMetadata, currentVaultColumns) { + if !slices.Equal(blobs, wantTables["blobs"]) || !slices.Equal(packs, wantTables["blob_packs"]) || + !slices.Equal(vaultMetadata, wantTables["vault_metadata"]) { return fmt.Errorf( "opening database: schema version %d has an unexpected layout (blobs=%s blob_packs=%s vault_metadata=%s)", currentStorageSchemaVersion, strings.Join(blobs, ","), strings.Join(packs, ","), strings.Join(vaultMetadata, ","), ) } - wantTables := map[string][]string{ - "blob_stores": { - "binding", metadataCreatedAtField, "kind", "lifecycle", "name", - "ownership_epoch", "role", "store_id", - }, - "blob_locations": { - columnBlobHash, "encoding", "generation", "kind", "pack_eligible", - "store_id", "stored_size", - }, - "blob_pack_entries": { - columnBlobHash, "crc32c", "flags", "pack_id", "pack_offset", - "raw_len", "store_id", "stored_len", - }, - } - for table, want := range wantTables { + for _, table := range []string{"blob_stores", "blob_locations", "blob_pack_entries"} { got, err := tableColumns(db, table) if err != nil { return err } - if !slices.Equal(got, want) { + if !slices.Equal(got, wantTables[table]) { return fmt.Errorf( "opening database: schema version %d has an unexpected %s layout (%s)", currentStorageSchemaVersion, table, strings.Join(got, ","), @@ -262,6 +260,75 @@ func validateCurrentSchemaColumns( return nil } +func canonicalCurrentSchemaColumns(driver docsqlite.Driver) (map[string][]string, error) { + key := currentSchemaColumnsCacheKey{driverName: driver.Name(), schemaSQL: schemaSQL} + currentSchemaColumnsCache.Lock() + columns, ok := currentSchemaColumnsCache.values[key] + currentSchemaColumnsCache.Unlock() + if ok { + return cloneSchemaColumns(columns), nil + } + + columns, err := deriveCurrentSchemaColumns(driver) + if err != nil { + return nil, err + } + currentSchemaColumnsCache.Lock() + if cached, ok := currentSchemaColumnsCache.values[key]; ok { + currentSchemaColumnsCache.Unlock() + return cloneSchemaColumns(cached), nil + } + currentSchemaColumnsCache.values[key] = cloneSchemaColumns(columns) + currentSchemaColumnsCache.Unlock() + return cloneSchemaColumns(columns), nil +} + +func deriveCurrentSchemaColumns(driver docsqlite.Driver) (columns map[string][]string, err error) { + tempFile, err := os.CreateTemp("", "docbank-current-schema-*.db") + if err != nil { + return nil, fmt.Errorf("creating temporary database: %w", err) + } + tempPath := tempFile.Name() + var db *sql.DB + defer func() { + var closeErr error + if db != nil { + closeErr = db.Close() + } + err = errors.Join(err, closeErr, os.Remove(tempPath)) + }() + if err := tempFile.Close(); err != nil { + return nil, fmt.Errorf("closing temporary database file: %w", err) + } + + db, err = driver.Open(tempPath, docsqlite.OpenOptions{ + Access: docsqlite.Create, TransactionMode: docsqlite.Immediate, + }) + if err != nil { + return nil, fmt.Errorf("opening temporary database with %s: %w", driver.Name(), err) + } + if _, err := db.Exec(schemaSQL); err != nil { + return nil, fmt.Errorf("applying current schema to temporary database: %w", err) + } + columns = make(map[string][]string, len(currentSchemaTables)) + for _, table := range currentSchemaTables { + columns[table], err = tableColumns(db, table) + if err != nil { + return nil, err + } + slices.Sort(columns[table]) + } + return columns, nil +} + +func cloneSchemaColumns(columns map[string][]string) map[string][]string { + clone := make(map[string][]string, len(columns)) + for table, names := range columns { + clone[table] = slices.Clone(names) + } + return clone +} + func validateV2Schema(db *sql.DB, blobs, packs []string) error { v2BlobColumns := []string{ metadataCreatedAtField, "hash", "loose_encoding", "loose_stored_size", @@ -1004,7 +1071,7 @@ func validateUpgradeStage(path string, driver docsqlite.Driver) error { if err != nil { return fmt.Errorf("opening interrupted upgrade staging database: %w", err) } - kind, classifyErr := classifyDatabaseSchema(db) + kind, classifyErr := classifyDatabaseSchema(driver, db) closeErr := db.Close() if classifyErr != nil || closeErr != nil { return errors.Join(classifyErr, closeErr) diff --git a/internal/store/upgrade_test.go b/internal/store/upgrade_test.go index cc7fd5c4..b8b7444b 100644 --- a/internal/store/upgrade_test.go +++ b/internal/store/upgrade_test.go @@ -9,6 +9,7 @@ import ( "errors" "os" "path/filepath" + "strings" "testing" "github.com/stretchr/testify/assert" @@ -206,6 +207,95 @@ func TestFreshStoresRecordCurrentStorageSchemaVersion(t *testing.T) { } } +func TestOpenAcceptsCurrentSchemaColumnAddedToEmbeddedSchema(t *testing.T) { + originalSchema := schemaSQL + t.Cleanup(func() { schemaSQL = originalSchema }) + + for _, table := range currentSchemaTables { + for _, test := range v090UpgradeDrivers() { + t.Run(table+"/"+test.name, func(t *testing.T) { + schemaSQL = schemaSQLWithAddedColumn(t, originalSchema, table, "synthetic_schema_269") + + dbPath := filepath.Join(t.TempDir(), "docbank.db") + s, err := Open(dbPath, test.driver) + require.NoError(t, err) + require.NoError(t, s.Close()) + + reopened, err := Open(dbPath, test.driver) + require.NoError(t, err) + require.NoError(t, reopened.Close()) + }) + } + } +} + +func TestCanonicalCurrentSchemaDerivationDoesNotCacheErrors(t *testing.T) { + originalSchema := schemaSQL + t.Cleanup(func() { schemaSQL = originalSchema }) + schemaSQL = schemaSQLWithAddedColumn(t, originalSchema, "blobs", "synthetic_derivation_retry_269") + driver := &flakySchemaDriver{Driver: DefaultSQLiteDriver()} + + _, err := canonicalCurrentSchemaColumns(driver) + require.ErrorContains(t, err, "synthetic derivation failure") + + columns, err := canonicalCurrentSchemaColumns(driver) + require.NoError(t, err) + assert.Contains(t, columns["blobs"], "synthetic_derivation_retry_269") +} + +type flakySchemaDriver struct { + docsqlite.Driver + + failed bool +} + +func (d *flakySchemaDriver) Open(path string, opts docsqlite.OpenOptions) (*sql.DB, error) { + if !d.failed { + d.failed = true + return nil, errors.New("synthetic derivation failure") + } + return d.Driver.Open(path, opts) +} + +func schemaSQLWithAddedColumn(t *testing.T, original, table, column string) string { + t.Helper() + lineEnding := "\n" + if strings.Contains(original, "\r\n") { + lineEnding = "\r\n" + } + result := strings.Replace( + original, + "CREATE TABLE IF NOT EXISTS "+table+" ("+lineEnding, + "CREATE TABLE IF NOT EXISTS "+table+" ("+lineEnding+ + " "+column+" TEXT,"+lineEnding, + 1, + ) + require.NotEqual(t, original, result) + return result +} + +func TestOpenRejectsCurrentDatabaseWithForeignColumn(t *testing.T) { + for _, test := range v090UpgradeDrivers() { + t.Run(test.name, func(t *testing.T) { + dbPath := filepath.Join(t.TempDir(), "docbank.db") + s, err := Open(dbPath, test.driver) + require.NoError(t, err) + require.NoError(t, s.Close()) + + db, err := test.driver.Open(dbPath, docsqlite.OpenOptions{ + Access: docsqlite.ReadWriteExisting, TransactionMode: docsqlite.Immediate, + }) + require.NoError(t, err) + _, err = db.Exec(`ALTER TABLE blobs ADD COLUMN synthetic_unexpected_269 TEXT`) + require.NoError(t, err) + require.NoError(t, db.Close()) + + _, err = Open(dbPath, test.driver) + require.ErrorContains(t, err, "schema version 4 has an unexpected layout") + }) + } +} + func TestOpenCutsOverEveryReleasedSchemaV2LayoutThroughJSONL(t *testing.T) { layouts := []struct { name string @@ -261,7 +351,7 @@ func TestOpenCutsOverEveryReleasedSchemaV2LayoutThroughJSONL(t *testing.T) { Access: docsqlite.ReadWriteExisting, TransactionMode: docsqlite.Deferred, }) require.NoError(t, err) - kind, err := classifyDatabaseSchema(backup) + kind, err := classifyDatabaseSchema(driver.driver, backup) require.NoError(t, err) assert.Equal(t, 2, kind.version) assert.NotNil(t, kind.source) @@ -300,7 +390,7 @@ func TestOpenCutsOverReleasedSchemaV3ThroughJSONL(t *testing.T) { Access: docsqlite.ReadWriteExisting, TransactionMode: docsqlite.Deferred, }) require.NoError(t, err) - kind, err := classifyDatabaseSchema(backup) + kind, err := classifyDatabaseSchema(test.driver, backup) require.NoError(t, err) assert.Equal(t, 3, kind.version) assert.NotNil(t, kind.source) @@ -622,7 +712,7 @@ func TestV090CutoverPublicationFailureRestoresReleasedDatabase(t *testing.T) { Access: docsqlite.ReadWriteExisting, TransactionMode: docsqlite.Deferred, }) require.NoError(t, err) - kind, err := classifyDatabaseSchema(db) + kind, err := classifyDatabaseSchema(driver, db) require.NoError(t, err) assert.Equal(t, 1, kind.version, "the released source is restored after publication fails") assert.NotNil(t, kind.source) From 08061ee1e175ec48a336e37b1d6373f450d77e85 Mon Sep 17 00:00:00 2001 From: Rod Boev Date: Thu, 3 Sep 2026 23:00:21 -0400 Subject: [PATCH 2/4] Cover current storage layout rejection --- internal/store/upgrade_test.go | 39 ++++++++++++++++++++-------------- 1 file changed, 23 insertions(+), 16 deletions(-) diff --git a/internal/store/upgrade_test.go b/internal/store/upgrade_test.go index b8b7444b..039c9f93 100644 --- a/internal/store/upgrade_test.go +++ b/internal/store/upgrade_test.go @@ -275,24 +275,31 @@ func schemaSQLWithAddedColumn(t *testing.T, original, table, column string) stri } func TestOpenRejectsCurrentDatabaseWithForeignColumn(t *testing.T) { - for _, test := range v090UpgradeDrivers() { - t.Run(test.name, func(t *testing.T) { - dbPath := filepath.Join(t.TempDir(), "docbank.db") - s, err := Open(dbPath, test.driver) - require.NoError(t, err) - require.NoError(t, s.Close()) + for _, table := range []struct { + name, expected string + }{ + {name: "blobs", expected: "schema version 4 has an unexpected layout"}, + {name: "blob_locations", expected: "schema version 4 has an unexpected blob_locations layout"}, + } { + for _, test := range v090UpgradeDrivers() { + t.Run(table.name+"/"+test.name, func(t *testing.T) { + dbPath := filepath.Join(t.TempDir(), "docbank.db") + s, err := Open(dbPath, test.driver) + require.NoError(t, err) + require.NoError(t, s.Close()) - db, err := test.driver.Open(dbPath, docsqlite.OpenOptions{ - Access: docsqlite.ReadWriteExisting, TransactionMode: docsqlite.Immediate, - }) - require.NoError(t, err) - _, err = db.Exec(`ALTER TABLE blobs ADD COLUMN synthetic_unexpected_269 TEXT`) - require.NoError(t, err) - require.NoError(t, db.Close()) + db, err := test.driver.Open(dbPath, docsqlite.OpenOptions{ + Access: docsqlite.ReadWriteExisting, TransactionMode: docsqlite.Immediate, + }) + require.NoError(t, err) + _, err = db.Exec(`ALTER TABLE ` + table.name + ` ADD COLUMN synthetic_unexpected_269 TEXT`) + require.NoError(t, err) + require.NoError(t, db.Close()) - _, err = Open(dbPath, test.driver) - require.ErrorContains(t, err, "schema version 4 has an unexpected layout") - }) + _, err = Open(dbPath, test.driver) + require.ErrorContains(t, err, table.expected) + }) + } } } From 20aac351c712df1a5d36b9a0411102df8c11b158 Mon Sep 17 00:00:00 2001 From: Rod Boev Date: Fri, 4 Sep 2026 18:24:18 -0400 Subject: [PATCH 3/4] Separate current-schema derivation by driver behavior --- internal/store/upgrade.go | 41 +--------------------------------- internal/store/upgrade_test.go | 32 ++++++++++++++++++++++++++ 2 files changed, 33 insertions(+), 40 deletions(-) diff --git a/internal/store/upgrade.go b/internal/store/upgrade.go index 817430dc..0bee83b8 100644 --- a/internal/store/upgrade.go +++ b/internal/store/upgrade.go @@ -10,7 +10,6 @@ import ( "path/filepath" "slices" "strings" - "sync" "go.kenn.io/kit/pack" @@ -71,17 +70,6 @@ var ( removeInvalidUpgradeStage = removeUpgradeFileSet ) -type currentSchemaColumnsCacheKey struct { - driverName string - schemaSQL string -} - -var currentSchemaColumnsCache = struct { - sync.Mutex - - values map[currentSchemaColumnsCacheKey]map[string][]string -}{values: make(map[currentSchemaColumnsCacheKey]map[string][]string)} - var currentSchemaTables = [...]string{ "blobs", "blob_packs", "vault_metadata", "blob_stores", "blob_locations", "blob_pack_entries", } @@ -261,26 +249,7 @@ func validateCurrentSchemaColumns( } func canonicalCurrentSchemaColumns(driver docsqlite.Driver) (map[string][]string, error) { - key := currentSchemaColumnsCacheKey{driverName: driver.Name(), schemaSQL: schemaSQL} - currentSchemaColumnsCache.Lock() - columns, ok := currentSchemaColumnsCache.values[key] - currentSchemaColumnsCache.Unlock() - if ok { - return cloneSchemaColumns(columns), nil - } - - columns, err := deriveCurrentSchemaColumns(driver) - if err != nil { - return nil, err - } - currentSchemaColumnsCache.Lock() - if cached, ok := currentSchemaColumnsCache.values[key]; ok { - currentSchemaColumnsCache.Unlock() - return cloneSchemaColumns(cached), nil - } - currentSchemaColumnsCache.values[key] = cloneSchemaColumns(columns) - currentSchemaColumnsCache.Unlock() - return cloneSchemaColumns(columns), nil + return deriveCurrentSchemaColumns(driver) } func deriveCurrentSchemaColumns(driver docsqlite.Driver) (columns map[string][]string, err error) { @@ -321,14 +290,6 @@ func deriveCurrentSchemaColumns(driver docsqlite.Driver) (columns map[string][]s return columns, nil } -func cloneSchemaColumns(columns map[string][]string) map[string][]string { - clone := make(map[string][]string, len(columns)) - for table, names := range columns { - clone[table] = slices.Clone(names) - } - return clone -} - func validateV2Schema(db *sql.DB, blobs, packs []string) error { v2BlobColumns := []string{ metadataCreatedAtField, "hash", "loose_encoding", "loose_stored_size", diff --git a/internal/store/upgrade_test.go b/internal/store/upgrade_test.go index 039c9f93..39959372 100644 --- a/internal/store/upgrade_test.go +++ b/internal/store/upgrade_test.go @@ -243,6 +243,21 @@ func TestCanonicalCurrentSchemaDerivationDoesNotCacheErrors(t *testing.T) { assert.Contains(t, columns["blobs"], "synthetic_derivation_retry_269") } +func TestCanonicalCurrentSchemaDerivationSeparatesSameNamedDrivers(t *testing.T) { + const extraColumn = "synthetic_driver_variant_269" + base := DefaultSQLiteDriver() + plain := &schemaVariantDriver{Driver: base} + variant := &schemaVariantDriver{Driver: base, extraColumn: extraColumn} + + columns, err := canonicalCurrentSchemaColumns(plain) + require.NoError(t, err) + assert.NotContains(t, columns["blobs"], extraColumn) + + columns, err = canonicalCurrentSchemaColumns(variant) + require.NoError(t, err) + assert.Contains(t, columns["blobs"], extraColumn) +} + type flakySchemaDriver struct { docsqlite.Driver @@ -257,6 +272,23 @@ func (d *flakySchemaDriver) Open(path string, opts docsqlite.OpenOptions) (*sql. return d.Driver.Open(path, opts) } +type schemaVariantDriver struct { + docsqlite.Driver + extraColumn string +} + +func (d *schemaVariantDriver) Open(path string, opts docsqlite.OpenOptions) (*sql.DB, error) { + db, err := d.Driver.Open(path, opts) + if err != nil || d.extraColumn == "" || opts.Access != docsqlite.Create { + return db, err + } + if _, err := db.Exec(schemaSQL + "\nALTER TABLE blobs ADD COLUMN " + d.extraColumn + " TEXT"); err != nil { + _ = db.Close() + return nil, err + } + return db, nil +} + func schemaSQLWithAddedColumn(t *testing.T, original, table, column string) string { t.Helper() lineEnding := "\n" From 86f3888fd8cf194157a1c6afd60c19029bc8a96d Mon Sep 17 00:00:00 2001 From: Rod Boev Date: Fri, 4 Sep 2026 18:26:27 -0400 Subject: [PATCH 4/4] Format same-named driver regression test --- internal/store/upgrade_test.go | 1 + 1 file changed, 1 insertion(+) diff --git a/internal/store/upgrade_test.go b/internal/store/upgrade_test.go index 39959372..f845ae7f 100644 --- a/internal/store/upgrade_test.go +++ b/internal/store/upgrade_test.go @@ -274,6 +274,7 @@ func (d *flakySchemaDriver) Open(path string, opts docsqlite.OpenOptions) (*sql. type schemaVariantDriver struct { docsqlite.Driver + extraColumn string }