#!/usr/bin/env python3 """engine.py -- THE ENGINE PROTOCOL: one running engine, its stream record, the SSE and NDJSON readers, the trap. What this file is, for a reader who opened it cold (SPEC Δ6, the bench-client lane §3.6). The ollama-vs-vLLM bench measures three engines with one client, one clock and one rule set: vLLM (engine_vllm.py), Ollama (engine_ollama.py) and llama-server (engine_llamaserver.py, llama.cpp's server as Ollama bundles it). Each adapter is an `Engine`: start() launch from servelines.argv_env, as a plain CHILD of the arm (no new session): the stop guard's signal to the arm's process group reaches it ready(timeout_s) seconds to ready; R-3 on an exit, a fatal log line or the timeout load() Ollama: one unscored request that loads the model (its load_duration kept); others: no-op readback() argv, env (whitelisted keys), the boot lines parsed, the internals (N-5), log offsets raw_body(S, ...) the parity path's body (item 1): never an `ignore_eos` key chat_row_body(t) the labelled chat row's body, with the engine's own thinking switch raw_path() / chat_path() decode_stream(resp, t0, rec) the frames -> the stream record (below), on the client's clock tokenize_raw(text, add_special) the IDs this engine evaluates for a raw string (R-P2) render_chat_ids(user_text) the chat render as IDs (R-P1) visible_count(text) the tokenizer count of visible text (the hidden-token gate's fallback) pieces(ids) each ID's piece, specials kept (the gate's exact method) server_clock(rec) / cached_tokens(rec) / prompt_count(rec) log_offset() / log_since(offset) the log bytes a level spans (all-layers, slot lines) stop() SIGINT -> SIGTERM -> SIGKILL, every descendant: seat.SeatProcess's own stop, inherited THE STREAM RECORD is the seat kit's `stream_once` record generalised; a field keeps its name where its meaning is unchanged (ttft_ms, t_first_token_s, t_finish_s, decode_tok_s, prompt_tokens, completion_tokens, finish_reason, completion_sha256, completion_head, token_frames, status, fail_kind, fail_reasons, unscored_*), and it adds generated_ids (when the engine sends them), prompt_ids_sha256 (parity.ids_sha of the prompt IDs the engine says it evaluated: R-P2's own sha, so the wire's IDs match the parity record of the same string), cached_tokens, server_clock{..., rule}, hidden_surplus / hidden_method, and path / engine / posture / nonce. ONE TIMING RULE for all three engines (the kit's Δ5 rule, seatlib.RULES["engine_timing"]): t_first_token_s is the first frame that carries a generated token, t_finish_s the frame carrying the last generated token, and decode_tok_s = (completion_tokens - 1) / (t_finish_s - t_first_token_s). Each adapter's `decode_stream` says which frame that is. """ import hashlib import http.client import json import os import re import signal import subprocess import time import pen import seat as S import seatlib as L import servelines as SL STREAM_TIMEOUT_S = 900.0 BARRIER_TIMEOUT_S = 120.0 CALL_TIMEOUT_S = 60.0 STREAM_TIMINGS = ("ttft_ms", "decode_tok_s", "t_first_token_s", "t_finish_s") UNSCORED_PREFIX = "unscored_" ENGINES = ("vllm", "ollama", "llama-server") _PY = re.compile(r"python[0-9.]*$") #: Every CPU this process could use when the kit was imported -- before the client pins itself (client.pin_client, #: N-12). An engine is launched on ALL of them: a child inherits its parent's affinity, and an engine squeezed onto the #: client's two SMT siblings would measure the pin, not the engine. ALL_CPUS = frozenset(os.sched_getaffinity(0)) if hasattr(os, "sched_getaffinity") else None def _release_pin(): """preexec_fn of every engine launch: the child runs on ALL_CPUS whatever the client pinned itself to.""" if ALL_CPUS: os.sched_setaffinity(0, ALL_CPUS) # ------------------------------------------------------------------------------------------------ # # processes # ------------------------------------------------------------------------------------------------ # def program_argv(argv): """An argv from the program token on: `python