tests: make the carstore lock tests deterministic, and silence slog in test binaries

TestCarstoreConcurrentWritesTwoOpeners failed in CI on 2026-09-13 with
"database is locked" while the code was correct. It hammered the file with
five concurrent writers and asserted that no lock error surfaced within the
5 s busy_timeout, which on a loaded runner with every package testing in
parallel is a statement about the disk, not the code. The property it guards
(busy_timeout applied to every connection of every pool on a hold database,
c44a874) is now tested directly: one connection holds the write lock via
BEGIN IMMEDIATE, a second writer is shown to block rather than fail, and to
succeed once the lock is released. Both topologies are covered (the shared
OpenHoldDB pool, and a second opener on the same file, in both directions),
and a control shows a pool without busy_timeout fails immediately under the
same lock, so the passing tests are known to observe the mechanism.

The failure was also buried under the INFO lines every hold and PDS test
emits while booting. internal/testlog.Quiet swaps the default slog handler
for a discard handler unless the run is verbose or ATCR_TEST_LOGS is set,
and every package that produced that output now calls it from TestMain.
`go test` only shows a package's output when it fails, so this changes
nothing for passing runs and leaves a failing one readable.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Hho5da4daoCoPBJ9tCrL7s
This commit is contained in:
Evan Jarrett
2026-09-12 22:11:54 -05:00
co-authored by Claude Fable 5.1
parent 83092d9aee
commit 0ae9ca1a96
25 changed files with 500 additions and 85 deletions
+33
View File
@@ -0,0 +1,33 @@
// Package testlog silences slog in test binaries.
//
// The hold, PDS and appview packages log at INFO while they boot, subscribe,
// and tear down, and every test that starts one of them produces dozens of
// lines. `go test` only prints a package's output when it fails, so the noise
// lands exactly where the failure message is, and buries it.
//
// Call Quiet from a package's TestMain. Logs come back with `go test -v` or
// ATCR_TEST_LOGS=1, so a failing run can be re-run with them on. Note that
// `go test -json` implies -v, so tooling that streams JSON sees the logs too.
package testlog
import (
"flag"
"io"
"log/slog"
"os"
"testing"
)
// Quiet replaces the default slog logger with one that discards everything,
// unless the test binary runs verbose (-v) or ATCR_TEST_LOGS is set. It parses
// the test flags if that has not happened yet, which TestMain otherwise leaves
// to m.Run.
func Quiet() {
if !flag.Parsed() {
flag.Parse()
}
if testing.Verbose() || os.Getenv("ATCR_TEST_LOGS") != "" {
return
}
slog.SetDefault(slog.New(slog.NewTextHandler(io.Discard, nil)))
}
+15
View File
@@ -0,0 +1,15 @@
package authgate
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package db
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package handlers
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package holdclient
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package jetstream
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package leases
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package appview
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package middleware
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package storage
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package webhooks
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package did
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+4
View File
@@ -6,6 +6,8 @@ import (
"path/filepath" "path/filepath"
"testing" "testing"
"atcr.io/internal/testlog"
"atcr.io/pkg/hold/pds" "atcr.io/pkg/hold/pds"
) )
@@ -20,6 +22,8 @@ var (
// TestMain sets up shared test fixtures // TestMain sets up shared test fixtures
func TestMain(m *testing.M) { func TestMain(m *testing.M) {
testlog.Quiet() // see internal/testlog; -v or ATCR_TEST_LOGS=1 restores logs
// Create temp directory for shared keys // Create temp directory for shared keys
var err error var err error
sharedTempDir, err = os.MkdirTemp("", "holdlocal_test") sharedTempDir, err = os.MkdirTemp("", "holdlocal_test")
+15
View File
@@ -0,0 +1,15 @@
package auth
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package oauth
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package token
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package admin
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package db
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+155 -85
View File
@@ -2,20 +2,34 @@ package db
import ( import (
"context" "context"
"database/sql"
"fmt" "fmt"
"path/filepath" "path/filepath"
"strings" "strings"
"sync"
"testing" "testing"
"time"
"github.com/bluesky-social/indigo/models" "github.com/bluesky-social/indigo/models"
blockformat "github.com/ipfs/go-block-format" blockformat "github.com/ipfs/go-block-format"
"github.com/ipfs/go-cid" "github.com/ipfs/go-cid"
) )
// These tests pin the one property that keeps a hold's writers from failing
// each other: every connection in every pool on a hold database waits for the
// write lock (busy_timeout) instead of failing the statement the moment
// another writer holds it. Before c44a874 only the first connection of a pool
// had the timeout, and production writes failed with "database is locked".
//
// They test the mechanism, not the weather. An earlier version hammered the
// file with concurrent writers for a few seconds and asserted that no lock
// error surfaced; that verdict depended on how fast the CI disk was that day
// and failed on a loaded runner (2026-09-13) with the code correct. Here one
// connection deliberately holds the write lock, a second writer is shown to
// block rather than fail, and it is shown to succeed once the lock is
// released. A regression fails every time, on any machine.
// isLockErr reports whether err is a SQLite busy/locked error rather than a // isLockErr reports whether err is a SQLite busy/locked error rather than a
// logic error. Contention tests must fail on lock errors specifically, so a // logic error, so a typo cannot masquerade as a reproduction.
// typo cannot masquerade as a reproduction.
func isLockErr(err error) bool { func isLockErr(err error) bool {
if err == nil { if err == nil {
return false return false
@@ -41,45 +55,91 @@ func makeBlocks(t *testing.T, n int, salt string) (cid.Cid, map[cid.Cid]blockfor
return root, blks return root, blks
} }
// writeShards hammers the carstore with concurrent shard writes and reports any // holdWriteLock pins one connection from db, opens a write transaction on it
// lock error it hits. This is the exact path every hold record write takes. // and leaves it open, so every other writer on the file must wait. The
func writeShards(t *testing.T, sqs *SQLiteStore, writers, iters int) error { // returned release commits and hands the connection back.
func holdWriteLock(t *testing.T, db *sql.DB) (release func()) {
t.Helper() t.Helper()
ctx := context.Background() ctx := context.Background()
var wg sync.WaitGroup conn, err := db.Conn(ctx)
var mu sync.Mutex if err != nil {
var firstErr error t.Fatalf("pin connection: %v", err)
}
for w := range writers { if _, err := conn.ExecContext(ctx, "CREATE TABLE IF NOT EXISTS lock_probe (k INTEGER PRIMARY KEY, v TEXT)"); err != nil {
wg.Add(1) t.Fatalf("create lock_probe: %v", err)
go func(w int) { }
defer wg.Done() // BEGIN IMMEDIATE takes the write lock now rather than at the first write,
for i := range iters { // so the lock is held from the moment this returns.
salt := fmt.Sprintf("w%d-i%d", w, i) if _, err := conn.ExecContext(ctx, "BEGIN IMMEDIATE"); err != nil {
root, blks := makeBlocks(t, 12, salt) t.Fatalf("begin immediate: %v", err)
rev := fmt.Sprintf("rev-%03d-%03d", w, i) }
_, err := sqs.writeNewShard(ctx, root, rev, models.Uid(1), i, blks, nil) if _, err := conn.ExecContext(ctx, "INSERT INTO lock_probe (v) VALUES ('held')"); err != nil {
if err != nil { t.Fatalf("insert under lock: %v", err)
mu.Lock() }
if firstErr == nil { return func() {
firstErr = err if _, err := conn.ExecContext(ctx, "COMMIT"); err != nil {
} t.Errorf("commit lock holder: %v", err)
mu.Unlock() }
return _ = conn.Close()
}
}
}(w)
} }
wg.Wait()
return firstErr
} }
// TestCarstoreConcurrentWritesSharedPool reproduces the production topology: // holdDuration is how long the lock is held before release. The blocked
// one *sql.DB pool (OpenHoldDB) shared by the carstore, records index, events // writer must take at least most of this, which proves it waited rather than
// and scan broadcaster. database/sql opens a fresh connection per concurrent // racing past an already-released lock, and it must finish well inside
// caller, and busy_timeout is per-connection, so any connection beyond the // DefaultBusyTimeoutMs, which proves busy_timeout is what let it through.
// first commits with busy_timeout=0 and fails immediately under contention. const holdDuration = 300 * time.Millisecond
func TestCarstoreConcurrentWritesSharedPool(t *testing.T) {
// awaitBlockedWrite runs write while the lock is held, releases the lock after
// holdDuration, and checks that write blocked and then succeeded.
func awaitBlockedWrite(t *testing.T, what string, release func(), write func() error) {
t.Helper()
done := make(chan error, 1)
start := time.Now()
go func() { done <- write() }()
select {
case err := <-done:
// Finished while the lock was still held: either it failed (the
// regression) or it did not actually contend for the lock (a broken
// test), and both are reported.
if err != nil {
if isLockErr(err) {
t.Fatalf("%s failed instead of waiting for the write lock: %v", what, err)
}
t.Fatalf("%s failed: %v", what, err)
}
t.Fatalf("%s completed in %s while another connection held the write lock; the test is not contending", what, time.Since(start))
case <-time.After(holdDuration):
}
release()
select {
case err := <-done:
if err != nil {
t.Fatalf("%s failed after the lock was released: %v", what, err)
}
case <-time.After(time.Duration(DefaultBusyTimeoutMs) * time.Millisecond):
t.Fatalf("%s still blocked %dms after the lock was released", what, DefaultBusyTimeoutMs)
}
if waited := time.Since(start); waited < holdDuration {
t.Fatalf("%s took %s, less than the %s the lock was held", what, waited, holdDuration)
}
}
func carstoreWrite(sqs *SQLiteStore, root cid.Cid, blks map[cid.Cid]blockformat.Block) func() error {
return func() error {
_, err := sqs.writeNewShard(context.Background(), root, "rev-001", models.Uid(1), 1, blks, nil)
return err
}
}
// TestCarstoreWriteWaitsForLockSharedPool is the production topology: one
// *sql.DB pool (OpenHoldDB) shared by the carstore, records index, events and
// scan broadcaster. database/sql hands a different connection to each
// concurrent caller, so the carstore's write runs on a connection that is not
// the one holding the lock and must have its own busy_timeout.
func TestCarstoreWriteWaitsForLockSharedPool(t *testing.T) {
path := filepath.Join(t.TempDir(), "db.sqlite3") path := filepath.Join(t.TempDir(), "db.sqlite3")
h, err := OpenHoldDB(path, LibsqlConfig{}) h, err := OpenHoldDB(path, LibsqlConfig{})
if err != nil { if err != nil {
@@ -91,20 +151,17 @@ func TestCarstoreConcurrentWritesSharedPool(t *testing.T) {
if err != nil { if err != nil {
t.Fatalf("NewSQLiteStoreWithDB: %v", err) t.Fatalf("NewSQLiteStoreWithDB: %v", err)
} }
root, blks := makeBlocks(t, 12, "shared")
if err := writeShards(t, sqs, 8, 25); err != nil { release := holdWriteLock(t, h.DB)
if isLockErr(err) { awaitBlockedWrite(t, "shared-pool carstore write", release, carstoreWrite(sqs, root, blks))
t.Fatalf("shared-pool carstore write hit a lock error: %v", err)
}
t.Fatalf("shared-pool carstore write failed: %v", err)
}
} }
// TestCarstoreConcurrentWritesTwoOpeners covers the topology where a second // TestCarstoreWriteWaitsForLockTwoOpeners covers a second subsystem opening
// subsystem (the scan broadcaster or the event broadcaster, when they are not // <dir>/db.sqlite3 again through its own pool, in both directions: the
// handed the shared pool) opens <dir>/db.sqlite3 again through its own pool // carstore waiting on the other pool's lock, and the other pool waiting on the
// while the carstore is writing. // carstore's.
func TestCarstoreConcurrentWritesTwoOpeners(t *testing.T) { func TestCarstoreWriteWaitsForLockTwoOpeners(t *testing.T) {
dir := t.TempDir() dir := t.TempDir()
sqs, err := NewSqliteStore(dir) sqs, err := NewSqliteStore(dir)
if err != nil { if err != nil {
@@ -112,49 +169,62 @@ func TestCarstoreConcurrentWritesTwoOpeners(t *testing.T) {
} }
defer sqs.Close() defer sqs.Close()
// Second, independent pool on the same file, opened the way every hold other, err := OpenLocalDB(filepath.Join(dir, "db.sqlite3"))
// subsystem now opens one.
ri, err := OpenLocalDB(filepath.Join(dir, "db.sqlite3"))
if err != nil { if err != nil {
t.Fatalf("open records index: %v", err) t.Fatalf("open second pool: %v", err)
} }
defer ri.Close() defer other.Close()
if _, err := ri.Exec("CREATE TABLE IF NOT EXISTS records (collection TEXT, rkey TEXT, cid TEXT, PRIMARY KEY (collection, rkey))"); err != nil { if _, err := other.Exec("CREATE TABLE IF NOT EXISTS records (collection TEXT, rkey TEXT, cid TEXT, PRIMARY KEY (collection, rkey))"); err != nil {
t.Fatalf("create records table: %v", err) t.Fatalf("create records table: %v", err)
} }
ctx := context.Background() root, blks := makeBlocks(t, 12, "two-openers")
stop := make(chan struct{}) release := holdWriteLock(t, other)
idxErr := make(chan error, 1) awaitBlockedWrite(t, "carstore write behind the second pool's lock", release, carstoreWrite(sqs, root, blks))
go func() {
var err error
for i := 0; ; i++ {
select {
case <-stop:
idxErr <- err
return
default:
}
_, e := ri.ExecContext(ctx,
"INSERT INTO records (collection, rkey, cid) VALUES (?, ?, ?) ON CONFLICT (collection, rkey) DO UPDATE SET cid=excluded.cid",
"io.atcr.hold.scan", fmt.Sprintf("rkey-%d", i), fmt.Sprintf("cid-%d", i))
if e != nil && err == nil {
err = e
}
}
}()
csErr := writeShards(t, sqs, 4, 30) release = holdWriteLock(t, sqs.db)
close(stop) awaitBlockedWrite(t, "second pool write behind the carstore's lock", release, func() error {
riResult := <-idxErr _, err := other.Exec(
"INSERT INTO records (collection, rkey, cid) VALUES (?, ?, ?) ON CONFLICT (collection, rkey) DO UPDATE SET cid=excluded.cid",
"io.atcr.hold.scan", "rkey-1", "cid-1")
return err
})
}
for _, e := range []error{csErr, riResult} { // TestRawPoolFailsWithoutBusyTimeout is the control: a pool opened around the
if e == nil { // bare driver, with no busy_timeout, fails immediately under the same held
continue // lock. It proves the two tests above are observing busy_timeout and not some
} // other reason the writes happened to get through.
if isLockErr(e) { func TestRawPoolFailsWithoutBusyTimeout(t *testing.T) {
t.Fatalf("two-opener write hit a lock error: %v", e) path := filepath.Join(t.TempDir(), "db.sqlite3")
} guarded, err := OpenLocalDB(path)
t.Fatalf("two-opener write failed: %v", e) if err != nil {
t.Fatalf("OpenLocalDB: %v", err)
}
defer guarded.Close()
raw, err := sql.Open("libsql", localDSN(path))
if err != nil {
t.Fatalf("raw open: %v", err)
}
defer raw.Close()
if _, err := raw.Exec("CREATE TABLE IF NOT EXISTS records (k INTEGER PRIMARY KEY, v TEXT)"); err != nil {
t.Fatalf("create table: %v", err)
}
release := holdWriteLock(t, guarded)
defer release()
start := time.Now()
_, err = raw.Exec("INSERT INTO records (v) VALUES ('raw')")
took := time.Since(start)
if err == nil {
t.Fatalf("raw pool wrote through a held write lock in %s; the control cannot distinguish a regression", took)
}
if !isLockErr(err) {
t.Fatalf("raw pool failed for a reason other than the lock: %v", err)
}
if took > holdDuration {
t.Fatalf("raw pool waited %s before failing; expected an immediate SQLITE_BUSY with no busy_timeout", took)
} }
} }
+15
View File
@@ -0,0 +1,15 @@
package gc
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package labeler
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+15
View File
@@ -0,0 +1,15 @@
package hold
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}
+4
View File
@@ -11,6 +11,8 @@ import (
"path/filepath" "path/filepath"
"testing" "testing"
"atcr.io/internal/testlog"
"atcr.io/pkg/atproto" "atcr.io/pkg/atproto"
"atcr.io/pkg/auth/oauth" "atcr.io/pkg/auth/oauth"
"atcr.io/pkg/hold/pds" "atcr.io/pkg/hold/pds"
@@ -25,6 +27,8 @@ var (
// TestMain sets up shared resources for all OCI tests // TestMain sets up shared resources for all OCI tests
func TestMain(m *testing.M) { func TestMain(m *testing.M) {
testlog.Quiet() // see internal/testlog; -v or ATCR_TEST_LOGS=1 restores logs
// Create a temporary directory for shared test key // Create a temporary directory for shared test key
tmpDir, err := os.MkdirTemp("", "oci-test-shared-*") tmpDir, err := os.MkdirTemp("", "oci-test-shared-*")
if err != nil { if err != nil {
+4
View File
@@ -11,6 +11,8 @@ import (
"testing" "testing"
"time" "time"
"atcr.io/internal/testlog"
"atcr.io/pkg/atproto" "atcr.io/pkg/atproto"
"atcr.io/pkg/auth/oauth" "atcr.io/pkg/auth/oauth"
"atcr.io/pkg/s3" "atcr.io/pkg/s3"
@@ -250,6 +252,8 @@ func init() {
// Cleanup function to remove test files // Cleanup function to remove test files
func TestMain(m *testing.M) { func TestMain(m *testing.M) {
testlog.Quiet() // see internal/testlog; -v or ATCR_TEST_LOGS=1 restores logs
// Create a temporary directory for shared test key // Create a temporary directory for shared test key
tmpDir, err := os.MkdirTemp("", "pds-test-shared-*") tmpDir, err := os.MkdirTemp("", "pds-test-shared-*")
if err != nil { if err != nil {
+15
View File
@@ -0,0 +1,15 @@
package labeler
import (
"os"
"testing"
"atcr.io/internal/testlog"
)
// TestMain silences slog for the package (see internal/testlog); run with -v
// or ATCR_TEST_LOGS=1 to see the logs again.
func TestMain(m *testing.M) {
testlog.Quiet()
os.Exit(m.Run())
}