diff --git a/appview/db/db.go b/appview/db/db.go index b344e09f8..56c59e472 100644 --- a/appview/db/db.go +++ b/appview/db/db.go @@ -5,11 +5,10 @@ import ( "database/sql" "fmt" "log/slog" - "strings" - _ "github.com/mattn/go-sqlite3" "tangled.org/core/log" "tangled.org/core/orm" + "tangled.org/core/sqlite" ) type DB struct { @@ -29,19 +28,10 @@ type Execer interface { } func Make(ctx context.Context, dbPath string) (*DB, error) { - // https://github.com/mattn/go-sqlite3#connection-string - opts := []string{ - "_foreign_keys=1", - "_journal_mode=WAL", - "_synchronous=NORMAL", - "_auto_vacuum=incremental", - "_busy_timeout=5000", - } - logger := log.FromContext(ctx) logger = log.SubLogger(logger, "db") - db, err := sql.Open("sqlite3", dbPath+"?"+strings.Join(opts, "&")) + db, err := sqlite.Open(dbPath) if err != nil { return nil, err } diff --git a/cmd/deliberi-migrate/main.go b/cmd/deliberi-migrate/main.go index e214d2514..33b051e60 100644 --- a/cmd/deliberi-migrate/main.go +++ b/cmd/deliberi-migrate/main.go @@ -7,7 +7,7 @@ import ( "fmt" "os" - _ "github.com/mattn/go-sqlite3" + "tangled.org/core/sqlite" ) func main() { @@ -24,13 +24,13 @@ func main() { } srcPath, dstPath := flag.Arg(0), flag.Arg(1) - src, err := sql.Open("sqlite3", srcPath+"?_journal_mode=WAL") + src, err := sqlite.Open(srcPath) if err != nil { fatal("open source: %v", err) } defer src.Close() - dst, err := sql.Open("sqlite3", dstPath+"?_foreign_keys=1&_journal_mode=WAL") + dst, err := sqlite.Open(dstPath) if err != nil { fatal("open dest: %v", err) } diff --git a/cmd/populatepipelines/populate_pipelines.go b/cmd/populatepipelines/populate_pipelines.go index 87466bc72..60c2d4a0f 100644 --- a/cmd/populatepipelines/populate_pipelines.go +++ b/cmd/populatepipelines/populate_pipelines.go @@ -9,7 +9,7 @@ import ( "time" "github.com/bluesky-social/indigo/atproto/syntax" - _ "github.com/mattn/go-sqlite3" + "tangled.org/core/sqlite" ) var ( @@ -68,7 +68,7 @@ func main() { log.Fatalf("Invalid repo format: %s (expected: did:plc:xyz/reponame)", *repo) } - db, err := sql.Open("sqlite3", *dbPath) + db, err := sqlite.Open(*dbPath) if err != nil { log.Fatalf("Failed to open database: %v", err) } diff --git a/deliberi/db/db.go b/deliberi/db/db.go index 12806332b..b0a7699db 100644 --- a/deliberi/db/db.go +++ b/deliberi/db/db.go @@ -6,9 +6,9 @@ import ( "log/slog" "strings" - _ "github.com/mattn/go-sqlite3" "tangled.org/core/log" "tangled.org/core/orm" + "tangled.org/core/sqlite" ) type DB struct { @@ -24,16 +24,9 @@ type Execer interface { } func Make(ctx context.Context, dbPath string) (*DB, error) { - opts := []string{ - "_foreign_keys=1", - "_journal_mode=WAL", - "_synchronous=NORMAL", - "_busy_timeout=5000", - } - logger := log.SubLogger(log.FromContext(ctx), "db") - db, err := sql.Open("sqlite3", dbPath+"?"+strings.Join(opts, "&")) + db, err := sqlite.Open(dbPath) if err != nil { return nil, err } diff --git a/eventconsumer/cursor/sqlite.go b/eventconsumer/cursor/sqlite.go index 4d1b74288..578134072 100644 --- a/eventconsumer/cursor/sqlite.go +++ b/eventconsumer/cursor/sqlite.go @@ -6,7 +6,7 @@ import ( "fmt" "log/slog" - _ "github.com/mattn/go-sqlite3" + "tangled.org/core/sqlite" ) type SqliteStore struct { @@ -23,7 +23,7 @@ func WithTableName(name string) SqliteStoreOpt { } func NewSQLiteStore(dbPath string, opts ...SqliteStoreOpt) (*SqliteStore, error) { - db, err := sql.Open("sqlite3", dbPath+"?_foreign_keys=1") + db, err := sqlite.Open(dbPath) if err != nil { return nil, fmt.Errorf("failed to open sqlite database: %w", err) } diff --git a/knotserver/db/db.go b/knotserver/db/db.go index ab88b45a6..e5e208ff4 100644 --- a/knotserver/db/db.go +++ b/knotserver/db/db.go @@ -6,13 +6,12 @@ import ( "fmt" "log/slog" "os" - "strings" "sync" securejoin "github.com/cyphar/filepath-securejoin" - _ "github.com/mattn/go-sqlite3" "tangled.org/core/log" "tangled.org/core/orm" + "tangled.org/core/sqlite" ) type DB struct { @@ -47,19 +46,10 @@ func (d *DB) Query(query string, args ...any) (*sql.Rows, error) { } func Setup(ctx context.Context, dbPath string) (*DB, error) { - // https://github.com/mattn/go-sqlite3#connection-string - opts := []string{ - "_foreign_keys=1", - "_journal_mode=WAL", - "_synchronous=NORMAL", - "_auto_vacuum=incremental", - "_busy_timeout=5000", - } - logger := log.FromContext(ctx) logger = log.SubLogger(logger, "db") - db, err := sql.Open("sqlite3", dbPath+"?"+strings.Join(opts, "&")) + db, err := sqlite.Open(dbPath) if err != nil { return nil, err } diff --git a/rbac/rbac.go b/rbac/rbac.go index 5908b0ca7..5b3ef8b6a 100644 --- a/rbac/rbac.go +++ b/rbac/rbac.go @@ -1,13 +1,13 @@ package rbac import ( - "database/sql" "slices" "strings" adapter "github.com/Blank-Xu/sql-adapter" "github.com/casbin/casbin/v2" "github.com/casbin/casbin/v2/model" + "tangled.org/core/sqlite" ) const ( @@ -43,7 +43,7 @@ func NewEnforcer(path string) (*Enforcer, error) { return nil, err } - db, err := sql.Open("sqlite3", path+"?_foreign_keys=1&_journal_mode=WAL&_busy_timeout=5000") + db, err := sqlite.Open(path) if err != nil { return nil, err } diff --git a/spindle/db/db.go b/spindle/db/db.go index 96ad17e8a..d4e91e8d0 100644 --- a/spindle/db/db.go +++ b/spindle/db/db.go @@ -5,11 +5,10 @@ import ( "database/sql" "log/slog" "slices" - "strings" - _ "github.com/mattn/go-sqlite3" "tangled.org/core/log" "tangled.org/core/orm" + "tangled.org/core/sqlite" ) type DB struct { @@ -22,19 +21,10 @@ type DBTX interface { } func Make(ctx context.Context, dbPath string) (*DB, error) { - // https://github.com/mattn/go-sqlite3#connection-string - opts := []string{ - "_foreign_keys=1", - "_journal_mode=WAL", - "_synchronous=NORMAL", - "_auto_vacuum=incremental", - "_busy_timeout=5000", - } - logger := log.FromContext(ctx) logger = log.SubLogger(logger, "db") - db, err := sql.Open("sqlite3", dbPath+"?"+strings.Join(opts, "&")) + db, err := sqlite.Open(dbPath) if err != nil { return nil, err } diff --git a/spindle/db/pipelines.go b/spindle/db/pipelines.go index 028f13f8c..1f2854467 100644 --- a/spindle/db/pipelines.go +++ b/spindle/db/pipelines.go @@ -71,25 +71,45 @@ func (d *DB) QueryPipelines(ctx context.Context, repoDid string, commits []strin if err != nil { return nil, "", 0, err } - defer rows.Close() - - var pipelines []*tangled.CiPipeline - var lastCreated int64 + // mapToCiPipeline issues nested queries, so rows must be fully drained and + // released before mapping; the pool is a single connection and would + // otherwise deadlock + var rawRows []struct { + rkey string + eventJson string + created int64 + } for rows.Next() { - var rkey, eventJson string - var created int64 - if err := rows.Scan(&rkey, &eventJson, &created); err != nil { + var row struct { + rkey string + eventJson string + created int64 + } + if err := rows.Scan(&row.rkey, &row.eventJson, &row.created); err != nil { + rows.Close() return nil, "", 0, err } - lastCreated = created + rawRows = append(rawRows, row) + } + if err := rows.Err(); err != nil { + rows.Close() + return nil, "", 0, err + } + rows.Close() + + var pipelines []*tangled.CiPipeline + var lastCreated int64 + + for _, row := range rawRows { + lastCreated = row.created var rawPipeline tangled.Pipeline - if err := json.Unmarshal([]byte(eventJson), &rawPipeline); err != nil { + if err := json.Unmarshal([]byte(row.eventJson), &rawPipeline); err != nil { continue } - p, err := d.mapToCiPipeline(rkey, created, rawPipeline) + p, err := d.mapToCiPipeline(row.rkey, row.created, rawPipeline) if err != nil { return nil, "", 0, err } diff --git a/spindle/secrets/sqlite.go b/spindle/secrets/sqlite.go index ff240590a..7e7ecc758 100644 --- a/spindle/secrets/sqlite.go +++ b/spindle/secrets/sqlite.go @@ -7,7 +7,7 @@ import ( "fmt" "time" - _ "github.com/mattn/go-sqlite3" + "tangled.org/core/sqlite" ) type SqliteManager struct { @@ -24,7 +24,7 @@ func WithTableName(name string) SqliteManagerOpt { } func NewSQLiteManager(dbPath string, opts ...SqliteManagerOpt) (*SqliteManager, error) { - db, err := sql.Open("sqlite3", dbPath+"?_foreign_keys=1") + db, err := sqlite.Open(dbPath) if err != nil { return nil, fmt.Errorf("failed to open sqlite database: %w", err) } diff --git a/spindle/startup_migrations.go b/spindle/startup_migrations.go index 6d59a66c6..24f9b037a 100644 --- a/spindle/startup_migrations.go +++ b/spindle/startup_migrations.go @@ -2,14 +2,13 @@ package spindle import ( "context" - "database/sql" "errors" "fmt" "log/slog" "os" - _ "github.com/mattn/go-sqlite3" "tangled.org/core/spindle/db" + "tangled.org/core/sqlite" ) const forceTapResyncFlag = "force-tap-repo-resync-v1" @@ -81,7 +80,7 @@ func nudgeTapForResync(ctx context.Context, d *db.DB, tapDBPath string, logger * return fmt.Errorf("stat tap db: %w", err) } - tdb, err := sql.Open("sqlite3", tapDBPath+"?_busy_timeout=5000") + tdb, err := sqlite.Open(tapDBPath) if err != nil { return fmt.Errorf("open tap db: %w", err) } diff --git a/spindle/tap_drain.go b/spindle/tap_drain.go index 71992dfa5..291da9d1f 100644 --- a/spindle/tap_drain.go +++ b/spindle/tap_drain.go @@ -2,7 +2,6 @@ package spindle import ( "context" - "database/sql" "fmt" "net/url" "sync/atomic" @@ -12,7 +11,7 @@ import ( "github.com/bluesky-social/indigo/events" "github.com/bluesky-social/indigo/events/schedulers/sequential" "github.com/gorilla/websocket" - _ "github.com/mattn/go-sqlite3" + "tangled.org/core/sqlite" ) const ( @@ -30,7 +29,7 @@ func (s *Spindle) watchTapDrain(ctx context.Context, stop context.CancelFunc) { s.l.Info("tap drain watcher: relay head seq at startup", "head", headSeq, "relay", s.cfg.Server.Tap.RelayUrl) } - conn, err := sql.Open("sqlite3", s.cfg.Server.Tap.DBPath) + conn, err := sqlite.Open(s.cfg.Server.Tap.DBPath) if err != nil { s.l.Warn("tap drain watcher: opening tap db failed", "err", err) return diff --git a/sqlite/sqlite.go b/sqlite/sqlite.go new file mode 100644 index 000000000..0464dc7f8 --- /dev/null +++ b/sqlite/sqlite.go @@ -0,0 +1,22 @@ +package sqlite + +import ( + "database/sql" + "strings" + + _ "github.com/mattn/go-sqlite3" +) + +// https://github.com/mattn/go-sqlite3#connection-string +var dsnOpts = []string{ + "_foreign_keys=1", + "_journal_mode=WAL", + "_synchronous=NORMAL", + "_auto_vacuum=incremental", + "_busy_timeout=5000", + "_txlock=immediate", +} + +func Open(dbPath string) (*sql.DB, error) { + return sql.Open("sqlite3", dbPath+"?"+strings.Join(dsnOpts, "&")) +}