Conversation
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.
There was a problem hiding this comment.
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
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 rungh pr ready --undo.
Click "Ready for review" or rungh pr readyto 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.
| 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) |
There was a problem hiding this comment.
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
| 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", |
There was a problem hiding this comment.
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>
| 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") |
There was a problem hiding this comment.
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.)
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") } |
There was a problem hiding this comment.
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>
| 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) |
There was a problem hiding this comment.
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.
| 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}) |
There was a problem hiding this comment.
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 |
There was a problem hiding this comment.
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>
| // 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 |
| _, err = queue.Take(ctx) | ||
| require.NoError(t, err) |
There was a problem hiding this comment.
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>
| _, 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}) |
There was a problem hiding this comment.
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>

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
runConnectsets 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 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.Attributesdepth,warn_at,recover_at,pause_at.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. Attributesdepth(what the queue holds, not counting the waiting offer),pause_at,resume_at.the backlog drained to half its pause depth; the feed is read at full speed again. Attributesdepth,pause_at,resume_at.the backlog fell back below its warning depth. Attributesdepth,warn_at,recover_at.recover_at, 500 by default).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.Paused()still reports the live state.runConnectdoes not change.OnWarn,OnRecover,OnPause,OnResume) follow the same edges.OnPausegets the held depth.connect statusstill reads only the ledger.Tests
TestIntakeReportsTheBacklogsEdgesOnItsLoggerbuilds intake the way production does (a queue with no callbacks, the connector's logger inOptions). 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 onmain: nothing is logged.TestABacklogHoveringAtItsThresholdsIsReportedOnceEachWayoscillates 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 onmain.TestQueuePausesTheCallerAtTheThreshold,TestPauseAndResumeAreObservedInTheOrderTheyHappened,TestAPanickingPauseCallbackLeavesNoPhantomBacklog,TestQueuePauseTracksEveryBlockedOffer) now drain to the resume band before they expect a resume.-race -count=20.bin/ciis green.Summary by cubic
The backlog between intake and admission warned, paused, resumed, and recovered only through optional callbacks, and
runConnectset none, so in production nothing reported it. The queue now reports each edge on its logger before running any callback.New Features
connect statusstill reads only the ledger.Written for commit 54fb6b0. Summary will update on new commits.