Sitelet https://github.com/basecamp/basecamp-cli/pull/820
Skip to content

Report when the connector's backlog warns, pauses and recovers - #820

Open
robzolkos wants to merge 2 commits into
mainfrom
connect-e2e/queue-report
Open

robzolkos wants to merge 2 commits into
mainfrom
connect-e2e/queue-report

Conversation

@robzolkos

@robzolkos robzolkos commented Oct 2, 2026 •

Copy link
Copy Markdown
Collaborator

Why

The connector's backlog sits between intake and admission. At 1,000 events it is worth a warning, and at 10,000 intake stops reading the feed faster than admission makes room. The connector package says this pause "is reported rather than hidden".

In production nothing reported it. The queue decides four edges (warn, pause, resume, recover) but handed them only to optional callbacks, and runConnect sets none. An operator watching stderr had no way to know the backlog was growing or that the feed read was held back.

Found while researching end-to-end tests for basecamp connect.

What

  • The queue now reports each edge on its logger, before it runs any callback:
    • Warn (WARN): the backlog reached its warning depth; admission is falling behind the feed. The usual cause is slow or failing Basecamp reads during admission; basecamp connect status shows the backlog of records by state. Attributes depth, warn_at, recover_at, pause_at.
    • Pause (WARN): the backlog reached its pause depth; the next event waits for room, and the feed is read only as fast as admission makes room, so the backlog cannot outgrow memory. plus the same next step. Attributes depth (what the queue holds, not counting the waiting offer), pause_at, resume_at.
    • Resume (INFO): the backlog drained to half its pause depth; the feed is read at full speed again. Attributes depth, pause_at, resume_at.
    • Recover (INFO): the backlog fell back below its warning depth. Attributes depth, warn_at, recover_at.
  • Each edge clears at a band, so a backlog that hovers at a threshold is one report each way, not one per event:
    • A warning clears at half the warning depth (recover_at, 500 by default).
    • A pause is raised when an offer first waits for room, and reported over only once no offer waits and the backlog has drained to half the pause depth (resume_at, 5,000 by default). At capacity each take lets one waiting offer in and the next offer waits again; that is now one pause line, not two lines per feed event.
    • Half needs no clock, a hovering backlog cannot cross it by itself, and both bands are zero or more, so a drained queue always ends recovered and resumed.
    • The feed's blocking at capacity does not change. Only which edges are reported moves. Paused() still reports the live state.
  • The report goes through the same ordered drain as the callbacks, so the lines come out in the order the edges were decided. When one take clears both a pause and a warning, the resume is told first.
  • Intake already gives the queue the connector's logger when it adopts the queue. So every caller that builds intake gets these lines, and none can forget to wire them. runConnect does not change.
  • A callback that panics is still contained, and the edge has already been reported before it runs.
  • The callbacks (OnWarn, OnRecover, OnPause, OnResume) follow the same edges. OnPause gets the held depth.
  • Depth is not persisted. connect status still reads only the ledger.

Tests

  • TestIntakeReportsTheBacklogsEdgesOnItsLogger builds intake the way production does (a queue with no callbacks, the connector's logger in Options). It fills the queue past warn and pause, drains it, and checks all four lines in order, with their levels, depths, thresholds and bands, and that no resume is reported while the backlog is still at its pause depth. It fails on main: nothing is logged.
  • TestABacklogHoveringAtItsThresholdsIsReportedOnceEachWay oscillates around the warning depth and then at capacity (one offer waiting, one take, the next offer waiting) 20 times each, and checks exactly one warn, one pause, one resume and one recover, in callbacks and in log lines, each at its band.
  • TestAPanickingCallbackDoesNotSilenceTheReport: a warning callback that panics still leaves the warning line, then the panic line. It fails on main.
  • Existing pause tests (TestQueuePausesTheCallerAtTheThreshold, TestPauseAndResumeAreObservedInTheOrderTheyHappened, TestAPanickingPauseCallbackLeavesNoPhantomBacklog, TestQueuePauseTracksEveryBlockedOffer) now drain to the resume band before they expect a resume.
  • The queue tests pass with -race -count=20. bin/ci is green.

Summary by cubic

The backlog between intake and admission warned, paused, resumed, and recovered only through optional callbacks, and runConnect set none, so in production nothing reported it. The queue now reports each edge on its logger before running any callback.

New Features

  • Warn and pause are logged at WARN, resume and recover at INFO, each with the depth and the threshold it crossed.
  • Intake already hands the queue the connector's logger, so every caller gets the reports without wiring one up.
  • A panicking callback can't silence the report; it's emitted first and the panic is contained.
  • Depth isn't persisted; connect status still reads only the ledger.

Written for commit 54fb6b0. Summary will update on new commits.

Review in cubic

The queue between intake and admission decides four edges: the backlog
crossing its warning depth, an offer waiting for room (which stops the
feed read), that wait ending, and the backlog falling back below the
warning depth. It handed them only to optional callbacks, and
runConnect sets none, so in production every edge went nowhere. The
package says the pause is reported rather than hidden; it was hidden.

The queue now reports each edge on its logger before running any
callback, with the depth and the threshold crossed, and says plainly
when reading the feed is paused and why. Intake already adopts the
connector's logger for the queue, so every caller that builds intake
gets the report and none can forget to wire it. A callback that panics
is contained after the report, so it cannot silence it.
Copilot AI balanced review requested due to automatic review settings October 2, 2026 17:16
@github-actions github-actions Bot added the tests Tests (unit and e2e) label Oct 2, 2026

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Pause and resume logs can incorrectly report live-feed state for background or canceled offers.

Review effort: Balanced
Findings: 1 Medium severity

Open (1)
What changed in this PR

Adds ordered backlog edge logging so connector operators can observe warnings, pauses, resumes, and recoveries.

Changes:

  • Logs queue transitions with depths and thresholds.
  • Preserves callback ordering and panic containment.
  • Adds integration and callback-panic coverage.

[!TIP]
If you aren't ready for review, convert to a draft PR.
Click "Convert to draft" or run gh pr ready --undo.
Click "Ready for review" or run gh pr ready to reengage.

File Description
internal/​connector/​queue.go Reports backlog transitions through the queue logger.
internal/​connector/​intake.go Documents intake’s logger adoption.
internal/​connector/​queue_test.go Tests reporting before a panicking callback.
internal/​connector/​seam_contract_test.go Tests ordered edge reporting through intake.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread internal/connector/queue.go Outdated
Comment on lines +336 to +337
log.Warn("the backlog reached its pause depth; reading the feed is paused until admission makes room, so the backlog cannot outgrow memory",
"depth", edge.depth, "pause_at", pauseAt)

@cubic-dev-ai cubic-dev-ai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

4 issues found across 4 files

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. If appropriate, use sub-agents to investigate and fix each issue separately.


<file name="internal/connector/queue.go">

<violation number="1" location="internal/connector/queue.go:311">
P3: The recover arm in `fire` is now registered after `q.logger()` and `q.report(log, edge)` run, so a panic in the reporting path is no longer contained. Before this change, `q.logger()` was only called inside the already-recovered region, so nothing untrusted ran unguarded. `report` calls into the caller-supplied slog handler (intake adopts the production logger via `adoptLogger`), and a handler that panics would escape `fire`, propagate through `deliver`'s drain into the `Offer`/`Take` caller, and take down the feed's delivery path — exactly what the panic containment in this file is for. The new comment on `fire` even claims "one that panics cannot take the report down with it," but it only protects against callback panics, not report panics. Arm the `defer` before the report so the report and callback both run contained.</violation>

<violation number="2" location="internal/connector/queue.go:336">
P2: A repair or requeue offer can trigger this edge while the feed is still being consumed, so this message can falsely report that reading the feed paused. Describe the blocked offer or full queue instead.</violation>
</file>

<file name="internal/connector/seam_contract_test.go">

<violation number="1" location="internal/connector/seam_contract_test.go:256">
P2: The pause and resume assertions only check that `depth` exists, so incorrect reported depths still pass. Assert 4 for the pause and 3 for the resume.

(Based on your team's feedback about structured JSON assertions.)</violation>
</file>

<file name="internal/connector/queue_test.go">

<violation number="1" location="internal/connector/queue_test.go:310">
P2: These post-`Offer` checks pass even if the warning report moves after the callback panic is recovered. Assert inside `OnWarn` that the warning record already exists before panicking to protect the report-before-callback guarantee.</violation>
</file>

Reply with feedback, questions, or to request a fix.

Re-trigger cubic

Comment thread internal/connector/queue.go Outdated
log.Info("the backlog fell back below its warning depth",
"depth", edge.depth, "warn_at", q.warnAt)
case edgePause:
log.Warn("the backlog reached its pause depth; reading the feed is paused until admission makes room, so the backlog cannot outgrow memory",

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: A repair or requeue offer can trigger this edge while the feed is still being consumed, so this message can falsely report that reading the feed paused. Describe the blocked offer or full queue instead.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At internal/connector/queue.go, line 336:

<comment>A repair or requeue offer can trigger this edge while the feed is still being consumed, so this message can falsely report that reading the feed paused. Describe the blocked offer or full queue instead.</comment>

<file context>
@@ -310,6 +320,27 @@ func (q *Queue) fire(edge queueEdge) {
+		log.Info("the backlog fell back below its warning depth",
+			"depth", edge.depth, "warn_at", q.warnAt)
+	case edgePause:
+		log.Warn("the backlog reached its pause depth; reading the feed is paused until admission makes room, so the backlog cannot outgrow memory",
+			"depth", edge.depth, "pause_at", pauseAt)
+	case edgeResume:
</file context>
Suggested change
log.Warn("the backlog reached its pause depth; reading the feed is paused until admission makes room, so the backlog cannot outgrow memory",
log.Warn("the backlog reached its pause depth; an offer is waiting for room",

assert.Equal(t, "WARN", lines[1]["level"])
assert.Contains(t, lines[1]["msg"], "reading the feed is paused")
assert.EqualValues(t, 3, lines[1]["pause_at"])
assert.Contains(t, lines[1], "depth")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: The pause and resume assertions only check that depth exists, so incorrect reported depths still pass. Assert 4 for the pause and 3 for the resume.

(Based on your team's feedback about structured JSON assertions.)

View Feedback

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At internal/connector/seam_contract_test.go, line 256:

<comment>The pause and resume assertions only check that `depth` exists, so incorrect reported depths still pass. Assert 4 for the pause and 3 for the resume.

(Based on your team's feedback about structured JSON assertions.) </comment>

<file context>
@@ -202,6 +204,85 @@ func TestIntakeGivesTheQueueItsLogger(t *testing.T) {
+	assert.Equal(t, "WARN", lines[1]["level"])
+	assert.Contains(t, lines[1]["msg"], "reading the feed is paused")
+	assert.EqualValues(t, 3, lines[1]["pause_at"])
+	assert.Contains(t, lines[1], "depth")
+
+	assert.Equal(t, "INFO", lines[2]["level"])
</file context>

require.NoError(t, err)
var logs bytes.Buffer
queue.SetLogger(slog.New(slog.NewTextHandler(&logs, nil)))
queue.OnWarn = func(int) { panic("a warning callback panics") }

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: These post-Offer checks pass even if the warning report moves after the callback panic is recovered. Assert inside OnWarn that the warning record already exists before panicking to protect the report-before-callback guarantee.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At internal/connector/queue_test.go, line 310:

<comment>These post-`Offer` checks pass even if the warning report moves after the callback panic is recovered. Assert inside `OnWarn` that the warning record already exists before panicking to protect the report-before-callback guarantee.</comment>

<file context>
@@ -299,3 +299,19 @@ func TestAPanickingPauseCallbackLeavesNoPhantomBacklog(t *testing.T) {
+	require.NoError(t, err)
+	var logs bytes.Buffer
+	queue.SetLogger(slog.New(slog.NewTextHandler(&logs, nil)))
+	queue.OnWarn = func(int) { panic("a warning callback panics") }
+
+	require.NoError(t, queue.Offer(context.Background(), 1))
</file context>
Suggested change
queue.OnWarn = func(int) { panic("a warning callback panics") }
queue.OnWarn = func(int) {
assert.Contains(t, logs.String(), "the backlog reached its warning depth")
panic("a warning callback panics")
}

// change it.
func (q *Queue) fire(edge queueEdge) {
log := q.logger()
q.report(log, edge)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P3: The recover arm in fire is now registered after q.logger() and q.report(log, edge) run, so a panic in the reporting path is no longer contained. Before this change, q.logger() was only called inside the already-recovered region, so nothing untrusted ran unguarded. report calls into the caller-supplied slog handler (intake adopts the production logger via adoptLogger), and a handler that panics would escape fire, propagate through deliver's drain into the Offer/Take caller, and take down the feed's delivery path — exactly what the panic containment in this file is for. The new comment on fire even claims "one that panics cannot take the report down with it," but it only protects against callback panics, not report panics. Arm the defer before the report so the report and callback both run contained.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At internal/connector/queue.go, line 311:

<comment>The recover arm in `fire` is now registered after `q.logger()` and `q.report(log, edge)` run, so a panic in the reporting path is no longer contained. Before this change, `q.logger()` was only called inside the already-recovered region, so nothing untrusted ran unguarded. `report` calls into the caller-supplied slog handler (intake adopts the production logger via `adoptLogger`), and a handler that panics would escape `fire`, propagate through `deliver`'s drain into the `Offer`/`Take` caller, and take down the feed's delivery path — exactly what the panic containment in this file is for. The new comment on `fire` even claims "one that panics cannot take the report down with it," but it only protects against callback panics, not report panics. Arm the `defer` before the report so the report and callback both run contained.</comment>

<file context>
@@ -299,9 +307,11 @@ func (q *Queue) deliver() {
 // change it.
 func (q *Queue) fire(edge queueEdge) {
+	log := q.logger()
+	q.report(log, edge)
 	defer func() {
 		if p := recover(); p != nil {
</file context>

A backlog hovering at the warning depth wrote a warning and a recovery
for every event, and a saturated queue wrote a pause and a resume for
every event the feed delivered. A warning now clears at half the warning
depth, and a pause is reported over only once no offer waits and the
backlog has drained to half the pause depth. The feed's blocking at
capacity is unchanged; only the reported edges move.

The pause line now reports what the queue holds, not counting the offer
waiting for room, and both warning lines point at the usual cause and at
basecamp connect status.
Copilot AI balanced review requested due to automatic review settings October 2, 2026 19:22

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Concurrent offers can cause the reported pause depth to exceed the queue capacity.

Review effort: Balanced
Findings: 2 Medium severity

Open (2)

case q.waiting == 0 && q.paused:
// The held depth: a waiting offer has counted its id but cannot put
// it in, so it is not part of what the queue holds.
q.pending = append(q.pending, queueEdge{kind: edgePause, fire: q.OnPause, depth: q.depth - q.waiting})

@cubic-dev-ai cubic-dev-ai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

1 existing issue remains and 3 new issues found across 4 files (changes from recent commits).

Prompt for AI agents (unresolved issues)

Check if these issues are valid — if so, understand the root cause of each and fix them. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. If appropriate, use sub-agents to investigate and fix each issue separately.


<file name="internal/connector/seam_contract_test.go">

<violation number="1" location="internal/connector/seam_contract_test.go:209">
P2: This test does not establish that the feed read is held back: it triggers the edge through `Queue.Offer`, which intake also uses for retry and requeue sweeps while live ingest can continue. Describe a generic blocked offer here, or test a source-specific live-ingest pause.</violation>
</file>

<file name="internal/connector/review2_test.go">

<violation number="1" location="internal/connector/review2_test.go:254">
P2: The remaining `Offer` goroutine can deliver the resume edge while this `Take` returns, making the immediate `resumes` assertion below flaky. Wait for `resumes.Load() == 1` with `require.Eventually` before asserting the count.</violation>
</file>

<file name="internal/connector/queue.go">

<violation number="1" location="internal/connector/queue.go:286">
P2: Concurrent offers can be counted by `q.stage(1)` before any registers in `q.waiting`, so this subtraction can report a pause depth above the channel capacity. Track the number of IDs actually held in the queue for this callback and log field.</violation>
</file>

Requires human review: Auto-approval blocked because this review re-detected 1 unresolved issue already reported by Cubic.

Re-trigger cubic


// Production builds its queue with no callbacks and hands it to intake with
// the connector's logger, so the logger is the only place an operator can
// learn that the backlog is growing or that the feed read is held back. Each

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: This test does not establish that the feed read is held back: it triggers the edge through Queue.Offer, which intake also uses for retry and requeue sweeps while live ingest can continue. Describe a generic blocked offer here, or test a source-specific live-ingest pause.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At internal/connector/seam_contract_test.go, line 209:

<comment>This test does not establish that the feed read is held back: it triggers the edge through `Queue.Offer`, which intake also uses for retry and requeue sweeps while live ingest can continue. Describe a generic blocked offer here, or test a source-specific live-ingest pause.</comment>

<file context>
@@ -206,8 +206,9 @@ func TestIntakeGivesTheQueueItsLogger(t *testing.T) {
 // the connector's logger, so the logger is the only place an operator can
-// learn that the backlog is growing or that the feed read has stopped. Each
-// edge is reported there, with the depth and the threshold it crossed.
+// learn that the backlog is growing or that the feed read is held back. Each
+// edge is reported there, with the depth, the threshold it crossed and the
+// band at which it clears.
</file context>
Suggested change
// learn that the backlog is growing or that the feed read is held back. Each
// learn that the backlog is growing or that an offer is waiting for room. Each

Comment on lines +254 to +255
_, err = queue.Take(ctx)
require.NoError(t, err)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: The remaining Offer goroutine can deliver the resume edge while this Take returns, making the immediate resumes assertion below flaky. Wait for resumes.Load() == 1 with require.Eventually before asserting the count.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At internal/connector/review2_test.go, line 254:

<comment>The remaining `Offer` goroutine can deliver the resume edge while this `Take` returns, making the immediate `resumes` assertion below flaky. Wait for `resumes.Load() == 1` with `require.Eventually` before asserting the count.</comment>

<file context>
@@ -244,11 +244,15 @@ func TestQueuePauseTracksEveryBlockedOffer(t *testing.T) {
 	require.Eventually(t, func() bool { return !queue.Paused() }, time.Second, time.Millisecond)
+	assert.Zero(t, resumes.Load(), "no resume until the backlog drains to its band")
+
+	_, err = queue.Take(ctx)
+	require.NoError(t, err)
 	assert.Equal(t, int32(1), pauses.Load())
</file context>
Suggested change
_, err = queue.Take(ctx)
require.NoError(t, err)
_, err = queue.Take(ctx)
require.NoError(t, err)
require.Eventually(t, func() bool { return resumes.Load() == 1 }, time.Second, time.Millisecond)

case q.waiting == 0 && q.paused:
// The held depth: a waiting offer has counted its id but cannot put
// it in, so it is not part of what the queue holds.
q.pending = append(q.pending, queueEdge{kind: edgePause, fire: q.OnPause, depth: q.depth - q.waiting})

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: Concurrent offers can be counted by q.stage(1) before any registers in q.waiting, so this subtraction can report a pause depth above the channel capacity. Track the number of IDs actually held in the queue for this callback and log field.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At internal/connector/queue.go, line 286:

<comment>Concurrent offers can be counted by `q.stage(1)` before any registers in `q.waiting`, so this subtraction can report a pause depth above the channel capacity. Track the number of IDs actually held in the queue for this callback and log field.</comment>

<file context>
@@ -250,14 +263,36 @@ func (q *Queue) stageWait(delta int) {
-	case q.waiting == 0 && q.paused:
+		// The held depth: a waiting offer has counted its id but cannot put
+		// it in, so it is not part of what the queue holds.
+		q.pending = append(q.pending, queueEdge{kind: edgePause, fire: q.OnPause, depth: q.depth - q.waiting})
+	}
+	if q.paused && q.waiting == 0 && q.depth <= q.resumeAt {
</file context>

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

tests Tests (unit and e2e)

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants