From b2eddfa8b172e8152dcd40342f6d991ec0c9e61e Mon Sep 17 00:00:00 2001 From: des Date: Sat, 2 May 2026 13:36:13 +0200 Subject: [PATCH] Enable safe concurrent SQLite reads and writes Reconfigure the SQLite backend so multiple goroutines can read and write without hitting SQLITE_BUSY or busy-snapshot errors: - Build a URI DSN that applies journal_mode=WAL, busy_timeout=5000, foreign_keys=1, and synchronous=NORMAL on every pooled connection (the previous one-shot db.Exec only configured one pool member). - Force _txlock=immediate so every BeginTx acquires the writer lock up front, eliminating busy-snapshot in the read-then-write transactions used for collision resolution and file ops. - Bound the connection pool from CPU count; pin :memory: paths to a single connection since each new conn opens a separate in-memory DB. Add concurrent test coverage for the repository (per-conn pragma check, distinct concurrent writes, mixed reader/writer load, concurrent read-then-write transactions, expiration sweeper racing with writes) and the service layer (concurrent collision resolution producing unique resolved words, distinct concurrent writes via StoreEntry). Run go test with -race in CI and add a dedicated concurrency stress job that runs the concurrent suites with -count=5. Document the SQLite concurrency configuration and the steps to switch an existing deployment from PostgreSQL to SQLite. --- .github/workflows/ci.yml | 16 +- README.md | 39 +++ internal/repository/sqlite.go | 72 ++++- internal/repository/sqlite_concurrent_test.go | 304 ++++++++++++++++++ .../service/entry_service_concurrent_test.go | 114 +++++++ 5 files changed, 534 insertions(+), 11 deletions(-) create mode 100644 internal/repository/sqlite_concurrent_test.go create mode 100644 internal/service/entry_service_concurrent_test.go diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index bba2399..86ad4a0 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -39,8 +39,20 @@ jobs: - uses: actions/setup-go@v5 with: go-version-file: go.mod - - name: Test - run: go test ./... + - name: Test (race detector) + run: go test -race -timeout 5m ./... + + concurrency: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + - uses: actions/setup-go@v5 + with: + go-version-file: go.mod + - name: SQLite concurrency stress (repository) + run: go test -race -count=5 -run='Concurrent|PragmasOnEveryPoolConnection' -timeout 5m ./internal/repository/... + - name: SQLite concurrency stress (service) + run: go test -race -count=5 -run='Concurrent' -timeout 5m ./internal/service/... build: runs-on: ubuntu-latest diff --git a/README.md b/README.md index 040438e..4368222 100644 --- a/README.md +++ b/README.md @@ -147,6 +147,45 @@ All settings live in `config.yaml`. Every value can be overridden with environme | `database.sqlite.path` | `WORDSTORE_DATABASE_SQLITE_PATH` | `./data/word-store.db` | SQLite file path | | `database.postgres.dsn` | `WORDSTORE_DATABASE_POSTGRES_DSN` | | PostgreSQL connection string | +#### SQLite concurrency + +The SQLite backend is configured for safe concurrent reads and writes: + +- **WAL journal mode** — multiple readers can run at the same time as one writer. +- **`busy_timeout=5000`** — writers wait up to 5 s on lock contention instead of failing immediately. +- **`BEGIN IMMEDIATE` for every transaction** — prevents busy-snapshot in read-then-write flows (collision resolution, file ops). +- **`synchronous=NORMAL` + `foreign_keys=1`** applied to every pooled connection. +- **Bounded connection pool** sized from CPU count. + +SQLite still serializes writers globally — that is a SQLite invariant — but readers run in parallel and write contention is absorbed by the busy timeout. For most MCP workloads this is more than sufficient; reach for PostgreSQL only if you need cross-process writers or a centralized DB. + +#### Switching from PostgreSQL to SQLite + +Driver selection is a config flip; there is no automatic data migration between backends. + +**1. Update config** — either edit `config.yaml`: + +```yaml +database: + driver: "sqlite" + sqlite: + path: "./data/word-store.db" +``` + +…or override via env var (takes precedence over `config.yaml`): + +```bash +export WORDSTORE_DATABASE_DRIVER=sqlite +export WORDSTORE_DATABASE_SQLITE_PATH=./data/word-store.db +./watchword +``` + +The directory in `path` is created on startup, and migrations run on first boot. + +**2. Docker** — the default `docker-compose.yml` launches PostgreSQL alongside watchword. To run on SQLite, stop that compose stack (`docker compose down`) and either run the binary directly or use a compose override that drops the `postgres` service, sets `WORDSTORE_DATABASE_DRIVER=sqlite` plus `WORDSTORE_DATABASE_SQLITE_PATH=/data/word-store.db`, and mounts a named volume at `/data` so the DB file survives container restarts. + +**3. Migrating data (optional)** — switching driver starts from an empty database. If you need to carry entries across, dump the `entries` table from PostgreSQL (`COPY entries TO STDOUT (FORMAT csv, HEADER)`) and load it into SQLite with `.import`; the schemas are equivalent, but PostgreSQL `timestamptz` columns must be converted to RFC3339 strings for SQLite during the dump. + ### Authentication | Setting | Env var | Default | Description | diff --git a/internal/repository/sqlite.go b/internal/repository/sqlite.go index 879c5c3..2413569 100644 --- a/internal/repository/sqlite.go +++ b/internal/repository/sqlite.go @@ -4,6 +4,8 @@ import ( "context" "database/sql" "fmt" + "runtime" + "strings" "time" "github.com/google/uuid" @@ -17,19 +19,71 @@ type SQLiteRepo struct { db *sql.DB } +// buildSQLiteDSN turns a database path into a modernc.org/sqlite URI that +// applies the pragmas required for safe concurrent access on every pooled +// connection, and forces BEGIN IMMEDIATE for every transaction so the +// read-then-write flows in the service layer cannot hit busy-snapshot. +// +// - journal_mode=WAL multi-reader / single-writer concurrency +// - busy_timeout=5000 wait up to 5s on lock contention instead of +// failing immediately with SQLITE_BUSY +// - foreign_keys=1 enforce FKs on every connection (the prior +// single Exec only configured one pool member) +// - synchronous=NORMAL safe under WAL, much faster than FULL +// - _txlock=immediate every BeginTx issues BEGIN IMMEDIATE so the +// writer lock is acquired up front +func buildSQLiteDSN(path string) string { + const pragmas = "_pragma=journal_mode(WAL)" + + "&_pragma=busy_timeout(5000)" + + "&_pragma=foreign_keys(1)" + + "&_pragma=synchronous(NORMAL)" + + "&_txlock=immediate" + + switch { + case path == ":memory:": + return "file::memory:?" + pragmas + case strings.HasPrefix(path, "file:"): + sep := "?" + if strings.Contains(path, "?") { + sep = "&" + } + return path + sep + pragmas + default: + return "file:" + path + "?" + pragmas + } +} + +// sqlitePoolSize picks a reasonable bounded pool. SQLite serializes writes +// globally, so a deep pool buys nothing for writes — but WAL allows truly +// concurrent reads, so a handful of connections lets readers run in parallel. +func sqlitePoolSize(path string) int { + if path == ":memory:" { + // Each new connection to ":memory:" opens a separate in-memory + // database. Pin to a single shared connection so the schema and + // data stay consistent across queries. + return 1 + } + n := runtime.NumCPU() * 2 + if n < 4 { + n = 4 + } + if n > 16 { + n = 16 + } + return n +} + func NewSQLiteRepo(dbPath string) (*SQLiteRepo, error) { - db, err := sql.Open("sqlite", dbPath) + db, err := sql.Open("sqlite", buildSQLiteDSN(dbPath)) if err != nil { return nil, fmt.Errorf("opening sqlite: %w", err) } - if _, err := db.Exec("PRAGMA journal_mode=WAL"); err != nil { - db.Close() - return nil, fmt.Errorf("setting WAL mode: %w", err) - } - if _, err := db.Exec("PRAGMA foreign_keys=ON"); err != nil { - db.Close() - return nil, fmt.Errorf("enabling foreign keys: %w", err) - } + + pool := sqlitePoolSize(dbPath) + db.SetMaxOpenConns(pool) + db.SetMaxIdleConns(pool) + db.SetConnMaxIdleTime(5 * time.Minute) + return &SQLiteRepo{db: db}, nil } diff --git a/internal/repository/sqlite_concurrent_test.go b/internal/repository/sqlite_concurrent_test.go new file mode 100644 index 0000000..516fd77 --- /dev/null +++ b/internal/repository/sqlite_concurrent_test.go @@ -0,0 +1,304 @@ +package repository + +import ( + "context" + "database/sql" + "fmt" + "path/filepath" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/watchword/watchword/internal/domain" +) + +func newDiskRepo(t *testing.T) *SQLiteRepo { + t.Helper() + dbPath := filepath.Join(t.TempDir(), "test.db") + repo, err := NewSQLiteRepo(dbPath) + if err != nil { + t.Fatalf("NewSQLiteRepo: %v", err) + } + t.Cleanup(func() { repo.Close() }) + if err := repo.Migrate(context.Background()); err != nil { + t.Fatalf("Migrate: %v", err) + } + return repo +} + +func TestBuildSQLiteDSN(t *testing.T) { + const wantPragmas = "_pragma=journal_mode(WAL)" + + "&_pragma=busy_timeout(5000)" + + "&_pragma=foreign_keys(1)" + + "&_pragma=synchronous(NORMAL)" + + "&_txlock=immediate" + + cases := []struct { + name string + in string + want string + }{ + {"memory", ":memory:", "file::memory:?" + wantPragmas}, + {"absolute path", "/tmp/test.db", "file:/tmp/test.db?" + wantPragmas}, + {"relative path", "./data/test.db", "file:./data/test.db?" + wantPragmas}, + {"uri without query", "file:/tmp/test.db", "file:/tmp/test.db?" + wantPragmas}, + {"uri with query", "file:/tmp/test.db?cache=shared", "file:/tmp/test.db?cache=shared&" + wantPragmas}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + if got := buildSQLiteDSN(tc.in); got != tc.want { + t.Errorf("buildSQLiteDSN(%q)\n got: %q\n want: %q", tc.in, got, tc.want) + } + }) + } +} + +// TestSQLite_PragmasOnEveryPoolConnection verifies that the pragmas are +// applied per-connection rather than once at startup. Every connection in the +// pool must report WAL, busy_timeout=5000, and foreign_keys=1; otherwise +// concurrent writes can hit instant SQLITE_BUSY and FK constraints become +// silently inconsistent. +func TestSQLite_PragmasOnEveryPoolConnection(t *testing.T) { + repo := newDiskRepo(t) + ctx := context.Background() + + // Hold several connections open simultaneously so the pool is forced + // to allocate distinct underlying connections. + const n = 4 + conns := make([]*sql.Conn, n) + for i := 0; i < n; i++ { + c, err := repo.db.Conn(ctx) + if err != nil { + t.Fatalf("Conn(%d): %v", i, err) + } + conns[i] = c + defer c.Close() + } + + for i, c := range conns { + var journal string + if err := c.QueryRowContext(ctx, "PRAGMA journal_mode").Scan(&journal); err != nil { + t.Fatalf("conn %d journal_mode: %v", i, err) + } + if !strings.EqualFold(journal, "wal") { + t.Errorf("conn %d journal_mode=%q, want wal", i, journal) + } + + var busy int + if err := c.QueryRowContext(ctx, "PRAGMA busy_timeout").Scan(&busy); err != nil { + t.Fatalf("conn %d busy_timeout: %v", i, err) + } + if busy != 5000 { + t.Errorf("conn %d busy_timeout=%d, want 5000", i, busy) + } + + var fk int + if err := c.QueryRowContext(ctx, "PRAGMA foreign_keys").Scan(&fk); err != nil { + t.Fatalf("conn %d foreign_keys: %v", i, err) + } + if fk != 1 { + t.Errorf("conn %d foreign_keys=%d, want 1", i, fk) + } + } +} + +// TestSQLite_ConcurrentDistinctWrites stresses the writer-serialization path: +// many goroutines insert distinct words at the same time. With busy_timeout +// and BEGIN IMMEDIATE in place, every write must succeed even though SQLite +// only admits one writer at a time. +func TestSQLite_ConcurrentDistinctWrites(t *testing.T) { + repo := newDiskRepo(t) + ctx := context.Background() + + const goroutines = 32 + const perGoroutine = 25 + + var wg sync.WaitGroup + var failures atomic.Int32 + for g := 0; g < goroutines; g++ { + wg.Add(1) + go func(g int) { + defer wg.Done() + for i := 0; i < perGoroutine; i++ { + word := fmt.Sprintf("w_%d_%d", g, i) + if _, err := repo.Store(ctx, &domain.Entry{Word: word, Payload: "x"}); err != nil { + failures.Add(1) + t.Errorf("Store(%s): %v", word, err) + return + } + } + }(g) + } + wg.Wait() + if failures.Load() != 0 { + t.Fatalf("had %d store failures under concurrency", failures.Load()) + } + + _, total, err := repo.List(ctx, "active", 1, 0, "word", "asc") + if err != nil { + t.Fatalf("List: %v", err) + } + if want := goroutines * perGoroutine; total != want { + t.Errorf("expected %d entries, got %d", want, total) + } +} + +// TestSQLite_ConcurrentReadsAndWrites verifies that readers and writers can +// run concurrently — readers should never be blocked by an in-flight writer, +// and writers should not produce errors while readers iterate. +func TestSQLite_ConcurrentReadsAndWrites(t *testing.T) { + repo := newDiskRepo(t) + ctx := context.Background() + + for i := 0; i < 50; i++ { + if _, err := repo.Store(ctx, &domain.Entry{Word: fmt.Sprintf("seed%d", i), Payload: "p"}); err != nil { + t.Fatalf("seed %d: %v", i, err) + } + } + + deadline := time.Now().Add(2 * time.Second) + var wg sync.WaitGroup + var readErrs, writeErrs atomic.Int32 + + for i := 0; i < 8; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for time.Now().Before(deadline) { + if _, _, err := repo.List(ctx, "active", 50, 0, "word", "asc"); err != nil { + readErrs.Add(1) + t.Errorf("List: %v", err) + return + } + if _, err := repo.GetByWord(ctx, "seed10", false); err != nil { + readErrs.Add(1) + t.Errorf("GetByWord: %v", err) + return + } + } + }() + } + + var counter atomic.Int64 + for i := 0; i < 4; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for time.Now().Before(deadline) { + n := counter.Add(1) + w := fmt.Sprintf("conc%d", n) + if _, err := repo.Store(ctx, &domain.Entry{Word: w, Payload: "p"}); err != nil { + writeErrs.Add(1) + t.Errorf("Store(%s): %v", w, err) + return + } + } + }() + } + + wg.Wait() + if readErrs.Load() != 0 || writeErrs.Load() != 0 { + t.Fatalf("readErrs=%d writeErrs=%d", readErrs.Load(), writeErrs.Load()) + } + if counter.Load() == 0 { + t.Fatal("expected at least one concurrent write") + } +} + +// TestSQLite_ConcurrentTransactions runs many BEGIN IMMEDIATE transactions +// that read-then-write. Without _txlock=immediate this is the classic +// busy-snapshot trap; with it, every transaction must complete cleanly. +func TestSQLite_ConcurrentTransactions(t *testing.T) { + repo := newDiskRepo(t) + ctx := context.Background() + + const goroutines = 16 + var wg sync.WaitGroup + var failures atomic.Int32 + + for g := 0; g < goroutines; g++ { + wg.Add(1) + go func(g int) { + defer wg.Done() + err := repo.WithTx(ctx, func(txRepo Repository) error { + w := fmt.Sprintf("tx%d", g) + exists, err := txRepo.WordExistsActive(ctx, w) + if err != nil { + return err + } + if exists { + return fmt.Errorf("word already exists in fresh DB: %s", w) + } + _, err = txRepo.Store(ctx, &domain.Entry{Word: w, Payload: "p"}) + return err + }) + if err != nil { + failures.Add(1) + t.Errorf("tx %d: %v", g, err) + } + }(g) + } + wg.Wait() + if failures.Load() != 0 { + t.Fatalf("transaction failures: %d", failures.Load()) + } + + _, total, _ := repo.List(ctx, "active", 1, 0, "word", "asc") + if total != goroutines { + t.Errorf("expected %d, got %d", goroutines, total) + } +} + +// TestSQLite_ConcurrentExpirationAndWrites is closer to the real workload: +// the background expiration sweeper races with online writes. Both paths use +// UPDATE/INSERT and must not deadlock or fail. +func TestSQLite_ConcurrentExpirationAndWrites(t *testing.T) { + repo := newDiskRepo(t) + ctx := context.Background() + + past := time.Now().Add(-1 * time.Hour) + for i := 0; i < 20; i++ { + if _, err := repo.Store(ctx, &domain.Entry{Word: fmt.Sprintf("exp%d", i), Payload: "p", ExpiresAt: &past}); err != nil { + t.Fatalf("seed %d: %v", i, err) + } + } + + deadline := time.Now().Add(time.Second) + var wg sync.WaitGroup + var failures atomic.Int32 + + wg.Add(1) + go func() { + defer wg.Done() + for time.Now().Before(deadline) { + if _, err := repo.MarkExpiredBatch(ctx, 10); err != nil { + failures.Add(1) + t.Errorf("MarkExpiredBatch: %v", err) + return + } + } + }() + + var counter atomic.Int64 + for i := 0; i < 4; i++ { + wg.Add(1) + go func() { + defer wg.Done() + for time.Now().Before(deadline) { + n := counter.Add(1) + w := fmt.Sprintf("live%d", n) + if _, err := repo.Store(ctx, &domain.Entry{Word: w, Payload: "p"}); err != nil { + failures.Add(1) + t.Errorf("Store(%s): %v", w, err) + return + } + } + }() + } + wg.Wait() + if failures.Load() != 0 { + t.Fatalf("failures: %d", failures.Load()) + } +} diff --git a/internal/service/entry_service_concurrent_test.go b/internal/service/entry_service_concurrent_test.go new file mode 100644 index 0000000..f6cb787 --- /dev/null +++ b/internal/service/entry_service_concurrent_test.go @@ -0,0 +1,114 @@ +package service + +import ( + "context" + "fmt" + "log/slog" + "os" + "path/filepath" + "sync" + "sync/atomic" + "testing" + + "github.com/watchword/watchword/internal/repository" +) + +func newDiskService(t *testing.T) *EntryService { + t.Helper() + dbPath := filepath.Join(t.TempDir(), "test.db") + repo, err := repository.NewSQLiteRepo(dbPath) + if err != nil { + t.Fatalf("NewSQLiteRepo: %v", err) + } + t.Cleanup(func() { repo.Close() }) + if err := repo.Migrate(context.Background()); err != nil { + t.Fatalf("Migrate: %v", err) + } + logger := slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})) + return NewEntryService(repo, 168, logger) +} + +// TestStoreEntry_ConcurrentCollisionResolution drives many goroutines at the +// same base word simultaneously. Because the collision-resolution flow is a +// read-then-write transaction (WordExistsActive -> Store), concurrent calls +// must serialize cleanly via BEGIN IMMEDIATE + busy_timeout. Each call must +// succeed and produce a distinct resolved word. +func TestStoreEntry_ConcurrentCollisionResolution(t *testing.T) { + svc := newDiskService(t) + ctx := context.Background() + + const goroutines = 24 + var wg sync.WaitGroup + var failures atomic.Int32 + + results := make([]string, goroutines) + for i := 0; i < goroutines; i++ { + wg.Add(1) + go func(i int) { + defer wg.Done() + res, err := svc.StoreEntry(ctx, "rabbit", "payload", nil) + if err != nil { + failures.Add(1) + t.Errorf("StoreEntry %d: %v", i, err) + return + } + results[i] = res.Entry.Word + }(i) + } + wg.Wait() + if failures.Load() != 0 { + t.Fatalf("had %d failures", failures.Load()) + } + + seen := make(map[string]bool, goroutines) + hasBase := false + for _, w := range results { + if w == "" { + t.Fatal("empty result word") + } + if seen[w] { + t.Errorf("duplicate resolved word: %s", w) + } + seen[w] = true + if w == "rabbit" { + hasBase = true + } + } + if !hasBase { + t.Error("expected exactly one goroutine to win the base word 'rabbit'") + } + if len(seen) != goroutines { + t.Errorf("expected %d unique words, got %d", goroutines, len(seen)) + } +} + +// TestStoreEntry_ConcurrentDistinctWords confirms that unrelated writers +// don't interfere with each other under load. +func TestStoreEntry_ConcurrentDistinctWords(t *testing.T) { + svc := newDiskService(t) + ctx := context.Background() + + const goroutines = 32 + const perGoroutine = 10 + var wg sync.WaitGroup + var failures atomic.Int32 + + for g := 0; g < goroutines; g++ { + wg.Add(1) + go func(g int) { + defer wg.Done() + for i := 0; i < perGoroutine; i++ { + word := fmt.Sprintf("word_%d_%d", g, i) + if _, err := svc.StoreEntry(ctx, word, "p", nil); err != nil { + failures.Add(1) + t.Errorf("StoreEntry(%s): %v", word, err) + return + } + } + }(g) + } + wg.Wait() + if failures.Load() != 0 { + t.Fatalf("had %d failures", failures.Load()) + } +}