Add _txlock=immediate so batch writes wait instead of dropping (closes #25)
check / check (push) Failing after 0s
check / check (push) Failing after 0s
Batch flush paths read before they write, so a deferred transaction starts as a reader and must upgrade to the write lock on its first INSERT/UPDATE. When the background maintainer holds the write lock for a WAL checkpoint, that upgrade fails immediately with "database is locked" and the busy timeout does not apply, so the batch is dropped. Adding _txlock=immediate to the DSN makes every transaction take the write lock at BEGIN, so it waits up to busy_timeout instead of failing. A regression test drives batch writes against a running checkpoint loop and fails with "database is locked" without the change. Model: opus-4-8
This commit was merged in pull request #26.
This commit is contained in:
@@ -88,9 +88,13 @@ func New(cfg *config.Config, logger *logger.Logger) (*Database, error) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Per-connection SQLite settings go in the DSN so every pooled connection
|
// Per-connection SQLite settings go in the DSN so every pooled connection
|
||||||
// gets them, not just the one that runs the Initialize pragmas.
|
// gets them, not just the one that runs the Initialize pragmas. _txlock=
|
||||||
|
// immediate makes every transaction take the write lock at BEGIN. Without it
|
||||||
|
// a transaction that reads before writing starts as a reader and, when it
|
||||||
|
// then writes while another connection holds the write lock, fails at once
|
||||||
|
// with "database is locked" without waiting for _busy_timeout.
|
||||||
dsn := fmt.Sprintf(
|
dsn := fmt.Sprintf(
|
||||||
"file:%s?_cache_size=%d&_synchronous=OFF&_busy_timeout=%d&_journal_mode=WAL",
|
"file:%s?_cache_size=%d&_synchronous=OFF&_busy_timeout=%d&_journal_mode=WAL&_txlock=immediate",
|
||||||
dbPath,
|
dbPath,
|
||||||
sqliteCacheSizeKiB,
|
sqliteCacheSizeKiB,
|
||||||
sqliteBusyTimeoutMs,
|
sqliteBusyTimeoutMs,
|
||||||
|
|||||||
@@ -4,7 +4,9 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"database/sql"
|
"database/sql"
|
||||||
"net"
|
"net"
|
||||||
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
"git.eeqj.de/sneak/routewatch/internal/config"
|
"git.eeqj.de/sneak/routewatch/internal/config"
|
||||||
"git.eeqj.de/sneak/routewatch/internal/logger"
|
"git.eeqj.de/sneak/routewatch/internal/logger"
|
||||||
@@ -18,6 +20,18 @@ const tempStoreMemory = 2
|
|||||||
// once so each is a distinct SQLite connection that parsed the DSN.
|
// once so each is a distinct SQLite connection that parsed the DSN.
|
||||||
const heldConnections = 5
|
const heldConnections = 5
|
||||||
|
|
||||||
|
// Parameters for the checkpoint-contention regression test.
|
||||||
|
const (
|
||||||
|
// contentionIterations is how many batch writes race the checkpoint loop.
|
||||||
|
contentionIterations = 400
|
||||||
|
// contendedASNCount is the small set of ASNs the batches reuse, so most
|
||||||
|
// batches update existing rows and exercise the read-before-write path.
|
||||||
|
contendedASNCount = 16
|
||||||
|
// asnSecondBand offsets a second ASN per batch so each batch writes more
|
||||||
|
// than one row.
|
||||||
|
asnSecondBand = 100
|
||||||
|
)
|
||||||
|
|
||||||
func TestIPToUint32(t *testing.T) {
|
func TestIPToUint32(t *testing.T) {
|
||||||
tests := []struct {
|
tests := []struct {
|
||||||
name string
|
name string
|
||||||
@@ -361,6 +375,56 @@ func TestConnectionPoolPragmas(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// TestBatchWriteDuringCheckpoint reproduces issue #25. A batch write reads
|
||||||
|
// (SELECT) before it writes (INSERT/UPDATE). Under the default deferred locking
|
||||||
|
// the transaction begins as a reader and, when the maintainer's WAL checkpoint
|
||||||
|
// holds the write lock, its upgrade to writer fails immediately with "database
|
||||||
|
// is locked" without honouring busy_timeout, dropping the batch. With
|
||||||
|
// _txlock=immediate the transaction takes the write lock at BEGIN and waits, so
|
||||||
|
// no batch is dropped. The checkpoint runs without the Database mutex, exactly
|
||||||
|
// as the background maintainer does in production.
|
||||||
|
func TestBatchWriteDuringCheckpoint(t *testing.T) {
|
||||||
|
cfg := &config.Config{StateDir: t.TempDir()}
|
||||||
|
|
||||||
|
db, err := New(cfg, logger.New())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to create database: %v", err)
|
||||||
|
}
|
||||||
|
defer func() { _ = db.Close() }()
|
||||||
|
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
|
||||||
|
var wg sync.WaitGroup
|
||||||
|
wg.Add(1)
|
||||||
|
go func() {
|
||||||
|
defer wg.Done()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
_ = db.Checkpoint(ctx) // errors are the checkpoint's own to absorb
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
ts := time.Now().UTC()
|
||||||
|
for i := 0; i < contentionIterations; i++ {
|
||||||
|
asns := map[int]time.Time{
|
||||||
|
i % contendedASNCount: ts,
|
||||||
|
(i % contendedASNCount) + asnSecondBand: ts,
|
||||||
|
}
|
||||||
|
if err := db.GetOrCreateASNBatch(asns); err != nil {
|
||||||
|
cancel()
|
||||||
|
wg.Wait()
|
||||||
|
t.Fatalf("batch write failed under checkpoint contention: %v", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
cancel()
|
||||||
|
wg.Wait()
|
||||||
|
}
|
||||||
|
|
||||||
func BenchmarkIPToUint32(b *testing.B) {
|
func BenchmarkIPToUint32(b *testing.B) {
|
||||||
ip := net.ParseIP("192.168.1.1")
|
ip := net.ParseIP("192.168.1.1")
|
||||||
b.ResetTimer()
|
b.ResetTimer()
|
||||||
|
|||||||
Reference in New Issue
Block a user