Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
116 commits
Select commit Hold shift + click to select a range
74c95e5
Add Cua-S1 multimodal CUDA worker and parity recipe
Levius-Fubuki Sep 27, 2026
ee970c2
Merge main and align multimodal worker error contract
Levius-Fubuki Sep 28, 2026
767e211
Add reproducible multimodal profiling matrix
Levius-Fubuki Sep 28, 2026
ebd7adc
Distinguish profiler device annotations and support trace-only runs
Levius-Fubuki Sep 28, 2026
4c605e3
Record RTX 4090 profiling results and trace recovery
Levius-Fubuki Sep 28, 2026
159789a
docs: plan request-local image reuse experiment
Levius-Fubuki Sep 28, 2026
8c3aa20
test: add paired image reuse experiment and correctness checks
Levius-Fubuki Sep 28, 2026
7e7ed34
test: preserve unoptimized profiling baseline after reuse
Levius-Fubuki Sep 28, 2026
7874c6c
style: simplify profiling test environment stub
Levius-Fubuki Sep 28, 2026
cf899ff
Reuse image preprocessing and adapted vision features within requests
Levius-Fubuki Sep 28, 2026
c18a21b
test: cover multi-question JPEG reuse and document tensor contract
Levius-Fubuki Sep 28, 2026
4bb4035
Test image reuse cleanup after inference failures
Levius-Fubuki Sep 28, 2026
0abf31d
Avoid body-write race in transport rejection tests
Levius-Fubuki Sep 28, 2026
d0c4781
docs: link request reuse behavior and record completed reviews
Levius-Fubuki Sep 28, 2026
209a660
test: add reproducible result audit and CUDA HTTP postflight
Levius-Fubuki Sep 28, 2026
9f4a4fe
test: audit exact coverage of the paired workload matrix
Levius-Fubuki Sep 28, 2026
ccaa4fb
docs: record paired RTX 4090 image reuse results and exact parity
Levius-Fubuki Sep 28, 2026
95713a5
docs: record PR publication and verified experiment server shutdown
Levius-Fubuki Sep 28, 2026
224d88b
perf(cua-s1): project only final logits in reused path
Levius-Fubuki Sep 28, 2026
89e35e9
docs(cua-s1): publish RTX 4090 final-logits experiment
Levius-Fubuki Sep 28, 2026
c522e00
docs(cua-s1): record experiment publication and server shutdown
Levius-Fubuki Sep 28, 2026
a2abe7e
bench(cua-s1): measure multimodal CUDA Graph language forward
Levius-Fubuki Sep 28, 2026
809cacc
bench(cua-s1): validate graph replay with changed question text
Levius-Fubuki Sep 28, 2026
74b8c9e
docs(cua-s1): publish multimodal CUDA Graph evidence
Levius-Fubuki Sep 28, 2026
de00056
feat(cua-s1): add bounded segmented CUDA Graph runtime
Levius-Fubuki Sep 28, 2026
03c3acc
docs(cua-s1): publish segmented Graph runtime results
Levius-Fubuki Sep 28, 2026
9a818a7
docs(cua-s1): mark Graph runtime rollout complete
Levius-Fubuki Sep 28, 2026
1019581
Add mixed-shape CUDA Graph cost benchmark
Levius-Fubuki Sep 28, 2026
6e0a432
Parameterize Graph capture threshold for mixed workload comparison
Levius-Fubuki Sep 28, 2026
2336a08
Publish mixed-length CUDA Graph capture cost evidence
Levius-Fubuki Sep 28, 2026
af969d6
docs: plan request-aware Graph admission and validation
Levius-Fubuki Sep 28, 2026
7babb09
perf(cua-s1): gate Graph capture by request reuse and bounded work
Levius-Fubuki Sep 28, 2026
780ac4d
docs: explain request-aware Graph admission limits
Levius-Fubuki Sep 28, 2026
c07de3d
Add capture-inclusive admission benchmark and independent verifier
Levius-Fubuki Sep 28, 2026
383f9d6
test: add GPU input-change and admission fallback checks
Levius-Fubuki Sep 28, 2026
2969be1
Verify admission policy budgets and supported-path counters
Levius-Fubuki Sep 28, 2026
057736a
docs: record reproduced capture-stream lifetime hazard
Levius-Fubuki Sep 28, 2026
7f55152
Compare historical admission with shared exclusive-stream fix
Levius-Fubuki Sep 28, 2026
f285131
test: stress surviving graph replay across stream retirement
Levius-Fubuki Sep 28, 2026
0cdaad4
fix: give each CUDA graph an exclusively owned capture stream
Levius-Fubuki Sep 28, 2026
b094c92
style: satisfy admission test lint rules
Levius-Fubuki Sep 28, 2026
1b2b456
Share shape execution and pool accounting with legacy policy comparator
Levius-Fubuki Sep 29, 2026
97913fb
fix: share graph pools per shape and bound reserved memory
Levius-Fubuki Sep 29, 2026
dda1f0c
fix: require native allocator for bounded graph memory accounting
Levius-Fubuki Sep 29, 2026
9f2c7ef
test: reserve enough cache for distinct-layout replay coverage
Levius-Fubuki Sep 29, 2026
fb8a47d
docs: publish verified RTX 4090 admission and pool evidence
Levius-Fubuki Sep 29, 2026
5376cac
docs: record PR handoff and verified server shutdown
Levius-Fubuki Sep 29, 2026
be44f57
Fix shallow-checkout provenance test and image aspect validation
Levius-Fubuki Sep 29, 2026
3a3f235
docs: scope segmented graph bucketing experiment
Levius-Fubuki Sep 29, 2026
e8d73e8
experiment: bucket causal graph segments with strict per-length gates
Levius-Fubuki Sep 29, 2026
c810d29
experiment: locate padded segment numerical differences
Levius-Fubuki Sep 29, 2026
17e7044
experiment: preserve projection shapes and bucket DeltaNet rule only
Levius-Fubuki Sep 29, 2026
d2d416c
test: check changed bucket inputs and record capture failures
Levius-Fubuki Sep 29, 2026
9142e23
fix: accept cache-free rule metadata and verify experiment evidence
Levius-Fubuki Sep 29, 2026
d082d54
test: exercise bucket boundaries and cache fallback on GPU
Levius-Fubuki Sep 29, 2026
b9156ab
docs: record rule bucket results, tradeoffs and reproducible evidence
Levius-Fubuki Sep 29, 2026
d3bdc5a
docs: scope explicit rule-bucket worker integration
Levius-Fubuki Sep 29, 2026
d0fc83a
feat: integrate instance-local rule buckets into multimodal worker
Levius-Fubuki Sep 29, 2026
249b0e3
test: add worker HTTP checks and four-way bucket comparison
Levius-Fubuki Sep 29, 2026
74760be
fix: normalize rule bucket layouts at exact boundaries
Levius-Fubuki Sep 29, 2026
fd5c417
docs: record worker bucket parity and four-way GPU measurements
Levius-Fubuki Sep 29, 2026
71531bb
experiment(cua): compare tuned exact admission against rule buckets
Levius-Fubuki Sep 29, 2026
a67a506
experiment(cua): share graph cache and admission budgets across modes
Levius-Fubuki Sep 29, 2026
de1322a
test(cua): independently verify shared graph experiment evidence
Levius-Fubuki Sep 29, 2026
b33e690
experiment(cua): record controlled policy and shared-budget GPU results
Levius-Fubuki Sep 29, 2026
9051fb8
docs(cua): record final 241-test validation and evidence checksums
Levius-Fubuki Sep 29, 2026
0473de8
feat(cua): add opt-in automatic graph selection with shared budgets
Levius-Fubuki Sep 29, 2026
cb63dac
test(cua): exercise automatic worker full logits and lifecycle on CUDA
Levius-Fubuki Sep 29, 2026
e243a30
fix(cua): bound adaptive timing probes during eager fallback
Levius-Fubuki Sep 29, 2026
543fade
test(cua): record automatic worker parity and controlled GPU evidence
Levius-Fubuki Sep 29, 2026
5be2704
docs(cua): finalize 279-test validation and audited GPU report
Levius-Fubuki Sep 29, 2026
2180de0
docs(cua): record automatic worker PR delivery
Levius-Fubuki Sep 29, 2026
e6cf582
Merge main and fix multimodal validation review findings
Levius-Fubuki Sep 29, 2026
04acf99
Merge reviewed multimodal validation fixes into profiling
Levius-Fubuki Sep 29, 2026
642d14c
Merge reviewed validation and cleanup tests into image reuse
Levius-Fubuki Sep 29, 2026
068eb52
Sync reviewed validation fixes into post-reuse profiling
Levius-Fubuki Sep 29, 2026
8d20c44
Sync reviewed validation fixes into graph experiment
Levius-Fubuki Sep 29, 2026
8783e21
Sync reviewed protocol and deterministic cleanup tests into graph run…
Levius-Fubuki Sep 29, 2026
f9306db
Sync reviewed protocol and cleanup fixes into mixed-shape profiling
Levius-Fubuki Sep 29, 2026
b2d0fff
Reconcile reviewed validation and deterministic cleanup coverage for …
Levius-Fubuki Sep 29, 2026
8c0b470
Sync reviewed validation and cleanup regressions into bucket worker
Levius-Fubuki Sep 29, 2026
683d470
Sync reviewed validation and cleanup regressions into automatic graph…
Levius-Fubuki Sep 29, 2026
63edd9b
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
d608d39
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
88034eb
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
0ed9563
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
1e3df30
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
9f4626d
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
72f41b8
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
61597f6
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
0f0de52
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
22d6916
chore: remove local experiment artifacts from PR
Levius-Fubuki Sep 29, 2026
9093fd8
Move Cua-S1 HTTP serving out of models and reuse upstream weight mani…
Levius-Fubuki Sep 29, 2026
138c4d1
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
66d7cf5
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
70dad54
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
baaf32f
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
499db82
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
1f62815
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
a092759
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
9a7d6d3
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
e61ed17
Sync reviewed worker layout and weight manifest cleanup
Levius-Fubuki Sep 29, 2026
934e1ad
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
f298bbd
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
1a339e1
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
378d101
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
75aa051
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
e869ef5
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
ea5c19e
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
4862bd3
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
3254e90
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
164a345
chore: limit PR to multimodal runtime code
Levius-Fubuki Sep 30, 2026
552c02c
Fix graph capture ownership and retained pool accounting
Levius-Fubuki Sep 30, 2026
f4069e0
Integrate graph safety fixes from PR 22 before admission policy
Levius-Fubuki Sep 30, 2026
d598de3
Carry graph safety fixes into rule bucket worker
Levius-Fubuki Sep 30, 2026
10b1522
Carry graph safety fixes into automatic graph selection
Levius-Fubuki Sep 30, 2026
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
205 changes: 205 additions & 0 deletions src/frontend/cua_s1.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,205 @@
"""Small loopback HTTP worker; the Rust frontend remains the public serving layer."""

from __future__ import annotations

import argparse
import json
import logging
import socket
import threading
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer

from models.cua_s1.multimodal.protocol import (
MAX_BODY,
InvalidRequest,
MalformedJSON,
decode_request,
parse_request,
)

LOG = logging.getLogger(__name__)


class WorkerServer(ThreadingHTTPServer):
daemon_threads = False # Join accepted handlers before retiring GPU buffers.

def __init__(self, address, engine):
self.engine = engine
self.inference_lock = threading.Lock()
self._engine_closed = False
super().__init__(address, Handler)

def server_close(self):
super().server_close()
if not self._engine_closed:
close = getattr(self.engine, "close", None)
if close is not None:
close()
self._engine_closed = True


class Handler(BaseHTTPRequestHandler):
def setup(self):
super().setup()
self.connection.settimeout(15)

def log_message(self, format, *args):
# Do not log paths, input images, instructions or arbitrary request headers.
pass

def send_json(self, status, value):
raw = json.dumps(value, ensure_ascii=False, allow_nan=False).encode("utf-8")
self.send_response(status)
self.send_header("Content-Type", "application/json")
self.send_header("Content-Length", str(len(raw)))
self.end_headers()
try:
self.wfile.write(raw)
except (BrokenPipeError, ConnectionResetError):
pass

def do_GET(self):
if self.path == "/health":
self.send_json(200, {"status": "ready", "modality": "multimodal"})
else:
self.send_json(404, {"detail": "unknown route"})

def do_POST(self):
if self.path != "/v1/systemone":
self.send_json(404, {"detail": "unknown route"})
return
if self.headers.get("Transfer-Encoding"):
self.send_json(
411,
{
"detail": "Content-Length is required; chunked requests are unsupported"
},
)
return
lengths = self.headers.get_all("Content-Length", [])
if len(lengths) != 1:
self.send_json(411, {"detail": "one Content-Length is required"})
return
try:
length = int(lengths[0])
except ValueError:
self.send_json(400, {"detail": "invalid Content-Length"})
return
if length < 0 or length > MAX_BODY:
self.send_json(413, {"detail": "request exceeds body limit"})
return
if self.headers.get_content_type() != "application/json":
self.send_json(415, {"detail": "Content-Type must be application/json"})
return
if not self.server.inference_lock.acquire(blocking=False):
self.send_json(503, {"detail": "worker busy"})
return
try:
raw = self.rfile.read(length)
if len(raw) != length:
self.send_json(400, {"detail": "incomplete body"})
return
parsed = parse_request(decode_request(raw))
result = self.server.engine.predict(parsed)
self.send_json(200, result)
except MalformedJSON as exc:
self.send_json(400, {"detail": str(exc)})
except InvalidRequest as exc:
self.send_json(422, {"detail": str(exc)})
except (TimeoutError, socket.timeout):
self.send_json(408, {"detail": "request body timed out"})
except Exception as exc:
LOG.error("inference failed: %s", type(exc).__name__)
self.send_json(500, {"detail": "inference failed"})
finally:
self.server.inference_lock.release()


def parse_args(argv=None):
from models.cua_s1.multimodal.graph_runtime import GraphConfig

p = argparse.ArgumentParser(description=__doc__)
p.add_argument(
"--base", required=True, help="verified local base checkpoint directory"
)
p.add_argument(
"--adapter", required=True, help="verified local multimodal adapter directory"
)
p.add_argument("--port", type=int, default=8000)
p.add_argument(
"--graph", action="store_true", help="enable segmented CUDA Graph replay"
)
p.add_argument(
"--graph-mode",
choices=("exact", "rule-bucket", "auto"),
help="enable the selected Graph execution mode",
)
p.add_argument("--graph-bucket-width", type=int, default=None)
p.add_argument("--graph-max-shapes", type=int, default=8)
p.add_argument("--graph-max-memory-mib", type=int, default=1024)
p.add_argument(
"--graph-min-uses",
type=int,
default=2,
help="distinct requests needed before capture",
)
p.add_argument("--graph-max-tokens", type=int, default=2048)
p.add_argument("--graph-admission-window", type=int, default=8)
p.add_argument("--graph-cooldown-requests", type=int, default=32)
p.add_argument("--graph-capture-window", type=int, default=32)
p.add_argument("--graph-max-captures", type=int, default=4)
p.add_argument("--graph-capture-budget-ms", type=float, default=2000.0)
args = p.parse_args(argv)
mode = args.graph_mode or ("exact" if args.graph else None)
if args.graph_bucket_width is not None and mode not in {"rule-bucket", "auto"}:
p.error("--graph-bucket-width requires --graph-mode rule-bucket or auto")
try:
graph_config = (
GraphConfig(
mode=mode,
bucket_width=64
if args.graph_bucket_width is None
else args.graph_bucket_width,
max_shapes=args.graph_max_shapes,
max_bytes=args.graph_max_memory_mib * 1024 * 1024,
min_uses=args.graph_min_uses,
max_tokens=args.graph_max_tokens,
admission_window=args.graph_admission_window,
cooldown_requests=args.graph_cooldown_requests,
capture_window=args.graph_capture_window,
max_captures=args.graph_max_captures,
capture_budget_ms=args.graph_capture_budget_ms,
)
if mode is not None
else None
)
except ValueError as exc:
p.error(str(exc))
args.graph_config = graph_config
return args


def main():
from models.cua_s1.multimodal.model import MultimodalEngine

args = parse_args()
logging.basicConfig(level=logging.INFO)
engine = MultimodalEngine(args.base, args.adapter, graph_config=args.graph_config)
server = None
try:
engine.warmup()
# Bind only after model loading and representative inference succeed.
server = WorkerServer(("127.0.0.1", args.port), engine)
LOG.info("multimodal worker ready on 127.0.0.1:%s", args.port)
server.serve_forever()
except KeyboardInterrupt:
pass
finally:
if server is not None:
server.server_close()
else:
engine.close()


if __name__ == "__main__":
main()
71 changes: 71 additions & 0 deletions src/models/cua_s1/multimodal/graph_admission.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
"""Bounded request-frequency admission and sliding capture-work budget."""

from collections import OrderedDict, deque


class AdmissionPolicy:
"""Called under the runtime lock; counters advance once per predict request.

A completed attempt can exceed the time budget because capture is synchronous.
Every subsequent attempt waits for enough budget to expire. Cache hits never
spend capture budget. History is intentionally bounded and forgetting is cold.
"""

def __init__(self, config):
self.config = config
self.history = OrderedDict()
self.cooldowns = OrderedDict()
self.attempts = deque()
self.request_index = 0

def reset(self):
self.history.clear()
self.cooldowns.clear()
self.attempts.clear()
self.request_index = 0

def begin_request(self):
self.request_index += 1
while self.attempts and (
self.request_index - self.attempts[0][0] >= self.config.capture_window
):
self.attempts.popleft()

def reason(self, key):
"""Observe a cache miss and return its eager fallback reason, if any."""
if self.request_index <= self.cooldowns.get(key, -1):
return "cooldown"
self.cooldowns.pop(key, None)
last, count = self.history.get(key, (-self.config.admission_window, 0))
if self.request_index - last > self.config.admission_window:
count = 0
if last != self.request_index:
count += 1
self.history[key] = (self.request_index, count)
self.history.move_to_end(key)
while len(self.history) > 128:
self.history.popitem(last=False)
if count < self.config.min_uses:
return "warmup"
if (
len(self.attempts) >= self.config.max_captures
or sum(item[1] for item in self.attempts) >= self.config.capture_budget_ms
):
return "capture_budget"
return None

def start_capture(self):
ticket = [self.request_index, 0.0]
self.attempts.append(ticket)
return ticket

@staticmethod
def finish_capture(ticket, elapsed_ms):
ticket[1] = elapsed_ms

def evict(self, key):
self.history.pop(key, None)
self.cooldowns[key] = self.request_index + self.config.cooldown_requests
self.cooldowns.move_to_end(key)
while len(self.cooldowns) > 128:
self.cooldowns.popitem(last=False)
Loading
Loading