Skip to content
Open
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
106 changes: 99 additions & 7 deletions scripts/laya-sidecar.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,18 +9,75 @@
Run: ~/venvs/laya/bin/python scripts/laya-sidecar.py [--port 8091]
Probe: curl -X POST 127.0.0.1:8091/v1/systemone -H 'Content-Type: application/json'
-d '{"state":"...","model":"laya","questions":{"q":{"type":"score","instructions":"..."}}}'

This is the SOURCE copy. The one macOS actually runs is the installed
`~/.local/share/laya-sidecar/laya-sidecar.py`, started at login by the
launchd job `ai.hermes.laya-sidecar` on port 8092 and shared by hermes,
x-trader, sale-loop and xdev. Do not start a second instance: one resident
checkpoint is ~2.4 GB on a 16 GB Mac, and two hot ones ~3.2 GB. When you
change the request-handling path here, copy it over the installed file and
`launchctl kickstart -k gui/$(id -u)/ai.hermes.laya-sidecar`, then check
`curl -s http://127.0.0.1:8092/health`.

Tunables the launchd plist owns, not this script: LAYA_PRELOAD (which
checkpoints to prefer), LAYA_MAX_LOADED (how many stay resident) and
LAYA_EAGER (load at boot instead of on first request).
"""

from __future__ import annotations

import argparse
import json
import os
import threading
import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer

from laya import Router

router = Router(preload=True)

def _preload_names() -> list[str]:
"""Checkpoints to keep resident, from LAYA_PRELOAD (comma-separated)."""
names = [n.strip() for n in os.environ.get("LAYA_PRELOAD", "english").split(",") if n.strip()]
return names or ["english"]


def build_router() -> Router:
"""Build the router lazy, so only the checkpoint that answers stays resident.

`Router(preload=True)` loads ALL THREE: measured phys_footprint 5.5 GB
(peak 6.1 GB) for callers that only ever hit English, on a 16 GB Mac that
already swaps. Lazy routing with max_loaded=1 measures 145 MB idle and
2.4 GB once one checkpoint is live. The cost is a cold load on the first
request after a restart — paid once, instead of at every boot.

ponytail: keeps one checkpoint, so genuinely multilingual traffic pays a
reload per language switch. Raise LAYA_MAX_LOADED (and accept the memory)
if that traffic ever becomes real.
"""
names = _preload_names()
max_loaded = max(1, int(os.environ.get("LAYA_MAX_LOADED", "1")))
router = Router(preload=False, max_loaded=max_loaded)
if os.environ.get("LAYA_EAGER", "").strip().lower() in ("1", "true", "yes"):
router.preload(names)
print(f"laya eager: resident={','.join(router.loaded)}", flush=True)
else:
print(f"laya lazy: max_loaded={max_loaded} preferred={names} "
f"(first request pays the cold load)", flush=True)
return router


router = build_router()

# Serialises predict across request threads. ThreadingHTTPServer runs each
# request on its own thread and every one of them hits the same global
# `router`, but a Metal command buffer holds a single encoder: two threads
# encoding at once trip "A command encoder is already encoding to this command
# buffer" and kill the process (measured on M2 Pro, torch 2.14: 8 concurrent
# requests died after 2 replies, 20 sequential ones were clean). Serialising
# costs the queue, not the GPU — a request is 130-240 ms. GET /health takes no
# lock, so health checks stay responsive while predictions queue.
_PREDICT_LOCK = threading.Lock()


class Handler(BaseHTTPRequestHandler):
Expand All @@ -39,7 +96,17 @@ def _send(self, status: int, payload: object) -> None:

def do_GET(self) -> None: # noqa: N802 — http.server names the method
if self.path == "/health":
self._send(200, {"ok": True, "model": "laya"})
# `router.loaded` is a load ORDER, not a statement about which head
# is usable: preloading a hub subfolder pulls the whole bundle, so
# it can list names that are resident but not the ones we rely on.
# Report the configured preload alongside it so a silent fallback
# away from the expected head is visible from outside the process.
self._send(200, {
"ok": True,
"model": "laya",
"loaded": router.loaded,
"preload": _preload_names(),
})
else:
self._send(404, {"error": "not found"})

Expand All @@ -56,7 +123,12 @@ def do_POST(self) -> None: # noqa: N802 — http.server names the method
except (ValueError, OSError):
self._send(400, {"error": "invalid JSON body"})
return
state = body.get("state", "")
# System One callers send the input as `text`; Laya's own API calls it
# `state`. Accept either so the contract and the model agree — sending
# only `text` left `state` empty and the head answered a blank string.
state = body.get("state")
if state is None:
state = body.get("text", "")
questions = body.get("questions", {})
if not isinstance(questions, dict) or not questions:
self._send(400, {"error": "questions must be a non-empty map"})
Expand All @@ -76,7 +148,8 @@ def do_POST(self) -> None: # noqa: N802 — http.server names the method
qdef["criteria"] = ["low", "medium", "high", "critical"]
started = time.perf_counter()
try:
out = router.predict(state, questions)
with _PREDICT_LOCK:
out = router.predict(state, questions)
except Exception as e: # noqa: BLE001 — a model failure is a 500, not a crash
self._send(500, {"error": f"laya predict failed: {e}"})
return
Expand All @@ -91,12 +164,31 @@ def do_POST(self) -> None: # noqa: N802 — http.server names the method
print(f"systemone {len(questions)}q {ms:.0f}ms", flush=True)


class Server(ThreadingHTTPServer):
"""ThreadingHTTPServer with a deep accept backlog.

socketserver's default `request_queue_size` is 5, and that is the LISTEN
backlog, not the worker pool: a burst larger than 5 has its extra
connections refused by the kernel before any worker thread starts.
Measured against the deployed sidecar on this M2 Pro box: 120 concurrent
requests dropped 40 at ~9 rps, while 60 concurrent lost none at ~35 rps.
Deepening the backlog turns that loss into the queue the server already
provides, and it stays bounded by `request_queue_size` — a saturated
backlog still sheds load instead of growing threads without limit.
"""

request_queue_size = 128


def main() -> None:
parser = argparse.ArgumentParser(description="Laya local sidecar for POST /v1/systemone")
parser.add_argument("--port", type=int, default=8091)
parser.add_argument("--port", type=int, default=int(os.environ.get("PORT", "8091")))
parser.add_argument("--host", default=os.environ.get("LAYA_HOST", "127.0.0.1"),
help="bind host; a container must use 0.0.0.0")
args = parser.parse_args()
server = ThreadingHTTPServer(("127.0.0.1", args.port), Handler)
print(f"laya-sidecar listening on 127.0.0.1:{args.port} (preload warm)", flush=True)
server = Server((args.host, args.port), Handler)
print(f"laya-sidecar listening on {args.host}:{args.port} "
f"(resident: {','.join(router.loaded) or 'none — lazy'})", flush=True)
server.serve_forever()


Expand Down
102 changes: 102 additions & 0 deletions scripts/test_laya_sidecar.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
"""Regression check for the sidecar's accept backlog.

The fix that matters here is four lines: `Server.request_queue_size = 128`,
because socketserver's default LISTEN backlog of 5 refuses the connections of
any burst larger than 5 before a worker thread ever starts. Two ways that can
rot unnoticed:

1. someone lowers the constant, or
2. someone rebuilds the server as a plain `ThreadingHTTPServer` -- the class
survives, still correct, and now nothing uses it.

(2) is the one a reader cannot see, so this asserts the wiring too. The stub
Router's sleep stands in for one serialised forward pass; a backlog bug only
shows up when the handler is slower than the accept loop.

Run: ~/venvs/laya/bin/python scripts/test_laya_sidecar.py
"""

from __future__ import annotations

import importlib.util
import pathlib
import socketserver
import sys
import time
import types
import unittest
from unittest import mock

SCRIPT = pathlib.Path(__file__).resolve().parent / "laya-sidecar.py"
EXPECTED_BACKLOG = 128


def _load_sidecar():
"""Import the sidecar as a real module with `laya` stubbed out.

The real `laya` package loads checkpoints; this check only needs the HTTP
layer, so a stub Router stands in and keeps the test dependency-free. The
module name is not `__main__`, so the script's own `main()` call does not
fire on import.
"""
fake = types.ModuleType("laya")

class StubRouter:
loaded: list[str] = []

def __init__(self, *_a, **_k) -> None:
pass

def predict(self, _state, questions):
time.sleep(0.05) # one serialised forward pass
return {
"routing": {"model": "stub"},
"answers": {
q: {"type": "choice", "choice": "a", "probabilities": {"a": 1.0}, "confidence": 1.0}
for q in questions
},
}

fake.Router = StubRouter
spec = importlib.util.spec_from_file_location("laya_sidecar_under_test", SCRIPT)
module = importlib.util.module_from_spec(spec)
with mock.patch.dict(sys.modules, {"laya": fake}):
spec.loader.exec_module(module)
return module


class AcceptBacklog(unittest.TestCase):
def test_backlog_is_deeper_than_the_stdlib_default(self):
sidecar = _load_sidecar()
self.assertEqual(
sidecar.Server.request_queue_size,
EXPECTED_BACKLOG,
"the sidecar must set its own request_queue_size on Server",
)
self.assertGreater(
sidecar.Server.request_queue_size,
socketserver.TCPServer.request_queue_size,
"a backlog at or below the socketserver default sheds a burst before any worker runs",
)

def test_main_builds_the_server_with_the_deep_backlog_class(self):
"""(2): the class existing is not the same as main using it."""
sidecar = _load_sidecar()
built = []

class Recorder(sidecar.Server):
def __init__(self, addr, handler):
built.append(addr)
# Never bind: this asserts the wiring, not the socket.
socketserver.BaseServer.__init__(self, addr, handler)

with mock.patch.object(sys, "argv", ["laya-sidecar.py", "--port", "0"]):
with mock.patch.object(sidecar, "Server", Recorder):
with mock.patch.object(socketserver.BaseServer, "serve_forever", lambda self: None):
sidecar.main()
self.assertEqual(len(built), 1, "main() must construct exactly one server")
self.assertEqual(built[0], ("127.0.0.1", 0))


if __name__ == "__main__":
unittest.main(verbosity=2)
Loading