Skip to content

fix: preserve local in-process channels across peer DSM detachment - #108

Merged
stuhood merged 2 commits into
mainfrom
stuhood.task-multiplex-race
Oct 6, 2026
Merged

stuhood merged 2 commits into
mainfrom
stuhood.task-multiplex-race

Conversation

@stuhood

@stuhood stuhood commented Oct 5, 2026

Copy link
Copy Markdown

Problem

In multi-stage distributed queries where tasks are multiplexed across fewer worker processes than logical tasks (e.g. 8 tasks on 7 workers), peer workers can finish quickly, exit, and drop their DSM senders. This causes the remaining worker's DSM receiver to observe detachment (RecvBatchOutcome::Detached).

Previously, DrainHandle::detach_data_scope failed all registered channel buffers upon DSM detachment. When a worker hosted both a local producer and consumer fragment communicating via open_local_partition_sink, this prematurely failed the local in-process channel with:

Execution error: transport receiver detached before this channel's EOF; the producer went away

Solution

  • Track this_proc on DrainHandle to distinguish local in-process channels (sender_proc == this_proc) from remote DSM channels.
  • In detach_data_scope and register_data_channel, avoid failing local channels on DSM inbox detachment; local channels remain live until the local producer finishes.
  • Add failed: Option<String> to ChannelBufferRegistry so that fatal errors from fail_scope continue to fail both local and remote channels.
  • Implement Drop for LocalDrainPartitionSink to fail pending buffers if a local producer drops without calling finish() (e.g. due to panic or task error).
  • Wire local stream cancellation through MppMesh::cancel_stream and DrainHandle::cancel_stream.

Comment thread src/shm/transport.rs Outdated
Comment thread src/shm/transport.rs Outdated
Comment thread src/shm/transport.rs Outdated
Comment thread src/shm/transport.rs
Comment thread src/shm/transport.rs Outdated
Comment thread src/shm/transport.rs Outdated
Comment thread src/shm/transport.rs Outdated
Comment thread BUG_REPORT.md Outdated
@stuhood
stuhood force-pushed the stuhood.task-multiplex-race branch from ffe17f9 to 028ea1f Compare October 6, 2026 22:21
@stuhood
stuhood force-pushed the stuhood.task-multiplex-race branch from 028ea1f to 105b3ca Compare October 6, 2026 22:40
@stuhood
stuhood merged commit 2b82e8d into main Oct 6, 2026
37 checks passed
@stuhood
stuhood deleted the stuhood.task-multiplex-race branch October 6, 2026 22:48
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants