Repository navigation
fix(api): close the migration checks' source when the client goes away - #15618
ogabrielluiz wants to merge 2 commits into
Conversation
|
Important Review skippedAuto incremental reviews are disabled on this repository. Please check the settings in the CodeRabbit UI or the ⚙️ Run configuration
You can disable this status message by setting the Use the checkbox below for a quick retry:
No actionable comments were generated in the recent review. 🎉 ℹ️ Recent review info⚙️ Run configuration
📒 Files selected for processing (2)
Included review availability: This review used your included allowance. Your plan provides up to 10 included reviews per hour; 6 remain after this review. WalkthroughMigration source-check and database-backup responses now use ChangesMigration stream handling
Priority: ⬇️ Low Estimated code review effort: 2 (Simple) | ~10 minutes Change: Bug fix Merge Risk: ⚪ Minimal · up to No actionable merge-blocking issue was identified in the changed migration streams. 🚥 Pre-merge checks | ✅ 8 | ❌ 1❌ Failed checks (1 warning)✅ Passed checks (8 passed)✨ Finishing Touches 💡 2📝 Generate docstrings 💡
🛠️ Fix failing CI checks 💡
🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
✅ Test Coverage AdvisorNo source changes detected without accompanying tests. Thanks for keeping coverage up! 🎉
|
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## feat/migration-copy-destination #15618 +/- ##
===================================================================
+ Coverage 68.80% 68.96% +0.16%
===================================================================
Files 2770 2775 +5
Lines 298870 299572 +702
Branches 41158 40301 -857
===================================================================
+ Hits 205637 206604 +967
+ Misses 90739 90474 -265
Partials 2494 2494
Flags with carried forward coverage won't be shown. Click here to find out more.
🚀 New features to boost your workflow:
|
Cristhianzl
left a comment
There was a problem hiding this comment.
⚠️ Important (preferably this PR)
I1 — The late record write still clobbers a newer run, so the Jira's second symptom survives
_stream_checks' finally assigns the step unconditionally — record["steps"]["check_source"] = step at src/backend/base/langflow/api/v1/migration.py:762 — with no check that the record still belongs to this run. Before this PR the window was wide open: the abandoned stream was closed by the garbage collector at an arbitrary later time, which is the symptom LE-2965 describes as "when the abandoned stream was collected later it wrote cancelled over whatever run the record held by then". This PR moves that close inside the request, which removes the garbage-collector window, but it does not remove the unguarded write, and one path still reaches it late.
The surviving path needs a client that stops reading without closing the socket. The generator parks at yield _event(event) for the first event while await send(...) blocks on backpressure; migration-preflight writes the rest of its (small) output into the 64 KB stdout pipe and exits; _is_live at migration.py:217 then reads status == "running" with a dead pid and returns False, so the next POST /checks is admitted and writes its own check_source. When the stalled client finally goes away, this request's cleanup stamps the old run's status: "cancelled", old pid and report: None over the live one — the page shows the new run as cancelled with no report, which is the exact failure the description says the fix prevents.
The repo already has the idiom for this one function away: _settle_copies guards its write with if saved["steps"][step_id]["run_id"] == step["run_id"] at migration.py:652-654. check_source has no run_id, but started_at (or pid) identifies the run just as well, so the same compare-and-set closes it.
Code reference — src/backend/base/langflow/api/v1/migration.py:758-763
if step["status"] == "running":
step["status"] = "cancelled"
step["finished_at"] = _now()
record = _read_record()
record["steps"]["check_source"] = step
_write_record(record)Note on scope: the line itself comes from the base PR and this diff does not touch it. I am raising it anyway because the PR description and the Jira both present the overwrite as fixed by this change, and after reading the code it is narrowed rather than fixed. If you would rather keep this PR to the three lines it has, saying so in the description — and carrying the guard into 15621 — is a fine disposition; what I would not leave standing is the claim that the overwrite is gone.
93607bb to
0ee8f10
Compare
0ee8f10 to
2ece9bc
Compare
erichare
left a comment
There was a problem hiding this comment.
@ogabrielluiz Approving, with one small fix pushed for @Cristhianzl's I1. Thanks for fixing this by reusing the backup download's close-on-end response, now _ClosingStream, instead of adding a second mechanism, and for a test that holds the disconnect in the "event waiting in send" window without a sleep.
What I verified at 2ece9bcf24:
POST /checksnow returns_ClosingStream. In thefinallyof_stream_checks, the drain cancel, the kill, thecancelledstamp and the synchronous_saveall run before the first await. So a close from a cancelled scope cannot lose the record or leave the child running. No stale_Downloadreferences remain, and the backup download behaves the same.- The new test fails on the old code: with
run_checksreverted to a plainStreamingResponse, it fails withassert 'running' == 'cancelled'.
Fixes I pushed
- 85d19fc
fix(api): keep a late check cancel off a newer run's record. The late write still reached a newer run (I1), and the description says that overwrite is gone. A client that stops reading without closing the socket keeps its request open after its child exits, so_is_livelets the nextPOST /checksin. When the stalled client finally went away, its cleanup stamped the old run'scancelledover the new run:RUN2 record after it finished: done 25078 ...15:10:31 report=True FINAL record: cancelled 23128 ...15:09:59 report=Nonekeep()now saves only while the record'scheck_sourcestill has this run'sstarted_at, using the same compare-and-set that_settle_copiesuses withrun_id.test_a_run_that_ends_after_a_newer_run_took_the_record_leaves_it_alonedrives the response by hand, like your disconnect test. While the first event waits to be sent, another run takes the record and finishes, and then the client goes away. The test asserts that the newer step is untouched and the first child is gone. Without the guard, the record endscancelledinstead of the newerdone.
The PRs stacked above (#15621, #15622, #15623, #15624) will need a restack onto the new head.
Verification
test_migration.pyplustest_migration_copy_runs.py: 165 passed, 23 skipped (the skips need PostgreSQL or S3 servers).- Mutation check: run against the code before the guard, the new test fails as described above.
ruff checkandruff format --checkare clean on both changed files.
85d19fc to
19af5a7
Compare
6ee813a to
718072b
Compare
718072b to
28275ce
Compare
28275ce to
29e1d39
Compare
29e1d39 to
1a0886a
Compare
POST /migration/checks streams from a child process, and a client that goes away is meant to end the run: the child is stopped and the record says cancelled. That held only when the disconnect found the stream waiting for its child. When it found an event waiting to be sent, the response ended without closing its source. The child ran on and the record said running until the garbage collector closed the stream, which then wrote its cancelled run over whatever run the record held by that time. The backup download already had a response class that closes its source when the response ends, however it ends. It is now named for both uses and POST /checks returns it too.
A client that stops reading without closing keeps its POST /checks open after the child exits, so the next run is let in. When that client finally went away, its cancel was saved over the newer run's step. The save now compares started_at first, as _settle_copies compares run_id.
1a0886a to
ef34471
Compare
Stacked on the copy destination PR.
POST /api/v1/migration/checksstreams the source checks from a child process. A client that goes away is meant to end the run: the child is stopped and the record sayscancelled. That happened only when the disconnect found the run waiting for its child. When it found an event waiting to be sent, the response ended and left the run going.What that did. The record kept saying
runningand the child ran on to the end of its checks. For as long as the child lived, the nextPOST /checksanswered 409already_running. Later the garbage collector closed the abandoned stream, and the stream then wrote its own run into the record ascancelled, over whatever run the record held by that time. A check that had passed in between read as cancelled again, with no report.The fix. The backup download already returns a small response class that closes its source when the response ends, however it ends. This PR names that class for both uses (
_ClosingStream) and returns it fromPOST /checksas well, so the run ends inside its own request in both cases. Nothing else changes.Why it was rare. Starlette cancels a streamed response when its client disconnects. A cancellation that arrives while the stream waits inside its source is raised there, and the source cleans up. One that arrives while the stream waits to send never reaches the source. A check run spends nearly all of its time waiting for the child, so most disconnects land well. Each middleware layer of the app takes an event over before it passes it on, and for those moments the route is waiting to send.
test_a_second_run_waits_for_the_live_onehangs up just as the first event arrives, and it failed once in 33 runs of its file with the 409 above.The copy-run events stream.
GET /steps/{step_id}/runs/{run_id}/eventsis a plain streamed response too, and I left it as it is. A disconnect leaves its source open in the same way, but that source has nothing to clean up: it opens the run's log for each read and holds nothing in between, and a page that goes away is meant to stop nothing. Closing it inside the request would also take more than this class, because the route wraps the follower in a second generator and closing that one leaves the follower open.How I tested it
One new test in
test_migration.py. It takes the response that the route returns, with the realmigration-preflightchild behind it, and plays the client by hand: the first event is never taken and the client disconnects. When the response returns, the record on disk has to saycancelledwith an end time and the child has to be gone, with nothing awaited in between. Before the fix it failed 5 runs of 5 with the record sayingrunning.The test drives the response and leaves the app's middleware out. Through the whole app the same disconnect lands in either place as timing has it, and I found no way to hold it in the bad one without a sleep. That is also why counting runs of the old test cannot show the fix: on the old code it passed 30 runs of 30 alone.
Beyond the test:
runningand the child alive, and the nextPOST /checksanswered 409already_running. After the child had ended by itself a second run went through and readdone. Then I dropped the first response and ran the collector, and the record readcancelledwith no report. With the fix the record sayscancelledat the disconnect, the next run starts at once, and it still readsdoneafter the collector has run.test_migration.pyran 5 times in a row each way, with the drivers, a live PostgreSQL and an S3 server, and with no driver importable, as in the CI unit job. On this branch the five migration test files give 240 passed and 1 skipped against a local PostgreSQL and S3 server. The skip is the test that needs an S3 server that checks keys.ruff checkandruff format --checkare clean on the two files.Summary by CodeRabbit