diff --git a/spindle/db/mill_state.go b/spindle/db/mill_state.go index 95310fcf..1524e2aa 100644 --- a/spindle/db/mill_state.go +++ b/spindle/db/mill_state.go @@ -123,7 +123,12 @@ func (d *DB) AppendOutboxRow(payload []byte, control bool) (uint64, error) { var epoch string var nextSeqno uint64 - err = tx.QueryRow(`select epoch, next_seqno from mill_outbox_state limit 1`).Scan(&epoch, &nextSeqno) + 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 { @@ -143,13 +148,6 @@ func (d *DB) AppendOutboxRow(payload []byte, control bool) (uint64, error) { return 0, err } - if _, err := tx.Exec( - `update mill_outbox_state set next_seqno = ? where epoch = ?`, - nextSeqno+1, epoch, - ); err != nil { - return 0, err - } - if err := tx.Commit(); err != nil { return 0, err } diff --git a/spindle/db/mill_state_test.go b/spindle/db/mill_state_test.go index c4990b9d..9902c78c 100644 --- a/spindle/db/mill_state_test.go +++ b/spindle/db/mill_state_test.go @@ -4,6 +4,7 @@ import ( "context" "fmt" "path/filepath" + "sync" "testing" "tangled.org/core/notifier" @@ -268,6 +269,57 @@ func TestOutboxPrefixAck(t *testing.T) { } } +func TestOutboxConcurrentAppend(t *testing.T) { + d := newTestDB(t) + if err := d.SetOutboxEpoch("inc-1"); err != nil { + t.Fatalf("SetOutboxEpoch: %v", err) + } + + const count = 32 + type result struct { + seqno uint64 + err error + } + start := make(chan struct{}) + results := make(chan result, count) + var wg sync.WaitGroup + for i := range count { + wg.Add(1) + go func() { + defer wg.Done() + <-start + seqno, err := d.AppendOutboxRow([]byte(fmt.Sprintf("msg%d", i)), true) + results <- result{seqno: seqno, err: err} + }() + } + close(start) + wg.Wait() + close(results) + + seen := make(map[uint64]struct{}, count) + for result := range results { + if result.err != nil { + t.Errorf("AppendOutboxRow: %v", result.err) + continue + } + if _, ok := seen[result.seqno]; ok { + t.Errorf("duplicate outbox seqno %d", result.seqno) + } + seen[result.seqno] = struct{}{} + } + if len(seen) != count { + t.Fatalf("successful appends = %d, want %d", len(seen), count) + } + + rows, err := d.ListOutboxRows() + if err != nil { + t.Fatalf("ListOutboxRows: %v", err) + } + if len(rows) != count { + t.Fatalf("persisted outbox rows = %d, want %d", len(rows), count) + } +} + func TestCompositeCursors(t *testing.T) { d := newTestDB(t)