From 4832f63e010ec9c5fcb88c0be9b14566f246ea5b Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sat, 14 Mar 2026 09:08:39 -0400 Subject: [PATCH 1/9] fix(outbox): downgrade "Skipping message" log from INFO to DEBUG This log fires for every outbox record that was already processed by another pod or goroutine, which is expected behavior in multi-replica deployments. At INFO level it produces thousands of noisy log lines that obscure meaningful output. Co-Authored-By: Claude Opus 4.6 --- outbox.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/outbox.go b/outbox.go index 97b2ace..6d17319 100644 --- a/outbox.go +++ b/outbox.go @@ -249,7 +249,7 @@ func (o outbox[T]) processMessageTx(ctx context.Context, id xid.ID) func(s Store } if errors.Is(err, errSkippingRecord) { - logger.InfoContext( + logger.DebugContext( ctx, "Skipping message", slog.String("reason", err.Error()), From 5c2b26d9ddac1f3d8f54a7d1a0a9b3fa104d0cf5 Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sun, 15 Mar 2026 18:41:49 -0400 Subject: [PATCH 2/9] feat(pg): fix connection leak and add keyset pagination for scalability MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The getRecordIDs loop used defer inside a for-loop, accumulating unclosed rows and uncancelled contexts across pages — exhausting the connection pool and causing context deadline exceeded errors. This extracts a fetchPage method so defer scopes correctly per page. Also replaces LIMIT/OFFSET pagination with keyset pagination (WHERE id > $1 ORDER BY id) which is O(1) per page and stable under concurrent modifications. Increases the channel buffer from 1 to 100 to reduce backpressure blocking. Adds configurable WithPageSize and WithChannelBufferSize store options. Co-Authored-By: Claude Opus 4.6 --- .gitignore | 4 +- store/pg/pg.go | 112 ++++++++++++++++++++++++++++++++++++------------- 2 files changed, 85 insertions(+), 31 deletions(-) diff --git a/.gitignore b/.gitignore index b68996b..662310d 100644 --- a/.gitignore +++ b/.gitignore @@ -21,4 +21,6 @@ go.work # Mock gen files -mock*.go \ No newline at end of file +mock*.go + +.claude/settings.local.json diff --git a/store/pg/pg.go b/store/pg/pg.go index 088d825..be12b12 100644 --- a/store/pg/pg.go +++ b/store/pg/pg.go @@ -22,11 +22,13 @@ type execQuerier interface { } type Store struct { - db execQuerier - tableName string - connStr string - chanName string - logger *slog.Logger + db execQuerier + tableName string + connStr string + chanName string + logger *slog.Logger + pageSize int + chanBufferSize int } type option func(s *Store) @@ -47,15 +49,33 @@ func WithLogger(logger *slog.Logger) option { } } +func WithPageSize(n int) option { + return func(s *Store) { + if n > 0 { + s.pageSize = n + } + } +} + +func WithChannelBufferSize(n int) option { + return func(s *Store) { + if n > 0 { + s.chanBufferSize = n + } + } +} + var _ outbox.Store = &Store{} func NewStore(db execQuerier, connStr string, opts ...option) (*Store, error) { s := &Store{ - db, - "outbox", - connStr, - "", - slog.Default(), + db: db, + tableName: "outbox", + connStr: connStr, + chanName: "", + logger: slog.Default(), + pageSize: 20, + chanBufferSize: 100, } for _, o := range opts { @@ -121,24 +141,29 @@ func (s Store) Listen() <-chan xid.ID { s.chanName, ) - idChan := make(chan xid.ID, 1) + // Ping the listner every 90 seconds to ensure it stays connected and receives notifications in a timely manner. + pingTicker := time.NewTicker(90 * time.Second) + go func(l *pq.Listener) { + for range pingTicker.C { + if err := l.Ping(); err != nil { + logger.Error("error pinging listener", "error", err) + } + } + }(listener) + + idChan := make(chan xid.ID, s.chanBufferSize) go func(l *pq.Listener) { for { - ids, err := s.getRecordIDs() + err := s.getRecordIDs(idChan) if err != nil { s.logger.Error("unable to get record ids", "error", err) continue } - for _, i := range ids { - idChan <- i - } - select { case <-l.Notify: // New record(s) available to process case <-time.After(90 * time.Second): - go l.Ping() // Check if there's more work available, just in case it takes a while // for the Listener to notice connection loss and reconnect. } @@ -148,36 +173,60 @@ func (s Store) Listen() <-chan xid.ID { return idChan } -func (s Store) getRecordIDs() ([]xid.ID, error) { - var res []xid.ID +func (s Store) getRecordIDs(idChan chan xid.ID) error { + var lastID string + for { + n, err := s.fetchPage(idChan, &lastID) + if err != nil { + return err + } + if n == 0 { + return nil + } + } +} + +func (s Store) fetchPage(idChan chan xid.ID, lastID *string) (int, error) { ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() - query := fmt.Sprintf(` - SELECT id FROM %s; - `, s.tableName) + var query string + var args []any + if *lastID == "" { + query = fmt.Sprintf(`SELECT id FROM %s ORDER BY id LIMIT %d;`, s.tableName, s.pageSize) + } else { + query = fmt.Sprintf(`SELECT id FROM %s WHERE id > $1 ORDER BY id LIMIT %d;`, s.tableName, s.pageSize) + args = []any{*lastID} + } - rows, err := s.db.QueryContext(ctx, query) + rows, err := s.db.QueryContext(ctx, query, args...) if err != nil { - return nil, err + return 0, err } defer rows.Close() + numRows := 0 for rows.Next() { var rawID string if err := rows.Scan(&rawID); err != nil { - return nil, err + return 0, err } id, err := xid.FromString(rawID) if err != nil { - return nil, err + return 0, err } - res = append(res, id) + idChan <- id + *lastID = rawID + numRows++ + } + + if err := rows.Err(); err != nil { + return 0, err } - return res, nil + return numRows, nil } func (s Store) GetWithLock(ctx context.Context, id xid.ID) (*outbox.Record, error) { @@ -236,8 +285,11 @@ func (s Store) ProcessTx(ctx context.Context, fn func(outbox.Store) bool) error } store := Store{ - db: tx, - tableName: s.tableName, + db: tx, + tableName: s.tableName, + logger: s.logger, + pageSize: s.pageSize, + chanBufferSize: s.chanBufferSize, } if success := fn(store); !success { From a99055240e5acac5d3e265e7ecd617ebb776511d Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sun, 15 Mar 2026 20:04:48 -0400 Subject: [PATCH 3/9] fix(pg): separate DB query from channel send and tune buffer defaults Split fetchPage into queryPage (DB-scoped) and fetchPage (channel send) so that channel backpressure cannot trigger the 15s context timeout. This was causing "pq: canceling statement due to user request" errors every 15s under load when the consumer couldn't drain the channel fast enough. Reduce defaults to chanBufferSize=5, pageSize=6 so the producer blocks on the 6th item of each page. This keeps the consumer saturated (5 items of runway while the next page is fetched) while minimizing the window where another pod could read the same buffered-but-unprocessed items. Co-Authored-By: Claude Opus 4.6 --- store/pg/pg.go | 46 +++++++++++++++++++++++++++++++++------------- 1 file changed, 33 insertions(+), 13 deletions(-) diff --git a/store/pg/pg.go b/store/pg/pg.go index be12b12..1edce54 100644 --- a/store/pg/pg.go +++ b/store/pg/pg.go @@ -74,8 +74,8 @@ func NewStore(db execQuerier, connStr string, opts ...option) (*Store, error) { connStr: connStr, chanName: "", logger: slog.Default(), - pageSize: 20, - chanBufferSize: 100, + pageSize: 6, + chanBufferSize: 5, } for _, o := range opts { @@ -187,46 +187,66 @@ func (s Store) getRecordIDs(idChan chan xid.ID) error { } func (s Store) fetchPage(idChan chan xid.ID, lastID *string) (int, error) { + ids, err := s.queryPage(*lastID) + if err != nil { + return 0, err + } + + // Send to the channel outside of the DB context so that blocking on a + // full channel does not hold open rows or trigger a context timeout. + for _, id := range ids { + idChan <- id + } + + if len(ids) > 0 { + *lastID = ids[len(ids)-1].String() + } + + return len(ids), nil +} + +// queryPage executes a single keyset-paginated query and returns up to +// pageSize IDs. The context timeout only covers the DB round-trip; channel +// backpressure cannot cause it to expire. +func (s Store) queryPage(afterID string) ([]xid.ID, error) { ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) defer cancel() var query string var args []any - if *lastID == "" { + if afterID == "" { query = fmt.Sprintf(`SELECT id FROM %s ORDER BY id LIMIT %d;`, s.tableName, s.pageSize) } else { query = fmt.Sprintf(`SELECT id FROM %s WHERE id > $1 ORDER BY id LIMIT %d;`, s.tableName, s.pageSize) - args = []any{*lastID} + args = []any{afterID} } rows, err := s.db.QueryContext(ctx, query, args...) if err != nil { - return 0, err + return nil, err } defer rows.Close() - numRows := 0 + var ids []xid.ID for rows.Next() { var rawID string if err := rows.Scan(&rawID); err != nil { - return 0, err + return nil, err } id, err := xid.FromString(rawID) if err != nil { - return 0, err + return nil, err } - idChan <- id - *lastID = rawID - numRows++ + ids = append(ids, id) } if err := rows.Err(); err != nil { - return 0, err + return nil, err } - return numRows, nil + return ids, nil } func (s Store) GetWithLock(ctx context.Context, id xid.ID) (*outbox.Record, error) { From 618fec792adbdbc8ebaef1e4496f7534ff9fb59f Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sun, 15 Mar 2026 21:20:52 -0400 Subject: [PATCH 4/9] fix(outbox): deduplicate in-flight IDs and fix rollback noise MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Add sync.Map-based in-flight tracking in dispatch() to skip IDs already being processed. Without this, getRecordIDs rescans return records still being processed, sending duplicate IDs to the channel. Each duplicate spawns a goroutine that acquires a pool connection via BeginTx, saturating the pool and causing 30s context timeouts across the system — including spilling into the sites-api connection pool. Use defer tx.Rollback() in ProcessTx (idiomatic Go pattern) so that when context cancellation auto-rolls back the tx, the deferred Rollback is a silent no-op instead of returning "sql: transaction has already been committed or rolled back". Co-Authored-By: Claude Opus 4.6 --- outbox.go | 12 +++++++++++- store/pg/pg.go | 3 ++- 2 files changed, 13 insertions(+), 2 deletions(-) diff --git a/outbox.go b/outbox.go index 6d17319..a83b145 100644 --- a/outbox.go +++ b/outbox.go @@ -10,6 +10,7 @@ import ( "log/slog" "os" "reflect" + "sync" "time" "github.com/rs/xid" @@ -147,11 +148,20 @@ func (o outbox[T]) SendTx(ctx context.Context, tx *sql.Tx, msg T) error { func (o outbox[T]) dispatch() { tokens := make(chan struct{}, o.numRoutines) + var inFlight sync.Map + for id := range o.store.Listen() { + if _, loaded := inFlight.LoadOrStore(id, struct{}{}); loaded { + continue + } + tokens <- struct{}{} go func(id xid.ID) { + defer func() { + inFlight.Delete(id) + <-tokens + }() o.process(id) - <-tokens }(id) } } diff --git a/store/pg/pg.go b/store/pg/pg.go index 1edce54..83f7680 100644 --- a/store/pg/pg.go +++ b/store/pg/pg.go @@ -303,6 +303,7 @@ func (s Store) ProcessTx(ctx context.Context, fn func(outbox.Store) bool) error if err != nil { return fmt.Errorf("unable to create transaction: %v", err) } + defer tx.Rollback() // no-op after Commit; silently handles context cancellation store := Store{ db: tx, @@ -313,7 +314,7 @@ func (s Store) ProcessTx(ctx context.Context, fn func(outbox.Store) bool) error } if success := fn(store); !success { - return tx.Rollback() + return nil // rollback handled by defer; real error already logged in callback } return tx.Commit() From f9eb759182a6608ad41247dbb6694eceefc76663 Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sun, 15 Mar 2026 21:30:34 -0400 Subject: [PATCH 5/9] close the listener --- store/pg/pg.go | 12 ++---------- 1 file changed, 2 insertions(+), 10 deletions(-) diff --git a/store/pg/pg.go b/store/pg/pg.go index 83f7680..918a0f1 100644 --- a/store/pg/pg.go +++ b/store/pg/pg.go @@ -141,18 +141,9 @@ func (s Store) Listen() <-chan xid.ID { s.chanName, ) - // Ping the listner every 90 seconds to ensure it stays connected and receives notifications in a timely manner. - pingTicker := time.NewTicker(90 * time.Second) - go func(l *pq.Listener) { - for range pingTicker.C { - if err := l.Ping(); err != nil { - logger.Error("error pinging listener", "error", err) - } - } - }(listener) - idChan := make(chan xid.ID, s.chanBufferSize) go func(l *pq.Listener) { + defer l.Close() for { err := s.getRecordIDs(idChan) if err != nil { @@ -164,6 +155,7 @@ func (s Store) Listen() <-chan xid.ID { case <-l.Notify: // New record(s) available to process case <-time.After(90 * time.Second): + l.Ping() // Check if there's more work available, just in case it takes a while // for the Listener to notice connection loss and reconnect. } From 83ae52b261910cea516f919ba68d08ccb2507cb5 Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sun, 15 Mar 2026 21:35:13 -0400 Subject: [PATCH 6/9] restore goroutine --- store/pg/pg.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/store/pg/pg.go b/store/pg/pg.go index 918a0f1..860da8b 100644 --- a/store/pg/pg.go +++ b/store/pg/pg.go @@ -155,7 +155,7 @@ func (s Store) Listen() <-chan xid.ID { case <-l.Notify: // New record(s) available to process case <-time.After(90 * time.Second): - l.Ping() + go l.Ping() // Check if there's more work available, just in case it takes a while // for the Listener to notice connection loss and reconnect. } From 9e7cbf93d7c3f8ead7efef17d6f4e46839c8b9b4 Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sun, 15 Mar 2026 22:27:50 -0400 Subject: [PATCH 7/9] feat(pg): add dedicated consumer connection pool to isolate from API pressure The outbox consumer (queryPage, ProcessTx) shared the same *sql.DB pool as the API resolvers. Long-running GraphQL mutations exhausted the pool, starving the consumer and triggering 15s context timeouts. Open a dedicated consumer pool from the existing connStr with a small MaxOpenConns (default 8). queryPage and ProcessTx.BeginTx now use this isolated pool, so API traffic cannot block outbox processing. Also adds: - WithMaxConsumerConns option to configure the pool size - Close() method with done channel for graceful shutdown of the Listen goroutine and consumer pool - Cancellable channel sends via done channel to prevent goroutine leaks Fixes test connStr to point at the correct database (outbox_test). Co-Authored-By: Claude Opus 4.6 --- store/pg/pg.go | 93 ++++++++++++++++++++++++++++++--------- store/pg/pg_suite_test.go | 3 +- store/pg/pg_test.go | 45 +++++++++++++++++++ 3 files changed, 118 insertions(+), 23 deletions(-) diff --git a/store/pg/pg.go b/store/pg/pg.go index 860da8b..faafa9c 100644 --- a/store/pg/pg.go +++ b/store/pg/pg.go @@ -22,13 +22,16 @@ type execQuerier interface { } type Store struct { - db execQuerier - tableName string - connStr string - chanName string - logger *slog.Logger - pageSize int - chanBufferSize int + db execQuerier + consumerDB *sql.DB // dedicated pool for consumer operations (queryPage, ProcessTx) + tableName string + connStr string + chanName string + logger *slog.Logger + pageSize int + chanBufferSize int + maxConsumerConns int + done chan struct{} } type option func(s *Store) @@ -65,17 +68,27 @@ func WithChannelBufferSize(n int) option { } } +func WithMaxConsumerConns(n int) option { + return func(s *Store) { + if n > 0 { + s.maxConsumerConns = n + } + } +} + var _ outbox.Store = &Store{} func NewStore(db execQuerier, connStr string, opts ...option) (*Store, error) { s := &Store{ - db: db, - tableName: "outbox", - connStr: connStr, - chanName: "", - logger: slog.Default(), - pageSize: 6, - chanBufferSize: 5, + db: db, + tableName: "outbox", + connStr: connStr, + chanName: "", + logger: slog.Default(), + pageSize: 6, + chanBufferSize: 5, + maxConsumerConns: 8, + done: make(chan struct{}), } for _, o := range opts { @@ -93,13 +106,41 @@ func NewStore(db execQuerier, connStr string, opts ...option) (*Store, error) { "_", ) + consumerDB, err := sql.Open("postgres", connStr) + if err != nil { + return nil, fmt.Errorf("open consumer db pool: %w", err) + } + consumerDB.SetMaxOpenConns(s.maxConsumerConns) + consumerDB.SetMaxIdleConns(s.maxConsumerConns) + if err := consumerDB.Ping(); err != nil { + consumerDB.Close() + return nil, fmt.Errorf("ping consumer db pool: %w", err) + } + s.consumerDB = consumerDB + if err := s.init(); err != nil { + consumerDB.Close() return nil, err } return s, nil } +// Close signals the Listen goroutine to stop and shuts down the dedicated +// consumer connection pool. It does not close the caller-owned db passed to NewStore. +func (s *Store) Close() error { + select { + case <-s.done: + // already closed + default: + close(s.done) + } + if s.consumerDB != nil { + return s.consumerDB.Close() + } + return nil +} + func (s Store) CreateRecordTx(ctx context.Context, tx *sql.Tx, r outbox.Record) (*outbox.Record, error) { query := fmt.Sprintf(` INSERT INTO %s VALUES ($1, $2); @@ -145,13 +186,20 @@ func (s Store) Listen() <-chan xid.ID { go func(l *pq.Listener) { defer l.Close() for { + select { + case <-s.done: + return + default: + } + err := s.getRecordIDs(idChan) if err != nil { s.logger.Error("unable to get record ids", "error", err) - continue } select { + case <-s.done: + return case <-l.Notify: // New record(s) available to process case <-time.After(90 * time.Second): @@ -187,7 +235,11 @@ func (s Store) fetchPage(idChan chan xid.ID, lastID *string) (int, error) { // Send to the channel outside of the DB context so that blocking on a // full channel does not hold open rows or trigger a context timeout. for _, id := range ids { - idChan <- id + select { + case idChan <- id: + case <-s.done: + return 0, nil + } } if len(ids) > 0 { @@ -213,7 +265,7 @@ func (s Store) queryPage(afterID string) ([]xid.ID, error) { args = []any{afterID} } - rows, err := s.db.QueryContext(ctx, query, args...) + rows, err := s.consumerDB.QueryContext(ctx, query, args...) if err != nil { return nil, err } @@ -284,14 +336,11 @@ func (s Store) Delete(ctx context.Context, id xid.ID) error { } func (s Store) ProcessTx(ctx context.Context, fn func(outbox.Store) bool) error { - db, ok := s.db.(interface { - BeginTx(context.Context, *sql.TxOptions) (*sql.Tx, error) - }) - if !ok { + if s.consumerDB == nil { return errors.New("process transaction can only be called at the parent level") } - tx, err := db.BeginTx(ctx, nil) + tx, err := s.consumerDB.BeginTx(ctx, nil) if err != nil { return fmt.Errorf("unable to create transaction: %v", err) } diff --git a/store/pg/pg_suite_test.go b/store/pg/pg_suite_test.go index 2583521..5e974a7 100644 --- a/store/pg/pg_suite_test.go +++ b/store/pg/pg_suite_test.go @@ -13,7 +13,7 @@ import ( var ( db *sql.DB - connStr = getConfig().GenerateAddress() + connStr string ) func TestIntegration(t *testing.T) { @@ -28,6 +28,7 @@ var _ = BeforeEach(func() { var err error cfg := getConfig() cfg.Database = "outbox_test" + connStr = cfg.GenerateAddress() db, err = pghelpers.ConnectPostgres(*cfg) Expect(err).To(Succeed()) }) diff --git a/store/pg/pg_test.go b/store/pg/pg_test.go index 9d7a58f..1577750 100644 --- a/store/pg/pg_test.go +++ b/store/pg/pg_test.go @@ -21,6 +21,38 @@ var _ = Describe("pgStore", func() { subject, err := NewStore(db, connStr) Expect(err).To(Succeed()) Expect(subject).ToNot(BeNil()) + defer subject.Close() + }) + + It("should create a separate consumer pool", func() { + subject, err := NewStore(db, connStr) + Expect(err).To(Succeed()) + defer subject.Close() + + Expect(subject.consumerDB).ToNot(BeNil()) + Expect(subject.consumerDB.Ping()).To(Succeed()) + }) + + It("should respect WithMaxConsumerConns option", func() { + subject, err := NewStore(db, connStr, WithMaxConsumerConns(3)) + Expect(err).To(Succeed()) + defer subject.Close() + + stats := subject.consumerDB.Stats() + Expect(stats.MaxOpenConnections).To(Equal(3)) + }) + }) + + Describe("#Close", func() { + It("should close the consumer pool without error", func() { + subject, err := NewStore(db, connStr) + Expect(err).To(Succeed()) + + err = subject.Close() + Expect(err).To(Succeed()) + + err = subject.consumerDB.Ping() + Expect(err).To(HaveOccurred()) }) }) @@ -39,6 +71,10 @@ var _ = Describe("pgStore", func() { tx, _ = db.BeginTx(ctx, nil) }) + AfterEach(func() { + subject.Close() + }) + It("should save the provided record on a successfull transaction", func() { res, err := subject.CreateRecordTx(ctx, tx, record) Expect(err).To(Succeed()) @@ -74,6 +110,10 @@ var _ = Describe("pgStore", func() { Expect(err).To(Succeed()) }) + AfterEach(func() { + subject.Close() + }) + It("should return a record on a valid id", func() { subject.ProcessTx(ctx, func(s outbox.Store) bool { res, err := s.GetWithLock(ctx, id) @@ -102,6 +142,7 @@ var _ = Describe("pgStore", func() { subject = createStore() ids = []xid.ID{xid.New(), xid.New(), xid.New()} ) + defer subject.Close() for _, id := range ids { _ = insertRecord(db, outbox.Record{ID: id, Message: []byte("data")}) @@ -132,6 +173,7 @@ var _ = Describe("pgStore", func() { ctx = context.Background() tx, _ = db.BeginTx(ctx, nil) ) + defer subject.Close() res, err := subject.CreateRecordTx(ctx, tx, outbox.Record{}) tx.Commit() @@ -147,6 +189,7 @@ var _ = Describe("pgStore", func() { subject = createStore() id = xid.New() ) + defer subject.Close() _ = insertRecord(db, outbox.Record{ID: id, Message: []byte("data")}) @@ -164,6 +207,7 @@ var _ = Describe("pgStore", func() { subject = createStore() id = xid.New() ) + defer subject.Close() _ = insertRecord(db, outbox.Record{ID: id, Message: []byte("data")}) @@ -220,6 +264,7 @@ var _ = Describe("pgStore", func() { ) _ = insertRecord(db, record) + defer subject.Close() err := subject.Update(context.Background(), &record) Expect(err).To(Succeed()) From fa780be8ffe9d4972810e3ec115e4689bfe99880 Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sun, 15 Mar 2026 22:50:09 -0400 Subject: [PATCH 8/9] refactor(pg): route all parent store operations through consumer pool Set s.db = consumerDB on the parent store so init(), Delete, and Update bypass the caller's shared pool entirely. The db parameter in NewStore is kept for backward compatibility but no longer used. Also removes redundant pre-check select on done channel in Listen loop. Co-Authored-By: Claude Opus 4.6 --- store/pg/pg.go | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/store/pg/pg.go b/store/pg/pg.go index faafa9c..88ac0c6 100644 --- a/store/pg/pg.go +++ b/store/pg/pg.go @@ -117,6 +117,7 @@ func NewStore(db execQuerier, connStr string, opts ...option) (*Store, error) { return nil, fmt.Errorf("ping consumer db pool: %w", err) } s.consumerDB = consumerDB + s.db = consumerDB if err := s.init(); err != nil { consumerDB.Close() @@ -186,12 +187,6 @@ func (s Store) Listen() <-chan xid.ID { go func(l *pq.Listener) { defer l.Close() for { - select { - case <-s.done: - return - default: - } - err := s.getRecordIDs(idChan) if err != nil { s.logger.Error("unable to get record ids", "error", err) From 9b230ec227d8df259070d9bd8bd10924e697f3a6 Mon Sep 17 00:00:00 2001 From: Jason White <22136798+strobus@users.noreply.github.com> Date: Sun, 15 Mar 2026 22:59:55 -0400 Subject: [PATCH 9/9] remove consumerDB field --- store/pg/pg.go | 21 ++++++++++++++------- store/pg/pg_test.go | 8 ++++---- 2 files changed, 18 insertions(+), 11 deletions(-) diff --git a/store/pg/pg.go b/store/pg/pg.go index 88ac0c6..9686ddc 100644 --- a/store/pg/pg.go +++ b/store/pg/pg.go @@ -23,7 +23,6 @@ type execQuerier interface { type Store struct { db execQuerier - consumerDB *sql.DB // dedicated pool for consumer operations (queryPage, ProcessTx) tableName string connStr string chanName string @@ -106,6 +105,8 @@ func NewStore(db execQuerier, connStr string, opts ...option) (*Store, error) { "_", ) + // dedicated pool for consumer operations (queryPage, ProcessTx) so that Listen can keep polling + // for new work even if the caller is doing long-running work or has a slow connection consumerDB, err := sql.Open("postgres", connStr) if err != nil { return nil, fmt.Errorf("open consumer db pool: %w", err) @@ -116,7 +117,6 @@ func NewStore(db execQuerier, connStr string, opts ...option) (*Store, error) { consumerDB.Close() return nil, fmt.Errorf("ping consumer db pool: %w", err) } - s.consumerDB = consumerDB s.db = consumerDB if err := s.init(); err != nil { @@ -136,8 +136,12 @@ func (s *Store) Close() error { default: close(s.done) } - if s.consumerDB != nil { - return s.consumerDB.Close() + if s.db != nil { + if db, ok := s.db.(interface { + Close() error + }); ok { + db.Close() + } } return nil } @@ -260,7 +264,7 @@ func (s Store) queryPage(afterID string) ([]xid.ID, error) { args = []any{afterID} } - rows, err := s.consumerDB.QueryContext(ctx, query, args...) + rows, err := s.db.QueryContext(ctx, query, args...) if err != nil { return nil, err } @@ -331,11 +335,14 @@ func (s Store) Delete(ctx context.Context, id xid.ID) error { } func (s Store) ProcessTx(ctx context.Context, fn func(outbox.Store) bool) error { - if s.consumerDB == nil { + db, ok := s.db.(interface { + BeginTx(context.Context, *sql.TxOptions) (*sql.Tx, error) + }) + if !ok { return errors.New("process transaction can only be called at the parent level") } - tx, err := s.consumerDB.BeginTx(ctx, nil) + tx, err := db.BeginTx(ctx, nil) if err != nil { return fmt.Errorf("unable to create transaction: %v", err) } diff --git a/store/pg/pg_test.go b/store/pg/pg_test.go index 1577750..aedb943 100644 --- a/store/pg/pg_test.go +++ b/store/pg/pg_test.go @@ -29,8 +29,8 @@ var _ = Describe("pgStore", func() { Expect(err).To(Succeed()) defer subject.Close() - Expect(subject.consumerDB).ToNot(BeNil()) - Expect(subject.consumerDB.Ping()).To(Succeed()) + Expect(subject.db).ToNot(BeNil()) + Expect(subject.db.(interface{ Ping() error }).Ping()).To(Succeed()) }) It("should respect WithMaxConsumerConns option", func() { @@ -38,7 +38,7 @@ var _ = Describe("pgStore", func() { Expect(err).To(Succeed()) defer subject.Close() - stats := subject.consumerDB.Stats() + stats := subject.db.(interface{ Stats() sql.DBStats }).Stats() Expect(stats.MaxOpenConnections).To(Equal(3)) }) }) @@ -51,7 +51,7 @@ var _ = Describe("pgStore", func() { err = subject.Close() Expect(err).To(Succeed()) - err = subject.consumerDB.Ping() + err = subject.db.(interface{ Ping() error }).Ping() Expect(err).To(HaveOccurred()) }) })