Skip to content

fix: avoid quadratic buffering of long SSE lines - #76

Merged
keelerm84 merged 3 commits into
launchdarkly:mainfrom
bach-ta:fix/quadratic-line-buffering
Sep 17, 2026
Merged

keelerm84 merged 3 commits into
launchdarkly:mainfrom
bach-ta:fix/quadratic-line-buffering

Conversation

@bach-ta

@bach-ta bach-ta commented Sep 16, 2026 •

Copy link
Copy Markdown
Contributor

Requirements

  • I have added test coverage for new or changed functionality
  • I have followed the repository's pull request submission guidelines
  • I have validated my changes against all supported platform versions

Related issues

No linked issue.

Describe the solution you've provided

When a long SSE line spans many chunks, _BufferedLineReader.lines_from() repeatedly concatenates the accumulated partial line with the next fragment. Each concatenation copies the growing prefix, making buffering quadratic in line length for a fixed chunk size.

This change collects fragments in a list and joins them once the line ends. It preserves the existing splitting and decoding logic, including CR/LF/CRLF handling and discarding an unterminated final line.

The tests add coverage for fragmented lines with empty chunks, CRLF split across chunks, and UTF-8 characters split across byte boundaries.

Describe alternatives you've considered

Larger HTTP chunks reduce the number of copies but leave the quadratic behavior. A mutable byte buffer is another option; retaining the existing fragments and joining once keeps the change small and preserves the current parsing logic.

Additional context

The benchmark compares the published launchdarkly-eventsource==1.7.2 reader (d10fd70) against this patch (745b247) on Python 3.12.14, macOS 15.7.3, Apple M4 Pro.

Both implementations read the same synthetic data: line in 10,000-byte chunks. Inputs are prepared before timing; the timed region consumes the reader's output, including decoding. Every result is checked against the exact expected list of lines.

Each case runs nine times per implementation, in separate fresh processes with an identical warm-up. Case and implementation order are randomized with seed 0. Values below are median wall times, with Q1–Q3 in brackets. Growth is relative to the preceding row, calculated from unrounded medians.

Line size (MiB) Original ms [Q1–Q3] Growth Patched ms [Q1–Q3] Growth
1 1.55 [1.52–1.61] — 0.47 [0.47–0.48] —
2 8.49 [8.13–9.18] 5.47× 1.01 [0.87–1.04] 2.16×
4 37.05 [36.88–38.07] 4.36× 2.03 [2.03–2.05] 2.02×
8 154.24 [150.57–184.88] 4.16× 3.97 [3.94–3.99] 1.95×
16 857.70 [831.67–900.31] 5.56× 7.74 [7.70–7.86] 1.95×
32 5980.44 [5502.67–6105.69] 6.97× 17.16 [17.09–17.79] 2.22×

Line size excludes the two trailing LF bytes; 1 MiB = 1,048,576 bytes. Doubling the line size takes approximately 4–7× as long in the original reader and approximately 2× with the patch. A separate short-line control (about 1 MiB total) measured 1.30ms original / 1.32ms patched.

These measurements isolate line reading and decoding, not HTTP or full SDK initialization. Absolute timings will vary by machine.

Reproduction

Save the two scripts below as bench_patch.sh and bench_patch.py in the same directory. They require macOS or Linux, Bash, Git, tar, and Python with venv and pip. Package installation requires network access.

With a clean local checkout of this PR, run from the directory containing both scripts:

LD_BENCH_REPO=/path/to/python-eventsource \
LD_BENCH_PYTHON=python3.12 \
bash bench_patch.sh

The script compares release 1.7.2 with the checkout's committed HEAD and writes the summary, individual measurements, source snapshots, and environment metadata to reader-benchmark/. It also measures additional chunk sizes and the short-line control. Each run replaces the previous output.

bench_patch.sh
#!/usr/bin/env bash
# Usage: bash bench_patch.sh [--smoke]
set -euo pipefail

if [[ $# -gt 1 || (${1:-} != "" && ${1:-} != --smoke) ]]; then
    echo "Usage: bash $0 [--smoke]" >&2
    exit 2
fi

SCRIPT_DIR=$(cd -- "$(dirname -- "${BASH_SOURCE[0]}")" && pwd)
REPO=${LD_BENCH_REPO:-"$HOME/python-eventsource"}
PYTHON=${LD_BENCH_PYTHON:-python3.12}
BASE=d10fd70c7575d01e9e7f89625952be94a83ae301
PATCH=$(git -C "$REPO" rev-parse HEAD)

if [[ -n $(git -C "$REPO" status --porcelain --untracked-files=no) ]]; then
    echo "Commit your tracked changes first; this script benchmarks HEAD." >&2
    exit 1
fi
if git -C "$REPO" diff --quiet "$BASE" "$PATCH" -- ld_eventsource/reader.py; then
    echo "HEAD contains no reader change relative to release 1.7.2." >&2
    exit 1
fi
"$PYTHON" --version

RUN="$SCRIPT_DIR/reader-benchmark"
mkdir -p "$RUN"
rm -f "$RUN/summary.md" "$RUN/summary.json" "$RUN/manifest.json" "$RUN/results.jsonl" "$RUN/output.txt"
WORK=$(mktemp -d "${TMPDIR:-/tmp}/ld-reader-benchmark-work.XXXXXX")
trap 'rm -rf -- "$WORK"' EXIT
echo "Results: $RUN"
printf '%s\n' "$BASE" > "$RUN/original-commit.txt"
printf '%s\n' "$PATCH" > "$RUN/patched-commit.txt"
git -C "$REPO" show "$BASE:ld_eventsource/reader.py" > "$RUN/original_reader.py"
git -C "$REPO" show "$PATCH:ld_eventsource/reader.py" > "$RUN/patched_reader.py"
git -C "$REPO" diff "$BASE" "$PATCH" -- ld_eventsource/reader.py ld_eventsource/testing/test_reader.py > "$RUN/patch.diff"
mkdir "$WORK/source"
git -C "$REPO" archive "$PATCH" | tar -x -C "$WORK/source"
cp "$SCRIPT_DIR/bench_patch.py" "$SCRIPT_DIR/bench_patch.sh" "$RUN/"

echo "Installing original and patched packages in separate environments..."
"$PYTHON" -m venv "$WORK/original-venv"
"$PYTHON" -m venv "$WORK/patched-venv"
"$WORK/original-venv/bin/python" -I -m pip install --disable-pip-version-check 'launchdarkly-eventsource==1.7.2' > "$RUN/install-original.log" 2>&1 || { cat "$RUN/install-original.log"; exit 1; }
"$WORK/original-venv/bin/python" -I -m pip freeze > "$RUN/original-requirements.txt"
"$WORK/patched-venv/bin/python" -I -m pip install --disable-pip-version-check -r "$RUN/original-requirements.txt" > "$RUN/install-patched.log" 2>&1 || { cat "$RUN/install-patched.log"; exit 1; }
"$WORK/patched-venv/bin/python" -I -m pip install --disable-pip-version-check --no-deps --force-reinstall "$WORK/source" >> "$RUN/install-patched.log" 2>&1 || { cat "$RUN/install-patched.log"; exit 1; }
"$WORK/patched-venv/bin/python" -I -m pip freeze > "$RUN/patched-requirements.txt"

LD_BENCH_WORK="$WORK" "$WORK/original-venv/bin/python" -I "$RUN/bench_patch.py" "$@" | tee "$RUN/output.txt"
echo "Saved raw measurements, source snapshots, metadata, and summary.md in $RUN"
bench_patch.py
"""Compare installed original and patched readers using synthetic inputs."""

import gc
import hashlib
import importlib.metadata
import json
import os
import platform
import random
import resource
import statistics
import subprocess
import sys
import time
from datetime import datetime, timezone
from pathlib import Path

ROOT = Path(__file__).resolve().parent
MIB = 1024 * 1024
VARIANTS = ("original", "patched")


def sha(data):
    return hashlib.sha256(data).hexdigest()


def make_input(shape, size):
    if shape == "long":
        line = b"data: " + b"x" * (size - 6)
        return line + b"\n\n", [line.decode(), ""]
    line = b"data: " + b"x" * 72
    count = size // 80
    return (line + b"\n\n") * count, [line.decode(), ""] * count


def worker(variant, shape, size, chunk_size):
    import ld_eventsource.reader as module

    source_sha = sha(Path(module.__file__).read_bytes())
    if source_sha != sha((ROOT / f"{variant}_reader.py").read_bytes()):
        raise RuntimeError(f"Installed {variant} reader differs from its source snapshot")
    reader = module._BufferedLineReader.lines_from
    data, expected = make_input(shape, size)
    chunks = tuple(data[i:i + chunk_size] for i in range(0, len(data), chunk_size))
    if list(reader([b"x" * 1000] * 64 + [b"\n"])) != ["x" * 64000]:
        raise RuntimeError("Warm-up output mismatch")
    gc.collect()

    cpu_start = time.process_time_ns()
    wall_start = time.perf_counter_ns()
    output = list(reader(iter(chunks)))
    wall_ns = time.perf_counter_ns() - wall_start
    cpu_ns = time.process_time_ns() - cpu_start

    rss = resource.getrusage(resource.RUSAGE_SELF).ru_maxrss
    if output != expected:
        raise RuntimeError(f"Incorrect output from {variant}")
    return {
        "variant": variant, "shape": shape, "size": size, "chunk_size": chunk_size,
        "stream_bytes": len(data), "input_sha256": sha(data), "source_sha256": source_sha,
        "wall_ms": wall_ns / 1e6, "cpu_ms": cpu_ns / 1e6,
        "worker_peak_rss_bytes": rss if sys.platform == "darwin" else rss * 1024,
        "python": sys.version,
        "packages": {d.metadata["Name"]: d.version for d in importlib.metadata.distributions()},
    }


def cpu_model():
    if sys.platform == "darwin":
        return subprocess.check_output(["sysctl", "-n", "machdep.cpu.brand_string"], text=True).strip()
    path = Path("/proc/cpuinfo")
    if path.exists():
        for line in path.read_text().splitlines():
            if line.startswith("model name"):
                return line.partition(":")[2].strip()
    return platform.processor()


def summarize(rows, cases):
    summary = []
    for shape, size, chunk_size in cases:
        for variant in VARIANTS:
            samples = [r for r in rows if (r["shape"], r["size"], r["chunk_size"], r["variant"])
                       == (shape, size, chunk_size, variant)]
            item = dict(shape=shape, size=size, chunk_size=chunk_size, variant=variant, n=len(samples))
            for metric in ("wall_ms", "cpu_ms", "worker_peak_rss_bytes"):
                values = [r[metric] for r in samples]
                q1, _, q3 = statistics.quantiles(values, n=4, method="inclusive")
                item[metric] = dict(median=statistics.median(values), q1=q1, q3=q3)
            summary.append(item)
    return summary


def report(summary, manifest):
    index = {(r["shape"], r["size"], r["chunk_size"], r["variant"]): r for r in summary}
    sizes = sorted({r["size"] for r in summary if r["shape"] == "long" and r["chunk_size"] == 10000})
    lines = [
        "# Reader benchmark" + (" — SMOKE TEST ONLY" if manifest["smoke"] else ""), "",
        f"Python {platform.python_version()} · {manifest['cpu']} · {platform.platform()}", "",
        f"Original: `{manifest['commits']['original']}`; patched: `{manifest['commits']['patched']}`.", "",
        f"{manifest['repeats']} fresh-process pairs per case, seed 0. Timings are medians; brackets show Q1–Q3.",
        "Inputs are prepared before timing. Both readers must produce the exact expected output.", "",
        "| Line size (MiB) | Original ms [Q1–Q3] | Growth | Patched ms [Q1–Q3] | Growth |",
        "|---:|---:|---:|---:|---:|",
    ]
    previous = {}
    for size in sizes:
        cells = [f"{size / MIB:g}"]
        for variant in VARIANTS:
            stats = index["long", size, 10000, variant]["wall_ms"]
            median = stats["median"]
            growth = f"{median / previous[variant]:.2f}×" if variant in previous else "—"
            cells += [f"{median:.2f} [{stats['q1']:.2f}–{stats['q3']:.2f}]", growth]
            previous[variant] = median
        lines.append("| " + " | ".join(cells) + " |")
    lines += ["", "Growth is relative to the preceding size, using unrounded medians and 10,000-byte chunks.",
              "Line size excludes two trailing LF bytes; 1 MiB = 1,048,576 bytes.", "",
              "| Case | Chunk bytes | Original wall / CPU ms | Patched wall / CPU ms |",
              "|---|---:|---:|---:|"]
    for shape, size, chunk_size in manifest["cases"]:
        cells = [f"{shape}, {size / MIB:g} MiB", str(chunk_size)]
        for variant in VARIANTS:
            row = index[shape, size, chunk_size, variant]
            cells.append(f"{row['wall_ms']['median']:.2f} / {row['cpu_ms']['median']:.2f}")
        lines.append("| " + " | ".join(cells) + " |")
    lines += ["", "Raw samples are in results.jsonl; summary.json also includes whole-worker peak RSS.",
              "RSS includes the interpreter and prepared inputs, not just reader allocations.",
              "These timings isolate line reading and decoding; they do not measure full SDK initialization.", ""]
    return "\n".join(lines)


def run(smoke):
    venv_root = Path(os.environ["LD_BENCH_WORK"])
    repeats = 2 if smoke else 9
    sizes = [16 * 1024, 32 * 1024] if smoke else [n * MIB for n in (1, 2, 4, 8, 16, 32)]
    cases = [("long", size, 10000) for size in sizes]
    sensitivity_size = sizes[0] if smoke else 8 * MIB
    cases += [("long", sensitivity_size, n) for n in (2000, 50000, 200000)]
    cases.append(("short", 8192 if smoke else MIB, 10000))
    manifest = {
        "captured_at_utc": datetime.now(timezone.utc).isoformat(),
        "python": sys.version, "platform": platform.platform(), "machine": platform.machine(),
        "cpu": cpu_model(), "smoke": smoke, "seed": 0, "repeats": repeats, "cases": cases,
        "commits": {v: (ROOT / f"{v}-commit.txt").read_text().strip() for v in VARIANTS},
        "sources": {v: sha((ROOT / f"{v}_reader.py").read_bytes()) for v in VARIANTS},
        "driver_sha256": sha(Path(__file__).read_bytes()),
    }
    (ROOT / "manifest.json").write_text(json.dumps(manifest, indent=2) + "\n")
    rng = random.Random(0)
    rows = []
    total = repeats * len(cases) * 2
    with (ROOT / "results.jsonl").open("w") as stream:
        for repeat in range(repeats):
            order = list(cases)
            rng.shuffle(order)
            for shape, size, chunk_size in order:
                variants = list(VARIANTS)
                rng.shuffle(variants)
                pair = []
                for variant in variants:
                    command = [str(venv_root / f"{variant}-venv/bin/python"), "-I", str(Path(__file__).resolve()),
                               "--worker", variant, shape, str(size), str(chunk_size)]
                    row = json.loads(subprocess.check_output(command, text=True, timeout=300))
                    if row["python"] != sys.version or row["source_sha256"] != manifest["sources"][variant]:
                        raise RuntimeError("Worker runtime or source mismatch")
                    row.update(repeat=repeat, sequence=len(rows))
                    stream.write(json.dumps(row) + "\n")
                    stream.flush()
                    rows.append(row)
                    pair.append(row)
                if pair[0]["input_sha256"] != pair[1]["input_sha256"] or pair[0]["packages"] != pair[1]["packages"]:
                    raise RuntimeError("Inputs or installed dependency versions differ")
                print(f"{len(rows)}/{total}: {shape}, {size / MIB:g} MiB, {chunk_size}-byte chunks", flush=True)
    summary = summarize(rows, cases)
    (ROOT / "summary.json").write_text(json.dumps(summary, indent=2) + "\n")
    markdown = report(summary, manifest)
    (ROOT / "summary.md").write_text(markdown)
    print("\n" + markdown)


if __name__ == "__main__":
    if len(sys.argv) > 1 and sys.argv[1] == "--worker":
        print(json.dumps(worker(sys.argv[2], sys.argv[3], int(sys.argv[4]), int(sys.argv[5]))))
    else:
        if sys.argv[1:] not in ([], ["--smoke"]):
            raise SystemExit("Usage: bench_patch.py [--smoke]")
        run(smoke="--smoke" in sys.argv)

Validation

  • Unit tests: 136 passed on Python 3.14.7.
  • Type checking, import ordering, and style checks: passed (mypy, isort, pycodestyle).

Note

Overview
Fixes quadratic-time buffering when SSE lines arrive split across many byte chunks in both sync (_BufferedLineReader) and async (_AsyncBufferedLineReader) line readers.

Instead of prepending each new fragment onto a growing partial_line (re-copying the prefix every chunk), partial bytes are accumulated in pending_fragments and b"".join runs once when the line terminator appears. CR/LF/CRLF handling and withholding an unterminated final line are unchanged.

Tests add coverage for fragments with empty chunks, terminators split across chunks, and long UTF-8 lines split on byte boundaries; the sync mixed-terminator case is adjusted to exercise chunk boundaries more strictly.

Reviewed by Cursor Bugbot for commit 573d3c6. Bugbot is set up for automated code reviews on this repo. Configure here.

@bach-ta
bach-ta marked this pull request as ready for review September 16, 2026 06:02
@bach-ta
bach-ta requested a review from a team as a code owner September 16, 2026 06:02
The async reader is a hand-maintained parallel copy of the sync reader, so
it had the same quadratic buffering of a long line. It rebuilt the growing
partial line on every chunk. This applies the same fix: collect the
fragments in a list, and join them one time when the line ends.

CR, LF and CRLF handling is unchanged, and an unterminated final line is
still discarded. The new tests mirror the sync tests for fragmented input
with empty chunks and for UTF-8 characters split across chunk boundaries.

A 16 MiB line in 10,000-byte chunks now reads in about 10ms instead of
about 790ms, and the time per doubling of the line size drops from about
4x-6x to about 2x.
Both branches of the pending-fragment block started with the same append,
so the append moves above the branch. This also makes the sync and async
readers read the same, which matters because they are parallel files that
are maintained by hand.

The behavior does not change.
@keelerm84

Copy link
Copy Markdown
Member

Thanks for the thorough writeup and the reproducible benchmark. The analysis is right and the sync reader change is correct. I pushed two commits to the branch rather than asking you to round-trip on them.

What changed

6bfd331 -- the same fix for the async reader. ld_eventsource/async_reader.py is a hand-maintained parallel copy of the sync reader and carried the identical lines[0] = partial_line + lines[0], fed from the same _CHUNK_SIZE = 10000 in async_http.py. So the quadratic behavior described above survived on the async path, which is the half the async SDK client uses. The fix transplants verbatim. Tests mirroring the two you added (fragmented input with empty chunks, UTF-8 split across chunk boundaries) are now in test_async_reader.py, parametrized over \r, \n, and \r\n.

Worth flagging why this was easy to miss: test_sync_async_parity.py exists to catch sync/async drift, but it compares public attribute names only, so a divergence in a private implementation is invisible to it. Nothing to fix here, just a known gap.

573d3c6 -- a behavior-neutral cleanup. Both branches of the pending-fragment block started with the same pending_fragments.append(lines[0]), so that moved above the branch. The point is that the sync and async lines_from bodies are now line-for-line identical apart from the async/await keywords, which is what you want for two files kept in step by hand.

Async reader performance

Same shape as your benchmark: one long data: line in 10,000-byte chunks, input prepared before timing, output checked against the exact expected list, 9 fresh-process runs per case with an identical warm-up, case and variant order randomized with seed 0. Medians with Q1-Q3 in brackets; growth is relative to the preceding row from unrounded medians. "Before" is the async reader as of 745b247, "after" is 6bfd331.

Python 3.13.0, Intel Core Ultra 7 155H, Linux 7.1.5 x86_64.

Line size (MiB) Before ms [Q1-Q3] Growth After ms [Q1-Q3] Growth
1 2.60 [2.48-2.61] -- 1.01 [1.01-1.03] --
2 9.13 [8.96-9.36] 3.52x 2.10 [2.01-2.12] 2.08x
4 33.54 [33.32-34.76] 3.67x 4.02 [3.97-4.15] 1.92x
8 167.61 [165.46-171.98] 5.00x 8.38 [8.10-8.46] 2.08x
16 925.04 [924.00-941.57] 5.52x 16.75 [16.53-17.25] 2.00x
32 18203.95 [17757.21-18342.56] 19.68x 34.13 [33.80-35.01] 2.04x

Line size excludes the two trailing LF bytes; 1 MiB = 1,048,576 bytes. Doubling the line size costs about 2x after the fix, against 3.5x to 5.5x before it. Do not read too much into the 19.68x on the last row: at 32 MiB the unpatched reader is thrashing the allocator on top of the copying, so that figure is above the underlying quadratic trend rather than a truer measure of it.

Correctness verification

I checked both readers as behavior-preserving refactors rather than reading them closely, since the new fast path reorders when partial data is merged. Against the pre-PR implementations, over every possible chunking of every byte string in {a, \n, \r} up to length 7 (167,962 cases for the sync reader, to length 6 for the async one), plus 300,000 randomized cases for the sync reader and 20,000 for the async one, mixing \r\n, empty chunks, multi-byte UTF-8, and the bytes that str.splitlines treats as boundaries but bytes.splitlines does not (\v, \f, \x1c, \x85): no divergence in output or in exceptions raised.

Two invariants make the fast path safe, for the record. len(lines) == 1 and not terminated implies the chunk holds no terminator at all, so lines[0] is the whole chunk. And a non-empty pending_fragments is mutually exclusive with last_char_was_cr, because the flag is only set on the last_char == 13 branch, which is the one branch that never leaves a partial line, so the CRLF lines.pop(0) fixup can never coexist with a pending fragment. terminated = last_char in (10, 13) is also the correct boundary set: bytes.splitlines() splits on 10 and 13 and nothing else.

One trade-off worth recording

The fragment list costs roughly 41 bytes per fragment, so peak memory now depends on chunk size in a way it did not before. Peak above baseline for a 2 MiB line (tracemalloc, same figures for both readers):

Chunk size Before After
10,000 B 4.20 MB 4.20 MB
1,400 B 4.20 MB 4.20 MB
64 B 4.20 MB 5.00 MB
1 B 4.20 MB 187.00 MB

Identical at production chunk sizes, and 1-byte chunks are not reachable through urllib3 or aiohttp in practice. Not an objection to the approach, just something that belongs in the record next to the timing numbers.

@keelerm84
keelerm84 merged commit 8ede399 into launchdarkly:main Sep 17, 2026
17 checks passed
keelerm84 pushed a commit that referenced this pull request Sep 17, 2026
🤖 I have created a release *beep* *boop*
---


##
[1.7.3](1.7.2...1.7.3)
(2026-09-17)


### Bug Fixes

* avoid quadratic buffering of long SSE lines
([#76](#76))
([8ede399](8ede399))

---
This PR was generated with [Release
Please](https://github.com/googleapis/release-please). See
[documentation](https://github.com/googleapis/release-please#release-please).

Co-authored-by: github-actions[bot] <41898282+github-actions[bot]@users.noreply.github.com>
@keelerm84

Copy link
Copy Markdown
Member

@bach-ta thank you again for your contribution. This has been released in v1.7.3

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants