Shell reply sent raw on the ROUTER can drop the ZMQStream wake-up; a request arriving in that window is not read until the next connection

Abierto
#1,554 1 comentario 0 reacciones 0 asignados Ver en GitHub

Nadie ha tomado este issue todavía.

Evaluación

Dificultad
4/5
Tiempo estimado
3-5 días
Aptitud para principiantes
68/100
Tipo de issue
Error
Claridad
Bastante claro
Estado de actividad
Activo
Stack tecnológico
python

Línea de trabajo

Comienza en SubshellManager._send_on_shell_channel y en la configuración de shell-channel ZMQStream; después ejecuta zmq_lost_wakeup_det.py con --close-a y con --fix para reproducir y verificar la pérdida del wake-up. Se considera terminado cuando una solicitud que llega durante la respuesta raw ROUTER recibe respuesta sin necesitar una conexión posterior, preservando el comportamiento existente de shell-channel.

Escrito por el modelo de indexación a partir del texto del issue.

Descripción

Repo: ipykernel · Versions: ipykernel 7.2.0, pyzmq 27.1.0, libzmq 4.3.5, jupyter_server 2.21.0, macOS arm64 (also seen on Linux CI runners)

Summary

Since ipykernel 7, shell replies are written to the ROUTER socket with a raw
send_multipart from the shell channel thread
(SubshellManager._send_on_shell_channel), while shell requests are read
through a ZMQStream on that same socket. A raw send runs libzmq's
process_commands, which consumes the ROUTER's edge-triggered wake-up. pyzmq
re-reads ZMQ_EVENTS only after its own stream operations, so a request that
landed on the ROUTER while that reply was being sent is never noticed: it sits
in the socket until some later command on the socket wakes the stream (in
practice, the next client connection).

ipykernel 6 sent replies through the stream itself and did not have this window.

Why it shows up as "first message on a fresh connection hangs"

jupyter_server nudges every new channels websocket with a kernel_info_request
on a transient shell channel and closes that channel as soon as the control
channel answers. The kernel's shell reply to the nudge is therefore always sent
to a peer that has already gone (the ROUTER drops it silently), and a reply that
went nowhere produces no follow-up command. A client whose first real request
arrived during that send waits until another connection is opened. Symptoms:

  • a cell that never runs (kernel shows busy/idle never toggling),
  • a JupyterLab connection that hangs on first use and comes alive the moment a
    second tab connects to the same kernel.

Deterministic reproducer

zmq_lost_wakeup_det.py (below) rebuilds the shell-channel-thread shape —
one IOLoop, a ZMQStream on the ROUTER, a second ZMQStream on an inproc PAIR whose
callback writes the reply to the ROUTER raw — and forces the interleaving:
client B's request lands while A's reply callback is running.

$ python zmq_lost_wakeup_det.py --close-a --trials=3       # A closes like the nudge does
trial 0: B LOST (stuck in the ROUTER)
trial 1: B LOST (stuck in the ROUTER)
trial 2: B LOST (stuck in the ROUTER)
fix=False settle=0.2 trials=3 lost=3

$ python zmq_lost_wakeup_det.py --close-a --fix --trials=3  # re-read ZMQ_EVENTS after the raw send
trial 0: B answered
trial 1: B answered
trial 2: B answered
fix=True settle=0.2 trials=3 lost=0

Without --close-a the reply to A is delivered, A's peer eventually produces
another command, and B is rescued — which is why the bug looks intermittent in
the wild and always clears on the next connection.

(--settle must exceed libzmq's max_command_delay, ~1 ms on x86 TSC and
~125 ms on Apple Silicon; the default 0.2 s covers both.)

Suggested fix

After the raw send in SubshellManager._send_on_shell_channel, re-check the
shell stream's events, e.g. shell_stream._rebuild_io_state() — the same thing
pyzmq does after its own sends — or send the reply through the stream. We
currently carry the former as a wrapper installed before app.initialize()
and it removes the first-message hang in our CI.

I am happy to turn this into a PR if the maintainers prefer one approach over
the other.

zmq_lost_wakeup_det.py
"""Deterministic reproducer: ZMQStream wake-up lost to a raw send on the same socket.

Structure of ipykernel 7's shell channel thread: one IOLoop, a ZMQStream on the
ROUTER (requests in), and a second ZMQStream on an inproc PAIR (replies from the
main thread) whose callback writes the reply to the ROUTER with a RAW
``send_multipart``.

Interleaving forced here: while the PAIR callback is running (i.e. inside one
loop iteration, after the selector already took its readiness snapshot), client
B's request arrives at the ROUTER. libzmq queues an activate-read command and
signals the ROUTER's wake-up fd. The raw send then runs libzmq's
process_commands, which consumes that signal. Nothing re-reads ZMQ_EVENTS on the
ROUTER stream afterwards, so the loop never wakes for B: B's request sits in the
ROUTER until some later command on that socket wakes the stream again.

usage: python zmq_lost_wakeup_det.py [--fix] [--settle=SECONDS]
  --settle  time to wait inside the callback before the raw send, must exceed
            libzmq's max_command_delay (~1 ms on x86 TSC, ~125 ms on Apple Silicon)
"""
import select, sys, threading, time
import zmq
from tornado.ioloop import IOLoop
from zmq.eventloop.zmqstream import ZMQStream

FIX = "--fix" in sys.argv
SETTLE = float(next((a.split("=")[1] for a in sys.argv if a.startswith("--settle=")), 0.2))
CLOSE_A = "--close-a" in sys.argv
TRIALS = int(next((a.split("=")[1] for a in sys.argv if a.startswith("--trials=")), 5))

ctx = zmq.Context()
router = ctx.socket(zmq.ROUTER)
port = router.bind_to_random_port("tcp://127.0.0.1")
to_main_a, to_main_b = ctx.socket(zmq.PAIR), ctx.socket(zmq.PAIR)
to_main_a.bind("inproc://to-main"); to_main_b.connect("inproc://to-main")
from_main_a, from_main_b = ctx.socket(zmq.PAIR), ctx.socket(zmq.PAIR)
from_main_a.bind("inproc://from-main"); from_main_b.connect("inproc://from-main")

loop = IOLoop(make_current=False)
router_stream = ZMQStream(router, loop)
from_main_stream = ZMQStream(from_main_a, loop)

b_go = threading.Event()
b_sent = threading.Event()

def on_request(frames):
    to_main_a.send_multipart(frames)

def on_reply(frames):
    if frames[-1].startswith(b"A"):
        # A's reply is about to go out on the raw socket. Let B's request land first,
        # inside this same loop iteration, then wait out libzmq's send-side throttle.
        b_go.set(); b_sent.wait(); time.sleep(SETTLE)
        fd = router.fd
        print(f"    fd readable before raw send: {bool(select.select([fd], [], [], 0)[0])}")
    router.send_multipart(frames)          # raw send, as ipykernel 7 does
    if frames[-1].startswith(b"A"):
        print(f"    fd readable after raw send:  {bool(select.select([fd], [], [], 0)[0])}")
    if FIX:
        router_stream._rebuild_io_state()  # re-read ZMQ_EVENTS after touching the socket raw

router_stream.on_recv(on_request)
from_main_stream.on_recv(on_reply)
threading.Thread(target=loop.start, daemon=True, name="shell-channel").start()

def main_thread():
    while True:
        from_main_b.send_multipart(to_main_b.recv_multipart())
threading.Thread(target=main_thread, daemon=True, name="main-shell").start()

def dealer():
    s = ctx.socket(zmq.DEALER); s.linger = 0; s.connect(f"tcp://127.0.0.1:{port}"); return s

lost = 0
for trial in range(TRIALS):
    b_go.clear(); b_sent.clear()
    a, b = dealer(), dealer()
    time.sleep(0.05)                                   # both DEALERs connected
    a.send_multipart([b"", b"A-%d" % trial])
    if CLOSE_A:
        time.sleep(0.02); a.close()   # like jupyter_server's nudge: the transient channel is gone before the kernel replies
    b_go.wait()
    b.send_multipart([b"", b"B-%d" % trial]); time.sleep(0.02); b_sent.set()
    if not CLOSE_A:
        assert a.poll(2000), "A itself never answered"
        a.recv_multipart()
    if b.poll(1000):
        b.recv_multipart(); print(f"trial {trial}: B answered")
    else:
        lost += 1; print(f"trial {trial}: B LOST (stuck in the ROUTER)")
        c = dealer(); c.send_multipart([b"", b"C-%d" % trial])   # any new command wakes the stream
        print(f"    after a third client's request: B answered={bool(b.poll(1000))} C answered={bool(c.poll(1000))}")
        c.close()
    if not CLOSE_A: a.close()
    b.close()
print(f"fix={FIX} settle={SETTLE} trials={TRIALS} lost={lost}")
Lenguaje dominante
Python
Estrellas
734
Forks
411
Merge medio
1 d 2 h
PR fusionados (30 d)
9

Guía de contribución

Abrir la guía de contribución

Primeros pasos

  1. Lee el issue completo y luego la guía de contribución del proyecto.
  2. Comenta en el issue que vas a ocuparte — evita que dos personas hagan lo mismo.
  3. Haz un fork del repositorio y trabaja en una rama.
  4. Abre un pull request que haga referencia al número del issue.

Más de ipython/ipykernel

Todos los issues de ipython/ipykernel

Issues similares

Más issues de Python

Recibe los nuevos issues en tu correo

Un resumen breve de issues de GitHub para principiantes.