Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
98 changes: 75 additions & 23 deletions c/openai_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -3116,6 +3116,35 @@ class JOBOBJECT_EXTENDED_LIMIT_INFORMATION(ctypes.Structure):
return None # never let process bookkeeping break starting the engine


def _write_all(stream, data, frame):
"""Write every byte of `data` to `stream`, looping on short writes.

The production engine stdin is a raw, unbuffered pipe (bufsize=0 ->
io.FileIO), whose write() is a single os.write() and may transfer fewer
bytes than it was given (a signal landing mid-write, a full pipe buffer
on a large IMAGE frame). Discarding the return value would leave the
tail of a frame unsent and desynchronize the engine's stdin framing, so
the remainder is re-offered until it is all consumed.

Neither `None` nor 0 is progress. `RawIOBase.write` answers `None` when
the stream is non-blocking and could not take a single byte, and 0 says
the same thing with a count; re-offering the buffer after either would
spin forever, so both fail closed as the named engine-write error a
broken pipe raises."""
written = 0
total = len(data)
view = memoryview(data)
while written < total:
sent = stream.write(view[written:])
# None is RawIOBase's "not one byte went out", not an uncounted
# full write, so it fails closed exactly as a zero count does.
if sent is None or sent <= 0:
raise RuntimeError(
f"failed to write {frame} to the engine "
f"(stdin took {written} of {total} bytes)")
written += sent


class Engine:
# cap=None = "not explicitly set": a glm-arch model's engine resolves the
# 0 sentinel (8 historically, 1 on Metal+darwin+fast SSD -- colibri.c
Expand Down Expand Up @@ -3192,6 +3221,32 @@ def _fail_pending(self, error):
for events in requests:
events.put(("error", error))

def _write_frame(self, request_id, data, frame):
"""Checked server->engine protocol write for CANCEL/STOP: the write
and its flush happen under one write_lock acquisition. Any failure
here -- an OSError from the pipe itself, or _write_all's own
fail-closed RuntimeError on a None/zero-progress write -- drops this
request's pending-map entry: the dispatcher only does that on this
id's own DONE/ERROR frame, and neither arrives when the write that
would have solicited one never reached the engine. An OSError is
additionally re-raised as a named RuntimeError rather than left as
itself: BrokenPipeError is a ConnectionError subclass, so an
unwrapped failure here would fall into do_POST's client-hangup
handler (`except ConnectionError: pass`) and the client would see a
silent connection close instead of the 500 engine_error the failure
actually is. _write_all's own RuntimeError is already the named
error this raises for an OSError, so it is re-raised as-is."""
try:
with self.write_lock:
_write_all(self.process.stdin, data, frame)
self.process.stdin.flush()
except Exception as error:
with self.pending_lock:
self.pending.pop(request_id, None)
if isinstance(error, OSError):
raise RuntimeError(f"failed to write {frame} to the engine ({error})") from error
raise

def _read_exact(self, size):
chunks = []
remaining = size
Expand Down Expand Up @@ -3412,14 +3467,21 @@ def decode_tool(data):
# annunciato subito prima del SUBMIT a cui appartengono. Deve
# partire dentro lo stesso lock, o un'altra richiesta potrebbe
# infilarsi in mezzo e prendersi l'immagine di questa.
if image is not None:
patches, grid_h, grid_w = image
blob = patches.tobytes() if hasattr(patches, "tobytes") else patches
self.process.stdin.write(
f"IMAGE {request_id} {len(blob)} {grid_h} {grid_w}\n".encode()
+ blob + b"\n")
self.process.stdin.write(header + payload + xpayload + b"\n")
self.process.stdin.flush()
try:
if image is not None:
patches, grid_h, grid_w = image
blob = patches.tobytes() if hasattr(patches, "tobytes") else patches
try:
_write_all(
self.process.stdin,
f"IMAGE {request_id} {len(blob)} {grid_h} {grid_w}\n".encode()
+ blob + b"\n", "IMAGE")
except OSError as error:
raise RuntimeError(f"failed to write IMAGE to the engine ({error})") from error
_write_all(self.process.stdin, header + payload + xpayload + b"\n", "SUBMIT")
self.process.stdin.flush()
except OSError as error:
raise RuntimeError(f"failed to write SUBMIT to the engine ({error})") from error
except Exception:
with self.pending_lock:
self.pending.pop(request_id, None)
Expand Down Expand Up @@ -3459,9 +3521,7 @@ def _accept(info):
# DONE frame; ClientCancelled is raised when it arrives.
if not cancel_sent and not stop_sent and cancelled and cancelled():
cancel_sent = True
with self.write_lock:
self.process.stdin.write(f"CANCEL {request_id}\n".encode())
self.process.stdin.flush()
self._write_frame(request_id, f"CANCEL {request_id}\n".encode(), "CANCEL")
continue
if kind == "accept":
if accepted:
Expand All @@ -3473,17 +3533,13 @@ def _accept(info):
decode(value)
if stopped and stopped():
stop_sent = True
with self.write_lock:
self.process.stdin.write(f"STOP {request_id}\n".encode())
self.process.stdin.flush()
self._write_frame(request_id, f"STOP {request_id}\n".encode(), "STOP")
elif cancelled and cancelled():
# Same admission-holding rule as the idle branch above:
# send CANCEL, then keep consuming frames until the
# engine acknowledges with ERROR CANCELLED or DONE.
cancel_sent = True
with self.write_lock:
self.process.stdin.write(f"CANCEL {request_id}\n".encode())
self.process.stdin.flush()
self._write_frame(request_id, f"CANCEL {request_id}\n".encode(), "CANCEL")
elif kind == "echo":
# Lettura del prefill: arriva PRIMA di ogni DATA e non e' testo
# generato, quindi non passa da decode() e non entra nella
Expand All @@ -3497,14 +3553,10 @@ def _accept(info):
decode_tool(value)
if stopped and stopped():
stop_sent = True
with self.write_lock:
self.process.stdin.write(f"STOP {request_id}\n".encode())
self.process.stdin.flush()
self._write_frame(request_id, f"STOP {request_id}\n".encode(), "STOP")
elif cancelled and cancelled():
cancel_sent = True
with self.write_lock:
self.process.stdin.write(f"CANCEL {request_id}\n".encode())
self.process.stdin.flush()
self._write_frame(request_id, f"CANCEL {request_id}\n".encode(), "CANCEL")
elif kind == "done":
_accept({"prompt_tokens": None})
if cancel_sent:
Expand Down
Loading
Loading