fix(world-postgres): ignore stream rows written after the first EOF - #3712
Conversation
🦋 Changeset detectedLatest commit: 434b214 The changes in this PR will be included in the next version bump. This PR includes changesets to release 1 package
Not sure what this means? Click here to learn what changesets are. Click here if you're a maintainer who wants to add another changeset to this PR |
|
@himself65 is attempting to deploy a commit to the Vercel Labs Team on Vercel. A member of the Team first needs to authorize it. |
karthikscale3
left a comment
There was a problem hiding this comment.
See inline comment.
| controller.enqueue(new Uint8Array(msg.data)); | ||
| } | ||
| if (msg.eof) { | ||
| closed = true; |
There was a problem hiding this comment.
Thanks for submitting a fix for this. One edge case remains: because the offset > 0 branch runs before this EOF handling, a positive startIndex can consume the first EOF as if it were a data chunk. With five data rows + EOF + retried data/EOF, get(..., 6) skips the first EOF and returns the post-EOF duplicate (or hangs if no later EOF arrives). Could we handle EOF before decrementing offset—or only apply the offset branch when !msg.eof—and add a regression test for a start index beyond the valid data count?
There was a problem hiding this comment.
Good catch, thanks. Fixed in d70e0e8 by exempting EOF rows from the offset branch (if (offset > 0 && !msg.eof)) — offsets count data chunks, matching getInfo's tailIndex, so the marker is never consumed. Added a regression test for a start index at (5) and past (6) the data count with a post-EOF duplicate and no trailing EOF; the 6 case hung until timeout before the fix.
karthikscale3
left a comment
There was a problem hiding this comment.
See inline comment.
| // index at or past the data count must still close the | ||
| // stream rather than consume the marker and then hang, or | ||
| // surface rows written after it. | ||
| if (offset > 0 && !msg.eof) { |
There was a problem hiding this comment.
One remaining edge case: a negative startIndex is calculated below from chunks.length, which includes rows after the first EOF. For A…E, EOF, duplicate-E, get(..., -1) computes an offset from seven rows, then reaches the first EOF without returning the final valid chunk. Could we derive dataCount from the first EOF (for example, with findIndex) and add a startIndex = -1 regression test?
There was a problem hiding this comment.
Thanks — fixed in 4319027: dataCount now comes from chunks.findIndex((c) => c.eof) (falling back to chunks.length with no EOF), so a negative start index resolves against the same rows enqueue delivers. Regression test added for -1 (→ e) and -2 (→ d, e) on A…E + EOF + duplicate-E; -1 returned nothing before the fix.
4319027 to
6a05ac4
Compare
|
@himself65 in order for us to merge this, commits must have verified signatures. Can you squash and re-push a signed commit? |
8d31d56 to
3e35119
Compare
`streams.get()`'s catch-up loop closed the controller on the first EOF row and kept iterating. A stream with rows after its first EOF — a producer that retried its terminal write after a lost ACK or an overlapping attempt, so the frame was appended and the stream closed again — then hit `controller.enqueue()` on a closed controller. Node threw ERR_INVALID_STATE, the stream errored, and every chunk still queued was discarded: a reader saw an empty, errored stream instead of the data written before the EOF. Track the first EOF and ignore everything written after it: - `enqueue()` now ignores rows once the stream has reached a terminal state, via the same reader-lifecycle `cleanedUp` flag main already uses to detach a completed reader's listener. - The EOF marker itself is exempt from the start-index offset, so a start index at or past the data count still closes the stream instead of consuming the marker and hanging. - A negative start index resolves against the rows before the first EOF, not every row in the table. - `getChunks()` and `getInfo()` read the same table and had the same gap: both counted and returned rows written after the first EOF. They're now bounded by it too, via a shared `findFirstEofChunkId` lookup. Co-authored-by: Peter Wielander <peter.wielander@vercel.com> Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com> Signed-off-by: Alex Yang <himself65@outlook.com> Signed-off-by: Peter Wielander <peter.wielander@vercel.com>
3e35119 to
209aa44
Compare
VaguelySerious
left a comment
There was a problem hiding this comment.
AI review: no blocking issues
AI Review: Note
Out of scope for this PR (different package, not in this diff), but packages/world-local/src/streamer.ts around line 643 has the same class of bug this PR fixes for world-postgres: its negative-startIndex resolution only excludes a single trailing EOF chunk file, not every row written after the first EOF. Nothing in world-local's write/close prevents writing after close, so a retried terminal write there would inflate dataChunkCount and resolve a negative startIndex against the wrong position. Worth a follow-up issue if that producer-retry scenario is possible for world-local consumers too.
| // that retries a terminal write can append data and EOF rows after it; | ||
| // `getChunks`/`getInfo` bound their queries by this so those rows are | ||
| // ignored the same way `streams.get()` ignores them. | ||
| const findFirstEofChunkId = async ( |
There was a problem hiding this comment.
AI Review: Note
getChunks() and getInfo() read the same workflow_stream_chunks table as streams.get() and had the same gap: both counted/returned rows written after the first EOF (a retried terminal write never errors these two, it just leaks a duplicate chunk into a page and inflates tailIndex by one per duplicate). Verified against the original PR revision with a real-Postgres integration test (getChunks() returned ["a","b","c","d","e","e"] instead of 5 chunks, getInfo().tailIndex was 5 instead of 4) before this fix. Added findFirstEofChunkId and bounded both queries by it, plus test/stream-post-eof-duplicates.test.ts covering both against a real Postgres container.
Signed-off-by: Peter Wielander <mittgfu@gmail.com>
…3712) Signed-off-by: Alex Yang <himself65@outlook.com> Signed-off-by: Peter Wielander <peter.wielander@vercel.com> Co-authored-by: Peter Wielander <peter.wielander@vercel.com> Co-authored-by: Peter Wielander <mittgfu@gmail.com> Signed-off-by: Alex Yang <himself65@outlook.com>
|
Backport PR opened against |
Description
streams.get()in@workflow/world-postgrescloses the controller on the EOF row during its catch-up loop and keeps iterating. If the stream has rows after its first EOF — a producer that retried its terminal write after a lost ACK or an overlapping attempt, so the frame was appended and the stream closed a second time — the next row hitscontroller.enqueue()on a closed controller. Node throwsERR_INVALID_STATE(Invalid state: Controller is already closed) out ofstart(), the stream errors, and every chunk still queued is discarded: the reader sees an empty, errored stream instead of the data written before the EOF.We hit this in production: a long-running agent run pushed its output through
writeToStream/closeStream; an ingest retry duplicated the terminal frame (data ×3, EOF ×3, data ×2, EOF ×2inworkflow_stream_chunks), and the workflow's finalize step — which reads the stream from 0 — persisted the run as failed with empty output even though all the data was there. On our instance 143 streams had rows after their EOF and 1809 had more than one EOF row. Reproduces on 4.3.3, 4.3.4 andmain.The fix tracks the first EOF and ignores every later row (data or EOF). It does not change behaviour for well-formed streams.
How did you test your changes?
packages/world-postgres/src/streamer.test.tsdrivescreateStreameragainst a fake pool/drizzle (thepgClientis stubbed so no LISTEN socket is opened). It writes five data rows, an EOF, then a duplicate data row + EOF, and expects the stream to drain exactly the five rows. Without the fix it rejects withTypeError: Invalid state: Controller is already closed.enqueueblows away. With a single queued chunk the consumer drains it first by microtask order and the failure does not reproduce.tsc --noEmitandbiome checkon the package are clean (Biome's pre-existingnoExcessiveCognitiveComplexitywarning onenqueuegoes from 16 to 20 with the extra guard;getChunksalready warns at 21).PR Checklist - Required to merge
pnpm changesetwas run to create a changelog for this PR (@workflow/world-postgres: patch)git commit --signoffon your commits)@vercel/workflowin a comment once the PR is ready, and the above checklist is complete