{"record":{"id":"b683bb8180e5348a","repo":"risingwavelabs/risingwave","slug":"receiving-commit-request-from-non-running-handle","errorCode":null,"errorMessage":"receiving commit request from non-running handle {}, running handles: {:?}","messagePattern":"receiving commit request from non-running handle (.+?), running handles: (.+?)","errorType":"exception","errorClass":"anyhow::Error","httpStatus":null,"severity":"error","filePath":"src/meta/src/manager/sink_coordination/coordinator_worker.rs","lineNumber":799,"sourceCode":"                                \"committing\"\n                            );\n                        })\n                        .await;\n\n                    match commit_res {\n                        Ok(_) => {\n                            two_phase_handler.ack_committed(epoch).await?;\n                        }\n                        Err(e) => {\n                            two_phase_handler.failed_committed(epoch, e);\n                        }\n                    }\n\n                    continue;\n                }\n            };\n            if !running_handles.contains(&handle_id) {\n                bail!(\n                    \"receiving commit request from non-running handle {}, running handles: {:?}\",\n                    handle_id,\n                    running_handles\n                );\n            }\n            pending_epochs.entry(epoch).or_default().add_new_request(\n                handle_id,\n                commit_request,\n                self.handle_manager.vnode_bitmap(handle_id),\n            )?;\n            if pending_epochs\n                .first_key_value()\n                .expect(\"non-empty\")\n                .1\n                .aligned()\n            {\n                let (epoch, commit_requests) = pending_epochs.pop_first().expect(\"non-empty\");\n                let mut metadatas = Vec::with_capacity(commit_requests.requests.len());","sourceCodeStart":781,"sourceCodeEnd":817,"githubUrl":"https://github.com/risingwavelabs/risingwave/blob/6469eb736d691e8e9b8a419a57edd6429ca77417/src/meta/src/manager/sink_coordination/coordinator_worker.rs#L781-L817","documentation":"The steady-state loop tracks which handles are currently running; a commit request arriving from a handle not in `running_handles` is rejected. This means a stopped/aborted writer is still sending commits, which would write data from a handle that was logically torn down or replaced (e.g. during a parallelism change).","triggerScenarios":"A writer whose Stop was processed by the coordinator sends a buffered commit; a new writer reuses the HandleId before being registered as running; a writer continues emitting commits after being stopped in `alter_parallelisms`.","commonSituations":"Parallelism decrease where the removed writer has in-flight buffered data; races between Stop acknowledgement and the writer's commit send; duplicated HandleIds after a writer reconnect.","solutions":["Ensure the writer stops sending commits immediately upon receiving StopCoordination and drains in-flight sends before exiting.","Verify the coordinator's Stop path marks handles non-running only after the writer confirmed shutdown, and that the writer confirms before finishing.","Check for stale HandleId reuse on reconnect — issue a fresh HandleId per incarnation.","Restart the affected sink to clear the mismatched state, then investigate the writer's shutdown ordering."],"exampleFix":null,"handlingStrategy":"validation","validationCode":"// writer-side precheck before sending a commit\nif self.stop_requested { return Err(anyhow!(\"handle stopped; refusing to commit\")); }","typeGuard":"fn is_running(running: &HashSet<HandleId>, id: &HandleId) -> bool { running.contains(id) }","tryCatchPattern":"// restart the affected sink to clear state, then audit shutdown ordering\nif err.contains(\"non-running handle\") { restart_sink(); inspect_writer_shutdown(); }","preventionTips":["Writers must stop emitting commits immediately on StopCoordination and drain in-flight sends first.","Issue a fresh HandleId per writer incarnation to avoid stale-id reuse.","Only mark handles non-running after the writer confirms shutdown."],"tags":["rust","sink-coordination","lifecycle","race-condition"],"backgroundTag":"invalid-state-transition","analyzedSha":"6469eb736d691e8e9b8a419a57edd6429ca77417","analyzedAt":"2026-09-11T21:06:21.487Z","contentChangedAt":"2026-09-11T21:06:21.487Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}