Pull requests / #1562

#1562 serve: compact streamed line fragments into blocks

open · @InB4DevOps · 0 commentaires · Sur GitHub

Setup & installServer & APIWindowsLinux

Description

## Summary

Avoid copying the entire growing current line for each streamed text fragment in `OutputParser._track()`.

- Append newline-free fragments and compact every 64 pending fragments into a completed block.
- Keep completed blocks separate until a newline, parser decision or state snapshot needs the full line.
- Preserve string reads/assignments through the `line` property, including snapshot restoration and fence/inline-code decisions.
- Reset the internal buffers directly at newlines to limit reset overhead.
- Add parity tests against both the original tracking algorithm and original string storage.

This change is independent of #1542 and #1557. The branch contains only line tracking/storage and its tests. Partial-tag matching, tool-body closing-tag scans, GPU/native code and profiling infrastructure are unchanged.

## Standalone validation

Both the untouched baseline and isolated candidate started at upstream `fb58e0dbc8399662c0e47c76578c6e878b14f6cf`. The candidate contained no other local optimizations.

```text
python -m unittest serve.test_parser_track serve.test_frontend serve.test_reasoning_tools serve.test_server
Ran 298 tests — OK
```

Coverage includes per-fragment line/fence/backtick state, Unicode, CRLF, varied initial states, long lines, quoted tool markers, recovery, streamed tools, compaction boundaries, materialization followed by appends, and snapshot restoration. The reference parser overrides the new property with the old string attribute as well as retaining the old `_track()` implementation. `git diff --check` passed.

**Windows correctness and performance validation: pending.**

## Measurements and tradeoffs

Intel Core i7-12700KF, 20 logical CPUs; Linux 7.0.0-38-generic, glibc 2.39, CPython 3.12.3, Jinja2 3.1.6. Default scheduling/power policy, no affinity pinning. Seven fresh-process samples per checkout, alternating baseline/candidate order, two warmups per case, fixed 16-character chunks. Fixture construction and validation were outside timing. No model, GPU, or profiling timers were involved.

Median elapsed ms, untouched upstream → PR:

| Workload | Scale 4 | Scale 16 |
|---|---:|---:|
| Unbroken content line (128 / 512 KiB) | 19.93 → 10.68 | 325.67 → 41.92 |
| Ordinary multilingual code lines | 3.88 → 4.00 | 15.77 → 16.11 |
| Newline on every 16-character fragment | 13.34 → 14.27 | 54.47 → 58.09 |
| Many short streamed tools (128 / 512 calls) | 3.11 → 3.02 | 12.14 → 12.07 |

The 512 KiB unbroken line used **87.1% less elapsed time** (candidate range 41.56–43.59 ms, baseline 322.00–332.07 ms). Candidate time grew about 3.9× for 4× the line length in this workload. These are synthetic parser CPU results, **not end-to-end inference speedups**.

**Tradeoff:** ordinary multiline medians were 2–3% slower, and the deliberately newline-heavy case was about 7% slower. The table retains those regressions. Frequent explicit line inspections can still force repeated joins; this does not claim linear scaling for every parser input or fix tool-body accumulation.

All 112 measured outputs matched fixture content/final arguments/streamed JSON and had matching hashes across both checkouts and samples. Random tool IDs were excluded.

### Memory

A separate untimed `tracemalloc` run fed 32,768 distinct generated 16-character fragments (512 KiB) and discarded emitted events:

| Traced Python allocation | Upstream | PR |
|---|---:|---:|
| Retained before reading `line` | 513.0 KiB | 537.6 KiB |
| Peak including final join | 1,025.3 KiB | 1,049.7 KiB |
| Retained after reading `line` | 513.1 KiB | 513.1 KiB |

Compaction keeps retained memory near the original representation while avoiding per-fragment whole-line copies. This measures Python allocations, not RSS; total stored text still grows with line length.

## Minimal reproduction

With Jinja2 installed, run the following from each checkout root. It reproduces one 512 KiB long-line timing sample after two warmups, followed by a separate memory check. Repeat timing in fresh processes and alternate checkout order. Set `scale = 4` for the smaller timing case.

```python
import gc
import time
import tracemalloc
from serve.frontend import OutputParser

scale = 16
text = "x" * (32768 * scale)
chunks = [text[i:i + 16] for i in range(0, len(text), 16)]
tools = [{"name": "write", "parameters": {"properties": {"text": {"type": "string"}}}}]

def run():
    parser = OutputParser(thinking=False, tools=tools, stream_tools=False, recover=False)
    content, calls, args = [], [], []
    def collect(events):
        for event in events:
            if event.kind == "content":
                content.append(event.text)
            elif event.kind == "tool_call":
                calls.append((event.call.name, event.call.arguments))
            elif event.kind == "tool_args":
                args.append(event.text)
    for chunk in chunks:
        collect(parser.feed(chunk))
    collect(parser.finish())
    return "".join(content), calls, "".join(args)

for _ in range(2):
    assert run() == (text, [], "")
start = time.perf_counter_ns()
result = run()
elapsed_ms = (time.perf_counter_ns() - start) / 1e6
assert result == (text, [], "")
print("elapsed_ms", elapsed_ms)

gc.collect()
tracemalloc.start()
parser = OutputParser(thinking=False)
for i in range(32768):
    parser.feed(f"{i:016x}")
retained, peak = tracemalloc.get_traced_memory()
line = parser.line
after, joined_peak = tracemalloc.get_traced_memory()
tracemalloc.stop()
assert line == "".join(f"{i:016x}" for i in range(32768))
print("bytes: retained, peak, after read, peak including read", retained, peak, after, joined_peak)
```

Sur le site

Liens install, modèles, releases.