thegovind commited on
Commit
f7f0e34
·
verified ·
1 Parent(s): 400f1f6

Code revision v1.4: opt-in vLLM server (text only)

Browse files
Files changed (3) hide show
  1. README.md +34 -9
  2. VLLM.md +56 -0
  3. serve_vllm.py +682 -0
README.md CHANGED
@@ -29,8 +29,8 @@ Send a text or JSON `state` and typed questions: `choice` picks from up to 255 o
29
  |---|---|
30
  | Base model | [Qwen/Qwen3.5-4B](https://huggingface.co/Qwen/Qwen3.5-4B) (text model only; vision encoder and MTP head removed) |
31
  | Weights size | 8.4 GB (bf16, 4,205,751,296 parameters) |
32
- | Revision | v1.3 (code revision; weights identical to v1.0) |
33
- | License | Non-commercial research only ([LICENSE.md](LICENSE.md)); base model Apache-2.0 (`LICENSE-Qwen`) |
34
 
35
  ## Results
36
 
@@ -42,7 +42,7 @@ Send a text or JSON `state` and typed questions: `choice` picks from up to 255 o
42
  | JevK5 v0.2.0 | 76.1 | 79/111 (own runtime: 82/111) | 0.068 | 0.220 | 62.0 |
43
  | Jev 1.13.0 | — | — | — | — | 63.3 |
44
 
45
- The JevBench numbers are public-item development proxies, not official scores, and claim no rank or parity. Official scoring requires held-out, judge, and sealed items.
46
 
47
  <details><summary>How to read the public-item numbers</summary>
48
 
@@ -153,8 +153,8 @@ import os, sys
153
  from huggingface_hub import hf_hub_download
154
 
155
  os.environ["BLINK_MODEL"] = "thegovind/blink-4b"
156
- os.environ["BLINK_REVISION"] = "v1.3"
157
- sys.path.insert(0, os.path.dirname(hf_hub_download("thegovind/blink-4b", "blink.py", revision="v1.3")))
158
  import blink
159
 
160
  out = blink.decide(
@@ -178,7 +178,7 @@ print(out["answers"]["intent"]["probabilities"])
178
 
179
  ```sh
180
  pip install "torch==2.13.0" "transformers==5.17.0" "flash-linear-attention==0.5.2" "accelerate>=1.1.0" safetensors huggingface_hub
181
- hf download thegovind/blink-4b --revision v1.3 --local-dir blink-4b
182
  python blink-4b/serve.py --model ./blink-4b --port 8000
183
  # TypeSafe SDKs: export TYPESAFE_BASE_URL=http://127.0.0.1:8000 TYPESAFE_API_KEY=any
184
  ```
@@ -194,11 +194,36 @@ docker build -t blink-4b . && docker run --rm --gpus all -p 127.0.0.1:8000:8000
194
 
195
  To enable cross-request batching:
196
 
197
- ```sh
198
- hf download thegovind/blink-4b serve.py blink.py --revision v1.3 --local-dir blink-4b
199
  python blink-4b/serve.py --model ./blink-4b --port 8000 --batch-window-ms 5
200
  ```
201
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
202
  ## Screenshots (opt-in, self-hosted)
203
 
204
  Image input is off by default. Start `serve.py` with `--vision-tower Qwen/Qwen3.5-4B@851bf6e806efd8d0a36b00ddf55e13ccb7b8cd0a`, or set `BLINK_VISION_TOWER` to that value for direct `blink.py` calls. Cache the matching tower before serving a downloaded folder offline; image mode needs `torchvision==0.28.0`. The checkpoint's weights stay text-only.
@@ -301,7 +326,7 @@ These are source-repository licences; they do not settle rights in every underly
301
  - **Held-out and selection:** Candidate checkpoints and prompt format were selected on DI-S (3,000 requests). blink-4b was fixed using JevBench development proxies before its full-suite run. It scored 52.29 on the 129,422 requests outside DI-S. Public JevBench items (231 items) and JevK5's 65 hand-written hard items served as development selection sets; none were included in training. A one-time lockbox of 396 held-out authored items from domains unseen in any of the three checkpoints' training data (same generator families, not JevBench's sealed set) scored 0.861 accuracy and 0.026 ECE.
302
  - **Training overlap and audit:** Public train splits also used by the 0.1 index include ContractNLI, iSarcasmEval, VAST, Amazon ESCI, Humicroedit, and GSM8K train split. ANLI and BANKING77 train splits were also used. Public sources included 281 MMLU-Pro test-partition questions and 2 GPQA extended-set questions. An audit of all question rows against the 0.1 suite and JevBench public items found no content matches (all 36 flags were the BANKING77 template); no public JevBench items or chess positions matched. The 13-word passage check didn't search long-document bodies or option text. Semantic or pretraining overlap can't be ruled out, and private JevBench items weren't available to check.
303
  - **Temperature:** A split held out from T4 suggested T = 0.82, with negligible gain. But 166 of its 401 items were in T3 training, so it isn't held out from the released average. Temperature 1.0 is retained without fitting.
304
- - **Limits:** English-centric (Arabic task A scored 0.313, task C pairs 0.645 on 0.1; broader multilingual ability is unestablished); does not chat or explain answers; text in state can sway answers; date arithmetic and long policies are the weakest cases. Limits: 255 options per choice, 2–10 score levels, 131,072 input tokens per question (longest evaluated prompt: 37,906 tokens) and 512 questions per request; over-limit requests get HTTP 422 with the reason, never truncated.
305
 
306
  </details>
307
 
 
29
  |---|---|
30
  | Base model | [Qwen/Qwen3.5-4B](https://huggingface.co/Qwen/Qwen3.5-4B) (text model only; vision encoder and MTP head removed) |
31
  | Weights size | 8.4 GB (bf16, 4,205,751,296 parameters) |
32
+ | Revision | v1.4 (code revision; weights identical to v1.0) |
33
+ | License | Weights: non-commercial research and evaluation only ([LICENSE.md](LICENSE.md)); code: Apache-2.0. Base-model notice: Apache-2.0 (`LICENSE-Qwen`). |
34
 
35
  ## Results
36
 
 
42
  | JevK5 v0.2.0 | 76.1 | 79/111 (own runtime: 82/111) | 0.068 | 0.220 | 62.0 |
43
  | Jev 1.13.0 | — | — | — | — | 63.3 |
44
 
45
+ No official JevBench score for blink has been published. The JevBench numbers are public-item development proxies, not official scores, and claim no rank or parity. Official scoring requires held-out, judge, and sealed items.
46
 
47
  <details><summary>How to read the public-item numbers</summary>
48
 
 
153
  from huggingface_hub import hf_hub_download
154
 
155
  os.environ["BLINK_MODEL"] = "thegovind/blink-4b"
156
+ os.environ["BLINK_REVISION"] = "v1.4"
157
+ sys.path.insert(0, os.path.dirname(hf_hub_download("thegovind/blink-4b", "blink.py", revision="v1.4")))
158
  import blink
159
 
160
  out = blink.decide(
 
178
 
179
  ```sh
180
  pip install "torch==2.13.0" "transformers==5.17.0" "flash-linear-attention==0.5.2" "accelerate>=1.1.0" safetensors huggingface_hub
181
+ hf download thegovind/blink-4b --revision v1.4 --local-dir blink-4b
182
  python blink-4b/serve.py --model ./blink-4b --port 8000
183
  # TypeSafe SDKs: export TYPESAFE_BASE_URL=http://127.0.0.1:8000 TYPESAFE_API_KEY=any
184
  ```
 
194
 
195
  To enable cross-request batching:
196
 
197
+ ```sh
198
+ hf download thegovind/blink-4b serve.py blink.py --revision v1.4 --local-dir blink-4b
199
  python blink-4b/serve.py --model ./blink-4b --port 8000 --batch-window-ms 5
200
  ```
201
 
202
+ ## Higher throughput (opt-in)
203
+
204
+ `serve_vllm.py` (added in `v1.4`) is an opt-in, text-only server with higher throughput. `serve.py` stays the default.
205
+
206
+ Earlier paired loopback measurements (one self-hosted replica, 231 public items): c1 p50/p95 was 52/127 ms on vLLM vs 64/175 ms on `serve.py` (one request at a time); c16 was 34.7 vs 12.2 completed questions/s (16 concurrent requests).
207
+
208
+ Earlier peak memory was 68.0 GiB including cache and cold start 180 s; neither was retimed, nor were the other quality sets or c16 rerun. A release build re-passed the 231-item JevBench c1 quality check (80/111 hard, 199/231 total, hard ECE 0.069). Later changes touched only request checks, error handling and startup cleanup, not scoring.
209
+
210
+ Tested versions:
211
+
212
+ ```sh
213
+ python -m pip install "torch==2.13.0" "transformers==5.17.0" \
214
+ "vllm==0.30.0" "compressed-tensors==0.17.0" \
215
+ "accelerate>=1.1.0" safetensors huggingface_hub
216
+ python -m pip install "flash-linear-attention==0.5.2"
217
+ ```
218
+
219
+ ```sh
220
+ hf download thegovind/blink-4b --revision v1.4 --local-dir blink-4b
221
+ cd blink-4b
222
+ python serve_vllm.py --model . --port 8000 --quantization auto --max-concurrency 32 --max-num-seqs 32 --max-num-batched-tokens 8192 --gpu-memory-utilization 0.85
223
+ ```
224
+
225
+ [Setup, quality limits and caveats](https://huggingface.co/thegovind/blink-4b/blob/v1.4/VLLM.md). No official JevBench score for blink has been published.
226
+
227
  ## Screenshots (opt-in, self-hosted)
228
 
229
  Image input is off by default. Start `serve.py` with `--vision-tower Qwen/Qwen3.5-4B@851bf6e806efd8d0a36b00ddf55e13ccb7b8cd0a`, or set `BLINK_VISION_TOWER` to that value for direct `blink.py` calls. Cache the matching tower before serving a downloaded folder offline; image mode needs `torchvision==0.28.0`. The checkpoint's weights stay text-only.
 
326
  - **Held-out and selection:** Candidate checkpoints and prompt format were selected on DI-S (3,000 requests). blink-4b was fixed using JevBench development proxies before its full-suite run. It scored 52.29 on the 129,422 requests outside DI-S. Public JevBench items (231 items) and JevK5's 65 hand-written hard items served as development selection sets; none were included in training. A one-time lockbox of 396 held-out authored items from domains unseen in any of the three checkpoints' training data (same generator families, not JevBench's sealed set) scored 0.861 accuracy and 0.026 ECE.
327
  - **Training overlap and audit:** Public train splits also used by the 0.1 index include ContractNLI, iSarcasmEval, VAST, Amazon ESCI, Humicroedit, and GSM8K train split. ANLI and BANKING77 train splits were also used. Public sources included 281 MMLU-Pro test-partition questions and 2 GPQA extended-set questions. An audit of all question rows against the 0.1 suite and JevBench public items found no content matches (all 36 flags were the BANKING77 template); no public JevBench items or chess positions matched. The 13-word passage check didn't search long-document bodies or option text. Semantic or pretraining overlap can't be ruled out, and private JevBench items weren't available to check.
328
  - **Temperature:** A split held out from T4 suggested T = 0.82, with negligible gain. But 166 of its 401 items were in T3 training, so it isn't held out from the released average. Temperature 1.0 is retained without fitting.
329
+ - **Limits:** English-centric (Arabic task A scored 0.313, task C pairs 0.645 on 0.1; broader multilingual ability is unestablished); does not chat or explain answers; text in state can sway answers; date arithmetic and long policies are the weakest cases. For `blink.py` and default `serve.py`: 255 options per choice, 2–10 score levels, 131,072 input tokens per question (longest evaluated prompt: 37,906 tokens) and 512 questions per request; over-limit requests get HTTP 422 with the reason, never truncated.
330
 
331
  </details>
332
 
VLLM.md ADDED
@@ -0,0 +1,56 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ # vLLM serving (opt-in)
2
+
3
+ `serve_vllm.py` serves blink-4b text-only through `/v1/systemone`, `/v1/models` and `/healthz`. `serve.py` stays the default; it can take screenshots when enabled.
4
+
5
+ From the model folder at revision `v1.4`:
6
+
7
+ Tested versions:
8
+
9
+ ```sh
10
+ python -m pip install "torch==2.13.0" "transformers==5.17.0" \
11
+ "vllm==0.30.0" "compressed-tensors==0.17.0" \
12
+ "accelerate>=1.1.0" safetensors huggingface_hub
13
+ python -m pip install "flash-linear-attention==0.5.2"
14
+ ```
15
+
16
+ ```sh
17
+ python serve_vllm.py --model . --port 8000 \
18
+ --quantization auto --max-concurrency 32 --max-num-seqs 32 \
19
+ --max-num-batched-tokens 8192 --gpu-memory-utilization 0.85
20
+ ```
21
+
22
+ Tested: vLLM 0.30.0; Transformers 5.17.0; Torch 2.13.0; compressed-tensors 0.17.0; flash-linear-attention 0.5.2. BF16 weights, FP32 offered-label head; chunked prefill on, prefix caching off.
23
+
24
+ Tested vLLM context: `--max-model-len 32768`; on this server, a question over 32,768 tokens returns `422`. Default `serve.py` allows 131,072 tokens per question; another vLLM context size needs a new quality check.
25
+
26
+ ## Quality
27
+
28
+ The tables are from the earlier E3 study; capacity was not retimed and the other quality sets and c16 were not rerun. A release build re-passed the 231-item JevBench c1 quality check (80/111 hard, 199/231 total, hard ECE 0.069). Later changes touched only request checks, error handling and startup cleanup, not scoring.
29
+
30
+ | Task measure | vLLM result | Acceptance |
31
+ | --- | ---: | ---: |
32
+ | JevBench public hard correct | 80/111 | at least 78/111 |
33
+ | JevBench public total correct | 199/231 | at least 196/231 |
34
+ | JevBench public hard top-label ECE (rounded) | c1 0.069; c16 first 0.069 / repeat 0.074 | at most ~0.077 on each |
35
+ | Web actions, 5 options (c1): agreement with the same model's FP32 answer | 499/500 | at least 489/500 |
36
+ | Web actions, 9 options (c1): agreement with the same model's FP32 answer | 494/500 | at least 487/500 |
37
+ | Decision Index 0.1 latency-tail sample (160 requests, 1,851 questions): agreement with the same model's FP32 answer | c1 1846/1851; c16 first 1847/1851 / repeat 1849/1851 | at least 1836/1851 on each |
38
+
39
+ Both web-action sets have 500 text-only questions each from public Multimodal-Mind2Web; ECE is rounded here, but checked unrounded. The Decision Index 0.1 latency-tail sample is 160 length-selected requests (1,851 questions) from reference-covered u1000, excluding source mismatches; it is not full u1000, DI-S or 0.2.
40
+
41
+ ## Capacity
42
+
43
+ | Measure | `serve.py` control | `serve_vllm.py` |
44
+ | --- | ---: | ---: |
45
+ | JevBench c1 p50 / p95 | 64 / 175 ms | 52 / 127 ms |
46
+ | Longest TypeSafe documents (16 requests), c1 p95 | 22.2 s | 16.4 s |
47
+ | JevBench c16 completed q/s | 12.2 | 34.7 |
48
+ | Decision Index 0.1 latency-tail sample (160 requests, 1,851 questions), c16 completed q/s | 21.0 | 46.5 |
49
+ | Peak device memory while scoring | 21.4 GiB | 68.0 GiB |
50
+ | Startup through health and warm score | first not captured; later control under 19 s | first cold start 180 s |
51
+
52
+ One self-hosted replica, loopback HTTP. c1 uses one request at a time; c16 uses 16 concurrent requests and reports completed questions/s. Peak memory includes reserved cache. Quality was checked on these samples, not on every Decision Index item. Agreement with the same model's FP32 answer is not gold accuracy or bit-for-bit parity. For more traffic, use separate warm replicas; recheck quality after changing flags or runtime. No official JevBench score for blink has been published.
53
+
54
+ Any request with a top-level `images` field (even `[]`) or inline image data returns `422` (`"this model reads text only"`); the server reports `accepts_images: false`. `--quantization auto` keeps these BF16 weights. The INT8 builds that were tried did not pass the quality checks.
55
+
56
+ Code: Apache-2.0. Weights: non-commercial research and evaluation only; see `LICENSE.md`.
serve_vllm.py ADDED
@@ -0,0 +1,682 @@
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
+ """Opt-in, text-only vLLM server for the TypeSafe System One API.
2
+
3
+ Install vLLM 0.30.0 and the model's runtime dependencies. For the qualified
4
+ bf16 blink-4b checkpoint, run the V32 configuration:
5
+
6
+ python serve_vllm.py --model ./blink-4b --host 127.0.0.1 --port 8000 \\
7
+ --quantization none --max-concurrency 32 --max-num-seqs 32 \\
8
+ --max-num-batched-tokens 8192 --max-model-len 32768 \\
9
+ --tensor-parallel-size 1 --gpu-memory-utilization 0.85
10
+
11
+ Chunked prefill is on and prefix caching is off. Images receive a located 422.
12
+ Each question's prompt plus its one label token must fit --max-model-len
13
+ (default 32768); the default serve.py may use a larger context. Client sockets
14
+ time out after 60 seconds, with at most 2*max-concurrency+16 admitted connections.
15
+ Use a buffering reverse proxy for public ingress.
16
+ The released serve.py remains the default. This server uses the model folder's
17
+ blink.py to validate, render and assemble every typed decision. It asks vLLM
18
+ for processed logprobs over the offered single-token letters, with an FP32
19
+ offered-label head, and never returns a generated text completion.
20
+ """
21
+
22
+ from __future__ import annotations
23
+
24
+ import argparse
25
+ import asyncio
26
+ import concurrent.futures
27
+ import hashlib
28
+ import hmac
29
+ import importlib.metadata
30
+ import importlib.util
31
+ import json
32
+ import math
33
+ import os
34
+ import re
35
+ import signal
36
+ import socket
37
+ import sys
38
+ import threading
39
+ import uuid
40
+ from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
41
+ from pathlib import Path
42
+
43
+ WARMUP = (
44
+ "Order 4471 arrived with a cracked screen. The customer wants a replacement.",
45
+ {
46
+ "route": {
47
+ "type": "choice",
48
+ "instructions": "Which team should handle this?",
49
+ "criteria": {"returns": "Damaged items", "billing": "Charges", "shipping": "Late parcels"},
50
+ },
51
+ "urgent": {"type": "noul", "instructions": "Does this need a reply today?"},
52
+ },
53
+ )
54
+ MAX_BODY_BYTES = 20 * 1024 * 1024
55
+
56
+
57
+ def sha256(path: Path) -> str:
58
+ digest = hashlib.sha256()
59
+ with path.open("rb") as stream:
60
+ for block in iter(lambda: stream.read(1 << 24), b""):
61
+ digest.update(block)
62
+ return digest.hexdigest()
63
+
64
+
65
+ def verify_weights(root: Path) -> bool | None:
66
+ manifest = root / "weights.sha256"
67
+ if not manifest.is_file():
68
+ return None
69
+ checked = 0
70
+ for line in manifest.read_text(encoding="utf-8").splitlines():
71
+ if not line.strip():
72
+ continue
73
+ digest, *names = line.split()
74
+ if len(digest) != 64 or any(c not in "0123456789abcdef" for c in digest) or len(names) != 1:
75
+ raise ValueError("invalid weights.sha256 entry")
76
+ name = Path(names[0])
77
+ if name.is_absolute() or name.name != names[0] or not (root / name).is_file():
78
+ raise ValueError(f"invalid or missing weight manifest file: {name}")
79
+ if sha256(root / name) != digest:
80
+ raise ValueError(f"weight checksum mismatch: {name}")
81
+ checked += 1
82
+ if not checked:
83
+ raise ValueError("empty weights.sha256")
84
+ return True
85
+
86
+
87
+ def model_config(root: Path) -> dict:
88
+ config = json.loads((root / "config.json").read_text(encoding="utf-8"))
89
+ if not isinstance(config, dict):
90
+ raise TypeError("model config must be a JSON object")
91
+ return config
92
+
93
+
94
+ def detect_quantization(root: Path, requested: str) -> str | None:
95
+ config = model_config(root)
96
+ text_config = config.get("text_config") or {}
97
+ if not isinstance(text_config, dict):
98
+ raise TypeError("text model config must be a JSON object")
99
+ settings = config.get("quantization_config") or text_config.get("quantization_config")
100
+ if settings is not None and (not isinstance(settings, dict)
101
+ or settings.get("quant_method") != "compressed-tensors"):
102
+ raise ValueError("unsupported checkpoint quantization")
103
+ actual = "compressed-tensors" if settings else "none"
104
+ if requested != "auto" and requested != actual:
105
+ raise ValueError(f"requested {requested} but checkpoint quantization is {actual}")
106
+ return None if actual == "none" else actual
107
+
108
+
109
+ def has_vision_model(root: Path) -> bool:
110
+ config = model_config(root)
111
+ return config.get("model_type") == "qwen3_5" or config.get("vision_config") is not None
112
+
113
+
114
+ def load_blink(root: Path):
115
+ path = root / "blink.py"
116
+ if not path.is_file():
117
+ raise FileNotFoundError(f"model folder has no renderer: {path}")
118
+ spec = importlib.util.spec_from_file_location("blink_vllm_model", path)
119
+ if spec is None or spec.loader is None:
120
+ raise RuntimeError(f"cannot import renderer: {path}")
121
+ blink = importlib.util.module_from_spec(spec)
122
+ sys.modules[spec.name] = blink
123
+ try:
124
+ spec.loader.exec_module(blink)
125
+ except BaseException:
126
+ del sys.modules[spec.name]
127
+ raise
128
+ return blink
129
+
130
+
131
+ def bearer(value: str | None) -> str | None:
132
+ scheme, _, token = (value or "").strip().partition(" ")
133
+ return token.strip() or None if scheme.lower() == "bearer" else None
134
+
135
+
136
+ def api_key(value: str) -> str | None:
137
+ key = value.strip()
138
+ if not key:
139
+ return None
140
+ if not key.isascii() or not key.isprintable() or any(char.isspace() for char in key):
141
+ raise argparse.ArgumentTypeError("API key must be printable ASCII without whitespace")
142
+ return key
143
+
144
+
145
+ def int_range(minimum: int, maximum: int):
146
+ def parse(value: str) -> int:
147
+ try:
148
+ number = int(value)
149
+ except ValueError:
150
+ raise argparse.ArgumentTypeError(f"not an integer: {value!r}") from None
151
+ if not minimum <= number <= maximum:
152
+ raise argparse.ArgumentTypeError(f"must be {minimum}-{maximum}")
153
+ return number
154
+
155
+ return parse
156
+
157
+
158
+ def fraction(value: str) -> float:
159
+ try:
160
+ number = float(value)
161
+ except ValueError:
162
+ raise argparse.ArgumentTypeError(f"not a number: {value!r}") from None
163
+ if not 0 < number < 1:
164
+ raise argparse.ArgumentTypeError("must be strictly between 0 and 1")
165
+ return number
166
+
167
+
168
+ def error_loc(blink, request: dict, message: str) -> list:
169
+ questions = request.get("questions")
170
+ if isinstance(questions, dict) and 0 < len(questions) <= blink.MAX_QUESTIONS:
171
+ for key in questions:
172
+ if message.startswith(f"question {key!r} "):
173
+ return ["body", "questions", key]
174
+ for key, question in questions.items():
175
+ try:
176
+ blink.question_options(question)
177
+ except blink.BlinkError:
178
+ return ["body", "questions", key]
179
+ return ["body", "questions"]
180
+
181
+
182
+ def require_image_contract(blink) -> None:
183
+ missing = [name for name in ("contains_image_uri", "inspect_images")
184
+ if not callable(getattr(blink, name, None))]
185
+ blink_error = getattr(blink, "BlinkError", None)
186
+ if not isinstance(blink_error, type):
187
+ missing.append("BlinkError")
188
+ image_error = getattr(blink, "ImageError", None)
189
+ if (not isinstance(image_error, type) or not isinstance(blink_error, type)
190
+ or not issubclass(image_error, blink_error)):
191
+ missing.append("ImageError")
192
+ if missing:
193
+ raise RuntimeError(f"blink.py must provide the v1.3 image scanner: {', '.join(missing)}")
194
+
195
+
196
+ class EngineOwner:
197
+ def __init__(self):
198
+ self.lock = threading.Lock()
199
+ self.engine = None
200
+ self.abandoned = False
201
+
202
+ def adopt(self, engine) -> None:
203
+ with self.lock:
204
+ if not self.abandoned:
205
+ self.engine = engine
206
+ return
207
+ engine.shutdown()
208
+ raise RuntimeError("vLLM startup was interrupted")
209
+
210
+ def abandon(self) -> None:
211
+ with self.lock:
212
+ self.abandoned = True
213
+
214
+ def close(self) -> None:
215
+ with self.lock:
216
+ self.abandoned = True
217
+ engine, self.engine = self.engine, None
218
+ if engine is not None:
219
+ engine.shutdown()
220
+
221
+
222
+ def submit_initializer(loop, initialize):
223
+ result = concurrent.futures.Future()
224
+ scheduled = concurrent.futures.Future()
225
+
226
+ def start():
227
+ try:
228
+ task = loop.create_task(initialize())
229
+ except Exception as exc: # noqa: BLE001 - notify both waiters of a failed task creation
230
+ scheduled.set_exception(exc)
231
+ result.set_exception(exc)
232
+ return
233
+ scheduled.set_result(task)
234
+
235
+ def finish(done):
236
+ if done.cancelled():
237
+ result.cancel()
238
+ elif (error := done.exception()) is not None:
239
+ result.set_exception(error)
240
+ else:
241
+ result.set_result(done.result())
242
+
243
+ task.add_done_callback(finish)
244
+
245
+ loop.call_soon_threadsafe(start)
246
+ return result, scheduled
247
+
248
+
249
+ async def await_cleanup(awaitable):
250
+ pending = asyncio.ensure_future(awaitable)
251
+ while not pending.done():
252
+ try:
253
+ await asyncio.shield(pending)
254
+ except asyncio.CancelledError:
255
+ task = asyncio.current_task()
256
+ if task is not None:
257
+ task.uncancel()
258
+ return await pending
259
+
260
+
261
+ class VllmScorer:
262
+ """Use exactly the released renderer and E2's masked-logprob readout."""
263
+
264
+ def __init__(self, blink, renderer, engine, sampling_params, tokens_prompt, max_model_len=32768):
265
+ self.blink = blink
266
+ self.renderer = renderer
267
+ self.engine = engine
268
+ self.sampling_params = sampling_params
269
+ self.tokens_prompt = tokens_prompt
270
+ self.max_model_len = max_model_len
271
+
272
+ async def _read(self, item: dict) -> list[float]:
273
+ ids = item["cand"]
274
+ params = self.sampling_params(
275
+ max_tokens=1,
276
+ temperature=1.0,
277
+ logprobs=len(ids),
278
+ allowed_token_ids=list(ids),
279
+ detokenize=False,
280
+ )
281
+ request_id = uuid.uuid4().hex
282
+ output = None
283
+ try:
284
+ async for response in self.engine.generate(
285
+ self.tokens_prompt(prompt_token_ids=item["ids"]), params, request_id
286
+ ):
287
+ if response.finished:
288
+ output = response
289
+ except asyncio.CancelledError:
290
+ await await_cleanup(self.engine.abort(request_id))
291
+ raise
292
+ if output is None or len(output.outputs) != 1 or not output.outputs[0].logprobs:
293
+ raise RuntimeError(f"vLLM returned no complete label scores for {item['qkey']}")
294
+ logprobs = output.outputs[0].logprobs[0]
295
+ if any(token not in logprobs or not math.isfinite(float(logprobs[token].logprob))
296
+ for token in ids):
297
+ raise RuntimeError(f"vLLM omitted an offered-label logprob for {item['qkey']}")
298
+ return [float(logprobs[token].logprob) for token in ids]
299
+
300
+ async def decide(self, state, questions: dict) -> dict:
301
+ self.blink.validate(questions)
302
+ work = self.renderer.render(state, questions)
303
+ for item in work:
304
+ length = len(item["ids"])
305
+ if length + 1 > self.max_model_len:
306
+ raise self.blink.BlinkError(
307
+ f"question {item['qkey']!r} renders to {length} tokens; "
308
+ f"this server's max-model-len is {self.max_model_len}"
309
+ )
310
+ tasks = [asyncio.create_task(self._read(item)) for item in work]
311
+ try:
312
+ rows = await asyncio.gather(*tasks)
313
+ except BaseException:
314
+ for task in tasks:
315
+ task.cancel()
316
+ results = await await_cleanup(asyncio.gather(*tasks, return_exceptions=True))
317
+ for result in results:
318
+ if isinstance(result, Exception):
319
+ print(f"vLLM label cleanup failed: {type(result).__name__}: {result}",
320
+ file=sys.stderr, flush=True)
321
+ raise
322
+ answers = {
323
+ item["qkey"]: self.blink.answer_for(
324
+ questions[item["qkey"]], item["keys"], self.blink.softmax(logits, 1.0)
325
+ )
326
+ for item, logits in zip(work, rows)
327
+ }
328
+ return {"answers": answers, "meta": {"input_tokens": sum(len(item["ids"]) for item in work)}}
329
+
330
+
331
+ class BoundedHTTPServer(ThreadingHTTPServer):
332
+ def __init__(self, address, handler, max_connections: int):
333
+ self.connections = threading.BoundedSemaphore(max_connections)
334
+ super().__init__(address, handler)
335
+
336
+ def process_request(self, request, client_address):
337
+ if not self.connections.acquire(blocking=False):
338
+ print("vLLM HTTP connection limit reached", file=sys.stderr, flush=True)
339
+ self.shutdown_request(request)
340
+ return
341
+ try:
342
+ super().process_request(request, client_address)
343
+ except BaseException:
344
+ self.connections.release()
345
+ raise
346
+
347
+ def process_request_thread(self, request, client_address):
348
+ try:
349
+ super().process_request_thread(request, client_address)
350
+ finally:
351
+ self.connections.release()
352
+
353
+
354
+ def handler_for(blink, scorer: VllmScorer, loop: asyncio.AbstractEventLoop,
355
+ model_id: str, health: dict, key: str | None, max_concurrency: int):
356
+ require_image_contract(blink)
357
+ slots = threading.BoundedSemaphore(max_concurrency)
358
+ listing = {"models": [{
359
+ "name": model_id,
360
+ "description": "blink: typed decisions (noul, choice, score) with option probabilities from one prefill. "
361
+ "This server serves one model; a request's model field is accepted and not used.",
362
+ "release_date": "",
363
+ "accepts_images": False,
364
+ }]}
365
+
366
+ class Handler(BaseHTTPRequestHandler):
367
+ protocol_version = "HTTP/1.1"
368
+ timeout = 60
369
+
370
+ def log_message(self, fmt, *args):
371
+ pass
372
+
373
+ def setup(self):
374
+ super().setup()
375
+ self.connection.setsockopt(socket.IPPROTO_TCP, socket.TCP_NODELAY, 1)
376
+
377
+ def send_json(self, code: int, obj: dict, headers: dict | None = None,
378
+ status_text: str | None = None, request_id: str | None = None) -> None:
379
+ body = json.dumps(obj, ensure_ascii=False, allow_nan=False).encode("utf-8")
380
+ self.send_response(code, status_text)
381
+ self.send_header("Content-Type", "application/json")
382
+ self.send_header("Content-Length", str(len(body)))
383
+ self.send_header("x-typesafe-request-id", request_id or uuid.uuid4().hex)
384
+ for name, value in (headers or {}).items():
385
+ self.send_header(name, value)
386
+ try:
387
+ self.end_headers()
388
+ self.wfile.write(body)
389
+ except (BrokenPipeError, ConnectionResetError):
390
+ self.close_connection = True
391
+
392
+ def fail(self, code: int, reason: str, detail=None, headers: dict | None = None,
393
+ *, unread: bool = False, status_text: str | None = None,
394
+ request_id: str | None = None) -> None:
395
+ if unread:
396
+ self.close_connection = True
397
+ headers = {**(headers or {}), "Connection": "close"}
398
+ self.send_json(code, {"error": reason, "detail": reason if detail is None else detail},
399
+ headers, status_text, request_id)
400
+
401
+ def refused(self, code: int, reason: str, location: list, kind: str, *, unread: bool = False) -> None:
402
+ self.fail(code, reason, [{"loc": location, "msg": reason, "type": kind}], unread=unread)
403
+
404
+ def content_length(self) -> int | None:
405
+ if self.headers.get_all("Transfer-Encoding"):
406
+ self.refused(400, "Transfer-Encoding is not supported", ["body"], "value_error", unread=True)
407
+ return None
408
+ lengths = self.headers.get_all("Content-Length", [])
409
+ if len(lengths) > 1:
410
+ self.refused(400, "duplicate Content-Length", ["body"], "value_error", unread=True)
411
+ return None
412
+ if not lengths:
413
+ return 0
414
+ raw = lengths[0]
415
+ if re.fullmatch(r"[0-9]+", raw) is None:
416
+ self.refused(400, "invalid Content-Length", ["body"], "value_error", unread=True)
417
+ return None
418
+ digits = raw.lstrip("0") or "0"
419
+ if len(digits) > len(str(MAX_BODY_BYTES)):
420
+ return MAX_BODY_BYTES + 1
421
+ return int(digits)
422
+
423
+ def authorized(self) -> bool:
424
+ if key is None:
425
+ return True
426
+ token = bearer(self.headers.get("Authorization"))
427
+ return token is not None and hmac.compare_digest(token.encode(), key.encode())
428
+
429
+ def do_GET(self):
430
+ size = self.content_length()
431
+ if size is None:
432
+ return
433
+ if size:
434
+ return self.refused(400, "GET requests must not include a body",
435
+ ["body"], "value_error", unread=True)
436
+ path = self.path.rstrip("/")
437
+ if path in ("/healthz", "/health"):
438
+ return self.send_json(200 if health["ok"] else 503, health)
439
+ if path == "/v1/models":
440
+ if not self.authorized():
441
+ return self.fail(401, "missing or invalid API key: send Authorization: Bearer <key>",
442
+ headers={"WWW-Authenticate": "Bearer"})
443
+ return self.send_json(200, listing)
444
+ return self.fail(404, "not found")
445
+
446
+ def do_POST(self):
447
+ size = self.content_length()
448
+ if size is None:
449
+ return
450
+ if size > MAX_BODY_BYTES:
451
+ return self.fail(413, "request body exceeds 20 MiB", unread=True)
452
+ if self.path.rstrip("/") != "/v1/systemone":
453
+ return self.fail(404, "not found", unread=True)
454
+ if not self.authorized():
455
+ return self.fail(401, "missing or invalid API key: send Authorization: Bearer <key>",
456
+ headers={"WWW-Authenticate": "Bearer"}, unread=True)
457
+ try:
458
+ payload = self.rfile.read(size)
459
+ except (OSError, TimeoutError) as exc:
460
+ return self.refused(400, f"incomplete request body: {exc}", ["body"], "value_error", unread=True)
461
+ if len(payload) != size:
462
+ return self.refused(400, "incomplete request body", ["body"], "value_error", unread=True)
463
+ try:
464
+ request = json.loads(payload or b"{}")
465
+ except (ValueError, UnicodeDecodeError) as exc:
466
+ return self.refused(400, f"invalid JSON: {exc}", ["body"], "json_invalid")
467
+ if not isinstance(request, dict):
468
+ return self.refused(400, "the body must be a JSON object", ["body"], "value_error")
469
+ try:
470
+ if "images" in request:
471
+ if not isinstance(request["images"], list):
472
+ raise blink.ImageError("this model reads text only", ["body", "images"])
473
+ submission = blink.inspect_images(request.get("state"), request["images"])
474
+ location = (submission.loc if submission is not None and request["images"]
475
+ else ["body", "images"])
476
+ raise blink.ImageError("this model reads text only", location)
477
+ if blink.contains_image_uri(request.get("state")):
478
+ submission = blink.inspect_images(request.get("state"))
479
+ if submission is not None:
480
+ raise blink.ImageError("this model reads text only", submission.loc)
481
+ blink.validate(request.get("questions"))
482
+ except blink.BlinkError as exc:
483
+ loc = exc.loc if isinstance(exc, getattr(blink, "ImageError", ())) else error_loc(blink, request, str(exc))
484
+ return self.refused(422, str(exc), loc, "value_error")
485
+ if not slots.acquire(blocking=False):
486
+ return self.fail(529, "too many concurrent requests",
487
+ headers={"Retry-After": "1"}, status_text="Overloaded")
488
+ future = None
489
+ try:
490
+ future = asyncio.run_coroutine_threadsafe(
491
+ scorer.decide(request.get("state"), request["questions"]), loop
492
+ )
493
+ output = future.result(timeout=180)
494
+ except concurrent.futures.TimeoutError:
495
+ future.cancel()
496
+ return self.fail(504, "decision timed out")
497
+ except blink.BlinkError as exc:
498
+ loc = exc.loc if isinstance(exc, getattr(blink, "ImageError", ())) else error_loc(blink, request, str(exc))
499
+ return self.refused(422, str(exc), loc, "value_error")
500
+ except Exception as exc: # noqa: BLE001 - explicit HTTP error for a failed inference
501
+ request_id = uuid.uuid4().hex
502
+ print(f"vLLM inference failed request_id={request_id}: {type(exc).__name__}: {exc}",
503
+ file=sys.stderr, flush=True)
504
+ return self.fail(500, "internal error", request_id=request_id)
505
+ finally:
506
+ slots.release()
507
+ return self.send_json(200, {
508
+ "model": model_id,
509
+ "answers": output["answers"],
510
+ "usage": {"input_tokens": output["meta"]["input_tokens"], "output_tokens": 0},
511
+ })
512
+
513
+ return Handler
514
+
515
+
516
+ def parse_args(argv=None):
517
+ parser = argparse.ArgumentParser(description=__doc__)
518
+ parser.add_argument("--model", default=os.environ.get("BLINK_MODEL", "thegovind/blink-4b"))
519
+ parser.add_argument("--revision", default=os.environ.get("BLINK_REVISION"))
520
+ parser.add_argument("--host", default="127.0.0.1")
521
+ parser.add_argument("--port", type=int_range(1, 65535), default=8000)
522
+ parser.add_argument("--quantization", choices=("auto", "none", "compressed-tensors"), default="auto")
523
+ parser.add_argument("--max-concurrency", type=int_range(1, 1024), default=32)
524
+ parser.add_argument("--max-num-seqs", type=int_range(1, 1024), default=32)
525
+ parser.add_argument("--max-num-batched-tokens", type=int_range(512, 131072), default=8192)
526
+ parser.add_argument("--max-model-len", type=int_range(512, 131072), default=32768)
527
+ parser.add_argument("--tensor-parallel-size", type=int_range(1, 8), default=1)
528
+ parser.add_argument("--gpu-memory-utilization", type=fraction, default=0.85)
529
+ parser.add_argument("--prefix-cache", action="store_true", help="experimental; default off")
530
+ parser.add_argument("--no-chunked-prefill", action="store_true", help="experimental; default on")
531
+ parser.add_argument("--api-key", type=api_key, default=os.environ.get("BLINK_API_KEY", ""))
532
+ return parser.parse_args(argv)
533
+
534
+
535
+ def main(argv=None) -> None:
536
+ args = parse_args(argv)
537
+ local = Path(args.model).is_dir()
538
+ if local:
539
+ for name in ("HF_HUB_OFFLINE", "TRANSFORMERS_OFFLINE", "HF_HUB_DISABLE_TELEMETRY"):
540
+ os.environ.setdefault(name, "1")
541
+ root = Path(args.model).resolve()
542
+ else:
543
+ from huggingface_hub import snapshot_download
544
+
545
+ root = Path(snapshot_download(args.model, revision=args.revision))
546
+ verified = verify_weights(root)
547
+ quant = detect_quantization(root, args.quantization)
548
+ if quant and not verified:
549
+ raise ValueError("a quantized checkpoint requires a verified weights.sha256")
550
+ blink = load_blink(root)
551
+ require_image_contract(blink)
552
+ from transformers import AutoTokenizer
553
+ from vllm import AsyncEngineArgs, SamplingParams
554
+ from vllm.inputs import TokensPrompt
555
+ from vllm.v1.engine.async_llm import AsyncLLM
556
+
557
+ renderer = blink.TorchEngine.__new__(blink.TorchEngine)
558
+ renderer.tok = AutoTokenizer.from_pretrained(root, local_files_only=local)
559
+ renderer.labels, renderer.label_ids = renderer._verify_labels()
560
+ opts = {
561
+ "model": str(root), "tokenizer": str(root), "dtype": "bfloat16",
562
+ "seed": 0, "max_model_len": args.max_model_len,
563
+ "max_logprobs": 255, "logprobs_mode": "processed_logprobs",
564
+ "hf_overrides": {"head_dtype": "float32"},
565
+ "enable_prefix_caching": args.prefix_cache,
566
+ "mamba_cache_mode": "align" if args.prefix_cache else "none",
567
+ "enable_chunked_prefill": not args.no_chunked_prefill,
568
+ "max_num_seqs": args.max_num_seqs,
569
+ "max_num_batched_tokens": args.max_num_batched_tokens,
570
+ "gpu_memory_utilization": args.gpu_memory_utilization,
571
+ "tensor_parallel_size": args.tensor_parallel_size,
572
+ }
573
+ if quant:
574
+ opts["quantization"] = quant
575
+ if has_vision_model(root):
576
+ opts["limit_mm_per_prompt"] = {"image": 0, "video": 0}
577
+ loop = asyncio.new_event_loop()
578
+
579
+ def run_loop():
580
+ asyncio.set_event_loop(loop)
581
+ loop.run_forever()
582
+
583
+ worker = threading.Thread(target=run_loop, name="blink-vllm-async", daemon=True)
584
+ worker.start()
585
+ owner = EngineOwner()
586
+ stopping = threading.Event()
587
+
588
+ def terminate(_signum, _frame):
589
+ if stopping.is_set():
590
+ return
591
+ raise KeyboardInterrupt
592
+
593
+ previous_sigterm = signal.signal(signal.SIGTERM, terminate)
594
+
595
+ async def initialize():
596
+ try:
597
+ engine = AsyncLLM.from_engine_args(AsyncEngineArgs(**opts))
598
+ owner.adopt(engine)
599
+ effective = engine.vllm_config
600
+ if str(effective.model_config.head_dtype) != "torch.float32":
601
+ raise RuntimeError("vLLM did not retain the FP32 offered-label head")
602
+ if effective.model_config.max_model_len != args.max_model_len:
603
+ raise RuntimeError("vLLM did not retain the requested max-model-len")
604
+ if effective.cache_config.enable_prefix_caching != args.prefix_cache:
605
+ raise RuntimeError("vLLM did not retain the requested prefix-cache setting")
606
+ if effective.scheduler_config.enable_chunked_prefill != (not args.no_chunked_prefill):
607
+ raise RuntimeError("vLLM did not retain the requested chunked-prefill setting")
608
+ scorer = VllmScorer(blink, renderer, engine, SamplingParams, TokensPrompt,
609
+ max_model_len=args.max_model_len)
610
+ first = await scorer.decide(*WARMUP)
611
+ repeat = await scorer.decide(*WARMUP)
612
+ return engine, scorer, first["answers"] == repeat["answers"]
613
+ except BaseException:
614
+ owner.close()
615
+ raise
616
+
617
+ scheduled = None
618
+ try:
619
+ startup, scheduled = submit_initializer(loop, initialize)
620
+ _engine, scorer, repeated = startup.result(timeout=600)
621
+ versions = {}
622
+ for package in ("torch", "transformers", "vllm", "compressed-tensors"):
623
+ try:
624
+ versions[package] = importlib.metadata.version(package)
625
+ except importlib.metadata.PackageNotFoundError:
626
+ versions[package] = None
627
+ health = {
628
+ "ok": True, "model": args.model, "revision": args.revision,
629
+ "weights_verified": verified, "hub_offline": bool(local),
630
+ "warmup": {"repeat_identical": repeated},
631
+ "kernels": "vLLM processed masked logprobs; FP32 offered-label head",
632
+ "versions": versions, "quantization": quant or "none",
633
+ "accepts_images": False,
634
+ "batching": {
635
+ "max_concurrency": args.max_concurrency,
636
+ "max_num_seqs": args.max_num_seqs,
637
+ "max_num_batched_tokens": args.max_num_batched_tokens,
638
+ "chunked_prefill": not args.no_chunked_prefill,
639
+ "prefix_cache": args.prefix_cache,
640
+ "tensor_parallel_size": args.tensor_parallel_size,
641
+ "limit_mm_per_prompt": opts.get("limit_mm_per_prompt"),
642
+ },
643
+ "api_key_required": args.api_key is not None,
644
+ }
645
+ handler = handler_for(blink, scorer, loop, args.model, health, args.api_key, args.max_concurrency)
646
+ server = BoundedHTTPServer((args.host, args.port), handler, args.max_concurrency * 2 + 16)
647
+ server.daemon_threads = True
648
+ try:
649
+ print(f"blink vLLM serving {args.model} on http://{args.host}:{args.port}", flush=True)
650
+ server.serve_forever()
651
+ finally:
652
+ server.server_close()
653
+ finally:
654
+ stopping.set()
655
+ owner.abandon()
656
+ try:
657
+ if scheduled is not None:
658
+ task = scheduled.result()
659
+ if not task.done():
660
+ loop.call_soon_threadsafe(task.cancel)
661
+
662
+ async def settle():
663
+ return await asyncio.gather(task, return_exceptions=True)
664
+
665
+ results = asyncio.run_coroutine_threadsafe(settle(), loop).result()
666
+ if isinstance(results[0], BaseException) and not isinstance(results[0], asyncio.CancelledError):
667
+ print(f"vLLM initialization ended: {type(results[0]).__name__}: {results[0]}",
668
+ file=sys.stderr, flush=True)
669
+ finally:
670
+ try:
671
+ try:
672
+ owner.close()
673
+ finally:
674
+ loop.call_soon_threadsafe(loop.stop)
675
+ worker.join()
676
+ loop.close()
677
+ finally:
678
+ signal.signal(signal.SIGTERM, previous_sigterm)
679
+
680
+
681
+ if __name__ == "__main__":
682
+ main()