From d62e6b7fe832b99dd77d960a8decc826d6eef1ab Mon Sep 17 00:00:00 2001 From: dawn Date: Wed, 2 Sep 2026 22:31:51 +0900 Subject: [PATCH] spindle/db,spindle/mill: retry transient sqlite lock errors Signed-off-by: dawn --- spindle/db/mill_state.go | 217 +++++++++++++++++++--------------- spindle/db/mill_state_test.go | 50 ++++++++ spindle/db/retry.go | 30 +++++ spindle/db/retry_test.go | 75 ++++++++++++ spindle/mill/mill.go | 82 +++++++------ 5 files changed, 324 insertions(+), 130 deletions(-) create mode 100644 spindle/db/retry.go create mode 100644 spindle/db/retry_test.go diff --git a/spindle/db/mill_state.go b/spindle/db/mill_state.go index 7009a5376..c11becda5 100644 --- a/spindle/db/mill_state.go +++ b/spindle/db/mill_state.go @@ -111,96 +111,105 @@ func (d *DB) ListExecutorCursors() ([]ExecutorCursor, error) { } func (d *DB) SetOutboxEpoch(epoch string) error { - tx, err := d.Begin() - if err != nil { - return err - } - defer tx.Rollback() + return retrySQLite(func() error { + tx, err := d.Begin() + if err != nil { + return err + } + defer tx.Rollback() - if _, err := tx.Exec(`delete from mill_outbox_rows`); err != nil { - return err - } - if _, err := tx.Exec(`delete from mill_outbox_state`); err != nil { - return err - } - if _, err := tx.Exec(`insert into mill_outbox_state (epoch, next_seqno) values (?, 1)`, epoch); err != nil { - return err - } - return tx.Commit() + if _, err := tx.Exec(`delete from mill_outbox_rows`); err != nil { + return err + } + if _, err := tx.Exec(`delete from mill_outbox_state`); err != nil { + return err + } + if _, err := tx.Exec(`insert into mill_outbox_state (epoch, next_seqno) values (?, 1)`, epoch); err != nil { + return err + } + return tx.Commit() + }) } func (d *DB) AppendOutboxRow(payload []byte, control bool) (uint64, error) { - tx, err := d.Begin() - if err != nil { - return 0, err - } - defer tx.Rollback() - - var epoch string var nextSeqno uint64 - err = tx.QueryRow(` - update mill_outbox_state - set next_seqno = next_seqno + 1 - where epoch = (select epoch from mill_outbox_state limit 1) - returning epoch, next_seqno - 1 - `).Scan(&epoch, &nextSeqno) - if err == sql.ErrNoRows { - return 0, fmt.Errorf("no outbox epoch set") - } else if err != nil { - return 0, err - } + err := retrySQLite(func() error { + tx, err := d.Begin() + if err != nil { + return err + } + defer tx.Rollback() - byteSize := int64(len(payload)) - controlVal := 0 - if control { - controlVal = 1 - } + var epoch string + err = tx.QueryRow(` + update mill_outbox_state + set next_seqno = next_seqno + 1 + where epoch = (select epoch from mill_outbox_state limit 1) + returning epoch, next_seqno - 1 + `).Scan(&epoch, &nextSeqno) + if err == sql.ErrNoRows { + return fmt.Errorf("no outbox epoch set") + } else if err != nil { + return err + } - if _, err := tx.Exec( - `insert into mill_outbox_rows (epoch, seqno, payload, byte_size, control) values (?, ?, ?, ?, ?)`, - epoch, nextSeqno, payload, byteSize, controlVal, - ); err != nil { - return 0, err - } + byteSize := int64(len(payload)) + controlVal := 0 + if control { + controlVal = 1 + } - if err := tx.Commit(); err != nil { + if _, err := tx.Exec( + `insert into mill_outbox_rows (epoch, seqno, payload, byte_size, control) values (?, ?, ?, ?, ?)`, + epoch, nextSeqno, payload, byteSize, controlVal, + ); err != nil { + return err + } + + return tx.Commit() + }) + if err != nil { return 0, err } return nextSeqno, nil } func (d *DB) DeleteOutboxPrefix(ackedSeqno uint64) (OutboxDeletion, error) { - tx, err := d.Begin() - if err != nil { - return OutboxDeletion{}, err - } - defer tx.Rollback() + var deleted OutboxDeletion + err := retrySQLite(func() error { + tx, err := d.Begin() + if err != nil { + return err + } + defer tx.Rollback() - var epoch string - err = tx.QueryRow(`select epoch from mill_outbox_state limit 1`).Scan(&epoch) - if err == sql.ErrNoRows { - return OutboxDeletion{}, nil - } - if err != nil { - return OutboxDeletion{}, err - } + var epoch string + err = tx.QueryRow(`select epoch from mill_outbox_state limit 1`).Scan(&epoch) + if err == sql.ErrNoRows { + return nil + } + if err != nil { + return err + } - var deleted OutboxDeletion - if err := tx.QueryRow(` - select count(*), - coalesce(sum(byte_size), 0) - from mill_outbox_rows - where epoch = ? and seqno <= ? - `, epoch, ackedSeqno).Scan(&deleted.Rows, &deleted.Bytes); err != nil { - return OutboxDeletion{}, err - } - if _, err := tx.Exec( - `delete from mill_outbox_rows where epoch = ? and seqno <= ?`, - epoch, ackedSeqno, - ); err != nil { - return OutboxDeletion{}, err - } - if err := tx.Commit(); err != nil { + deleted = OutboxDeletion{} + if err := tx.QueryRow(` + select count(*), + coalesce(sum(byte_size), 0) + from mill_outbox_rows + where epoch = ? and seqno <= ? + `, epoch, ackedSeqno).Scan(&deleted.Rows, &deleted.Bytes); err != nil { + return err + } + if _, err := tx.Exec( + `delete from mill_outbox_rows where epoch = ? and seqno <= ?`, + epoch, ackedSeqno, + ); err != nil { + return err + } + return tx.Commit() + }) + if err != nil { return OutboxDeletion{}, err } return deleted, nil @@ -418,31 +427,53 @@ func (d *DB) ListPendingArtifacts() ([]PendingArtifact, error) { return res, rows.Err() } -func (d *DB) ApplyEventBatch(n *notifier.Notifier, fn func(tx *EventBatchTx) error) error { - err := eventstream.WithClock(func(insert eventstream.Inserter) error { - tx, err := d.Begin() - if err != nil { - return err - } - defer tx.Rollback() - - batchTx := &EventBatchTx{ - tx: tx, - db: d, - insert: insert, - } - if err := fn(batchTx); err != nil { - return err - } - return tx.Commit() +// ApplyEventBatch runs fn in a transaction and retries it while sqlite reports +// transient lock errors. fn may run more than once, so non-transactional side +// effects must be returned through its result: the result is published only +// from the attempt that commits, and a failed attempt's result is discarded. +func ApplyEventBatch[T any](d *DB, n *notifier.Notifier, fn func(tx *EventBatchTx) (T, error)) (T, error) { + var result T + err := retrySQLite(func() error { + return eventstream.WithClock(func(insert eventstream.Inserter) error { + tx, err := d.Begin() + if err != nil { + return err + } + defer tx.Rollback() + + batchTx := &EventBatchTx{ + tx: tx, + db: d, + insert: insert, + } + v, err := fn(batchTx) + if err != nil { + return err + } + if err := tx.Commit(); err != nil { + return err + } + result = v + return nil + }) }) if err != nil { - return err + var zero T + return zero, err } if n != nil { n.NotifyAll() } - return nil + return result, nil +} + +// ApplyEventBatch is the package-level ApplyEventBatch for callers without a +// result. +func (d *DB) ApplyEventBatch(n *notifier.Notifier, fn func(tx *EventBatchTx) error) error { + _, err := ApplyEventBatch(d, n, func(tx *EventBatchTx) (struct{}, error) { + return struct{}{}, fn(tx) + }) + return err } func (d *DB) DeleteMillLeasesByRepo(repoDid string) ([]string, error) { diff --git a/spindle/db/mill_state_test.go b/spindle/db/mill_state_test.go index 39adb1ddd..7ac1b7f66 100644 --- a/spindle/db/mill_state_test.go +++ b/spindle/db/mill_state_test.go @@ -8,6 +8,7 @@ import ( "sync" "testing" + "github.com/mattn/go-sqlite3" "tangled.org/core/notifier" "tangled.org/core/spindle/models" ) @@ -88,6 +89,55 @@ func TestExecutorCursors(t *testing.T) { } +func TestApplyEventBatchPublishesResultOnlyOnCommit(t *testing.T) { + d := newTestDB(t) + attempts := 0 + res, err := ApplyEventBatch(d, nil, func(tx *EventBatchTx) ([]string, error) { + attempts++ + if attempts == 1 { + // a result produced by a failed attempt must never escape + return []string{"ghost"}, sqlite3.Error{Code: sqlite3.ErrBusy} + } + if err := tx.AdvanceCursor("node-1", "inc-1", 7); err != nil { + return nil, err + } + return []string{"committed"}, nil + }) + if err != nil { + t.Fatalf("ApplyEventBatch: %v", err) + } + if attempts != 2 { + t.Fatalf("callback attempts = %d, want 2", attempts) + } + if len(res) != 1 || res[0] != "committed" { + t.Fatalf("result = %v, want [committed]", res) + } + if got, err := d.GetExecutorCursor("node-1", "inc-1"); err != nil || got != 7 { + t.Fatalf("cursor = %d, err = %v; want 7", got, err) + } +} + +func TestApplyEventBatchRetriesBusy(t *testing.T) { + d := newTestDB(t) + attempts := 0 + err := d.ApplyEventBatch(nil, func(tx *EventBatchTx) error { + attempts++ + if attempts == 1 { + return sqlite3.Error{Code: sqlite3.ErrBusy} + } + return tx.AdvanceCursor("node-1", "inc-1", 5) + }) + if err != nil { + t.Fatalf("ApplyEventBatch: %v", err) + } + if attempts != 2 { + t.Fatalf("callback attempts = %d, want 2", attempts) + } + if got, err := d.GetExecutorCursor("node-1", "inc-1"); err != nil || got != 5 { + t.Fatalf("cursor = %d, err = %v; want 5", got, err) + } +} + func TestCompleteMillLeaseIsAtomic(t *testing.T) { d := newTestDB(t) lease := MillLease{ diff --git a/spindle/db/retry.go b/spindle/db/retry.go new file mode 100644 index 000000000..7146c4705 --- /dev/null +++ b/spindle/db/retry.go @@ -0,0 +1,30 @@ +package db + +import ( + "errors" + "time" + + "github.com/mattn/go-sqlite3" +) + +const sqliteRetryAttempts = 4 + +func isSQLiteBusy(err error) bool { + var sqliteErr sqlite3.Error + if !errors.As(err, &sqliteErr) { + return false + } + return sqliteErr.Code == sqlite3.ErrBusy || sqliteErr.Code == sqlite3.ErrLocked +} + +func retrySQLite(fn func() error) error { + var err error + for attempt := 0; attempt < sqliteRetryAttempts; attempt++ { + err = fn() + if !isSQLiteBusy(err) || attempt == sqliteRetryAttempts-1 { + return err + } + time.Sleep(time.Duration(10*(1< m.nodeSeqno[currentKey] { - m.nodeSeqno[currentKey] = highestSeqno + if applied.highestSeqno > m.nodeSeqno[currentKey] { + m.nodeSeqno[currentKey] = applied.highestSeqno } m.mu.Unlock() - for _, leaseID := range terminalCancelIDs { + for _, leaseID := range applied.terminalCancelIDs { m.settleSessionCancel(sess, leaseID) } - for _, pt := range pendingTerminals { + for _, pt := range applied.pendingTerminals { if pt.lease.millRecordsTerminalMetrics && pt.ar.GetMillRecordsTerminalMetrics() && pt.ar.GetFailureClass() != "" && @@ -1555,7 +1563,7 @@ func (m *Mill) onEventBatch(sess *millSession, batch *millv1.EventBatch) error { } unlockLeases() - return m.sendAck(sess, highestSeqno) + return m.sendAck(sess, applied.highestSeqno) } func terminalMetricResult(status millv1.TerminalStatus) string { switch status { -- 2.51.2