Skip to content

Commit 1be4153

Browse files
m-messerclaude
andauthored
Keep the worker serving after a handler raises (#13)
A handler exception permanently killed the worker. jsonrpc_handler passed the exception object as JSON-RPC `data`, which ujson cannot serialize, so building the error response raised TypeError out of dispatch(); the serve loop caught it and broke out, closing the client. The process stayed resident but never read stdin again, so every later request on that worker timed out in the shim's RPC send. Pass only the exception message, and split the serve loop so a failed dispatch no longer ends the session. Read and write failures still stop serving, since a partial frame leaves no safe point to resume from. Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
1 parent ae52fa6 commit 1be4153

2 files changed

Lines changed: 54 additions & 25 deletions

File tree

‎lf_toolkit/io/rpc_handler.py‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,9 @@ async def wrapped(req: dict):
3232
result = await handler.handle(name, {"params": req})
3333
return Success(result)
3434
except Exception as e:
35-
return Error(0, str(e), e)
35+
# Pass only the message: the exception object is not JSON
36+
# serializable, so sending it as `data` makes serializing the
37+
# error response raise, which tears down the serve loop.
38+
return Error(0, str(e))
3639

3740
return wrapped

‎lf_toolkit/io/stream_io.py‎

Lines changed: 50 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,6 @@
1+
import sys
2+
import traceback
3+
14
from abc import ABC
25
from abc import abstractmethod
36

@@ -82,29 +85,52 @@ def wrap_io(self, client: StreamIO) -> StreamIO:
8285
async def _handle_client(self, client: StreamIO):
8386
io = self.wrap_io(client)
8487

85-
while True:
86-
try:
87-
import sys
88-
print("waiting for data...", file=sys.stderr, flush=True)
89-
data = await io.read(4096)
90-
print(f"got data: {data[:80]}", file=sys.stderr, flush=True)
91-
92-
if not data:
93-
break
94-
95-
print("dispatching...", file=sys.stderr, flush=True)
96-
response = await self.dispatch(data.decode("utf-8"))
97-
print(f"got response: {str(response)[:80]}", file=sys.stderr, flush=True)
98-
99-
await io.write(response.encode("utf-8"))
100-
print("wrote response", file=sys.stderr, flush=True)
101-
except anyio.EndOfStream:
102-
break
103-
except anyio.ClosedResourceError:
104-
break
105-
except Exception as e:
106-
import traceback
107-
traceback.print_exc(file=sys.stderr)
108-
break
88+
serving = True
89+
while serving:
90+
serving = await self._serve_once(io)
10991

11092
await client.close()
93+
94+
async def _serve_once(self, io: StreamIO) -> bool:
95+
"""Read one request, dispatch it and write the response.
96+
97+
Returns True when the session can carry on, and False once the
98+
stream has ended or is no longer safe to read from.
99+
"""
100+
try:
101+
print("waiting for data...", file=sys.stderr, flush=True)
102+
data = await io.read(4096)
103+
print(f"got data: {data[:80]}", file=sys.stderr, flush=True)
104+
except (anyio.EndOfStream, anyio.ClosedResourceError):
105+
return False
106+
except Exception:
107+
# A read failure may have consumed part of a frame, so there is
108+
# no point in the stream we can safely resume from.
109+
traceback.print_exc(file=sys.stderr)
110+
return False
111+
112+
if not data:
113+
return False
114+
115+
try:
116+
print("dispatching...", file=sys.stderr, flush=True)
117+
response = await self.dispatch(data.decode("utf-8"))
118+
print(f"got response: {str(response)[:80]}", file=sys.stderr, flush=True)
119+
except Exception:
120+
# One bad request must not end the session: the frame was read in
121+
# full, so the stream is still aligned and the next request can be
122+
# served. The caller gets no reply for this one and will time out,
123+
# which beats every later request on this worker timing out too.
124+
traceback.print_exc(file=sys.stderr)
125+
return True
126+
127+
try:
128+
await io.write(response.encode("utf-8"))
129+
print("wrote response", file=sys.stderr, flush=True)
130+
except (anyio.EndOfStream, anyio.ClosedResourceError):
131+
return False
132+
except Exception:
133+
traceback.print_exc(file=sys.stderr)
134+
return False
135+
136+
return True

0 commit comments

Comments
 (0)