Skip to content

perf: send parsed responses to record processors - #1180

Closed
jthomson04 wants to merge 3 commits into
mainfrom
jthomson04/parsed-response-ipc
Closed

jthomson04 wants to merge 3 commits into
mainfrom
jthomson04/parsed-response-ipc

Conversation

@jthomson04

@jthomson04 jthomson04 commented Jul 23, 2026 •

Copy link
Copy Markdown
Contributor

Summary

  • carry the parsed responses already produced by request workers through the internal InferenceResultsMessage
  • compact valid non-RAW built-in responses before worker-to-RecordProcessor serialization
  • retain the existing raw-response path for RAW export, invalid/error records, custom response subclasses, and older messages
  • consume supplied parsed responses in the RecordProcessor without repeating endpoint extraction
  • document the updated internal data flow

Why

After #1172, request workers already parsed completed responses for replay and metrics, but the worker still serialized every raw SSE response and the RecordProcessor parsed the same response stream again. This change reuses the worker-produced result while preserving all fallback and compatibility paths.

Performance

Matched 32K-context, four-turn runs used server token counts, summary export, two request workers and one RecordProcessor on CPUs 0-4, plus Dynamo and four accelerated mock workers on CPUs 5-23. Each selected run used concurrency 1024, a 20-second warmup, and a 60-second measurement.

Metric main candidate Change
Request throughput 916.45 req/s 969.03 req/s +5.7%
Whole-AIPerf CPU/request 4.542 ms 4.170 ms -8.2%
Request-worker CPU/request 2.211 ms 2.090 ms -5.5%
RecordProcessor CPU/request 1.120 ms 0.902 ms -19.5%
p95 latency 155.12 ms 147.04 ms -5.2%
Peak RSS 7.29 GiB 4.43 GiB -39.2%

Saturation probes at concurrency 512, 1024, and 2048 plateaued at 4.16 aggregate AIPerf cores. At 1024 the TimingManager, RecordProcessor, and both request workers were each individually CPU-bound; 2048 consumed no additional CPU and slightly reduced throughput, so 1024 was used instead of chasing an unfillable fifth-core topology gap.

Separate 49 Hz py-spy runs attributed a 52.6% reduction in the combined targeted serialization and duplicate-parsing CPU/request. RecordProcessor parsing/reconstruction fell 81.1%; worker envelope serialization fell 32.4% after accounting for the compact-payload encoding cost.

Both arms completed with zero request errors, cancellations, missing usage, affinity violations, Dynamo lag, skipped messages, or backend rejection.

Validation

  • 99 passed in focused worker, message, endpoint/parser, and RecordProcessor tests
  • 59 passed in property tests
  • pre-commit passed, including Ruff, generated-doc, schema, import, and ergonomics checks
  • full unit run: 15457 passed, 94 skipped, 1 xfailed; the two remaining TTY console failures reproduce unchanged on the pinned main control

Summary by CodeRabbit

  • Improvements
    • Inference result artifacts can now preserve supported parsed response details in IPC for non-raw export modes.
    • Workers now send additional inference metadata (timing, raw response counts, and compaction status).
    • Record processing can reuse pre-parsed responses to avoid redundant extraction.
  • Documentation
    • Updated architecture diagrams and messaging descriptions to clearly separate parsed vs raw end-to-end flows.
  • Tests
    • Added coverage for parsed-response encoding/round-tripping, legacy defaults, compact-message forwarding, and parser reuse behavior.

Signed-off-by: jthomson04 <jwillthomson19@gmail.com>
@github-actions

github-actions Bot commented Jul 23, 2026 •

Copy link
Copy Markdown

Try out this PR

Quick install:

pip install --upgrade --force-reinstall git+https://github.com/ai-dynamo/aiperf.git@9c256b765a6ab35ac770189d9158b782c74e05d2

Recommended with virtual environment (using uv):

uv venv --python 3.12 && source .venv/bin/activate
uv pip install --upgrade --force-reinstall git+https://github.com/ai-dynamo/aiperf.git@9c256b765a6ab35ac770189d9158b782c74e05d2

Last updated for commit: 9c256b7 • Browse code

@github-actions github-actions Bot added the perf label Jul 23, 2026
@github-actions

github-actions Bot commented Jul 23, 2026 •

Copy link
Copy Markdown

@codecov

codecov Bot commented Jul 23, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 97.46835% with 2 lines in your changes missing coverage. Please review.

Files with missing lines Patch % Lines
src/aiperf/records/record_processor_service.py 71.42% 1 Missing and 1 partial ⚠️

📢 Thoughts on this report? Let us know!

@coderabbitai

coderabbitai Bot commented Jul 23, 2026 •

Copy link
Copy Markdown
Contributor

Review Change Stack

Walkthrough

The change adds validated IPC payloads for built-in parsed responses, threads them through worker result publishing, and lets record processing reuse them without re-parsing. Raw responses remain available for raw exports, custom data, errors, or invalid timestamps. Architecture documentation and unit tests are updated.

Changes

Parsed inference IPC flow

Layer / File(s) Summary
Parsed response payload contract
src/aiperf/common/messages/inference_messages.py, tests/unit/common/messages/test_inference_messages.py
Built-in parsed responses can be encoded and reconstructed through validated payloads, while InferenceResultsMessage carries parsed responses and compaction metadata with legacy defaults.
Worker response propagation and IPC publishing
src/aiperf/workers/worker.py, tests/unit/workers/test_worker.py
Worker response paths return parsed responses and conditionally include encoded payloads based on export level, while setting response counts, timing, validation state, and raw-response retention.
Record processor parsed-response fast path
src/aiperf/records/inference_result_parser.py, src/aiperf/records/record_processor_service.py, tests/unit/records/*, docs/architecture.md
Record processing reconstructs supplied parsed responses, forwards IPC metadata, bypasses duplicate endpoint extraction, and documents the parsed-versus-raw message flow.

Estimated code review effort: 3 (Moderate) | ~25 minutes

Poem

I’m a rabbit with packets tucked neatly away,
Parsed hops travel the tunnel today.
Raw tails stay when export says “please,”
Built-in types cross with validated ease.
The record processor smiles at the stream—
Fewer re-hops in the inference dream.

🚥 Pre-merge checks | ✅ 4 | ❌ 1

❌ Failed checks (1 warning)

Check name Status Explanation Resolution
Docstring Coverage ⚠️ Warning Docstring coverage is 14.29% which is insufficient. The required threshold is 80.00%. Write docstrings for the functions missing them to satisfy the coverage threshold.
✅ Passed checks (4 passed)
Check name Status Explanation
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
Description Check ✅ Passed Check skipped - CodeRabbit’s high-level summary is enabled.
Title check ✅ Passed The title clearly matches the main change: workers now send parsed responses through to record processors.

Comment @coderabbitai help to get the list of available commands.

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick comments (1)
src/aiperf/workers/worker.py (1)

546-552: 🩺 Stability & Availability | 🔵 Trivial | ⚡ Quick win

Add a regression test for the try/finally safety net.

This try/finally is a meaningful correctness fix: if _finalize_session_response raises, _send_inference_result_message still fires (with parsed_responses=None, falling back to raw responses) instead of silently dropping the record. That's exactly the lockstep hazard record_processor_service.py's _on_inference_results docstring calls out (a dropped record permanently stalls the timeout-less completion barrier). No test in this diff exercises that raise-and-still-send path.

Consider adding a test that makes _finalize_session_response raise and asserts _send_inference_result_message is still awaited (with parsed_responses=None).

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@src/aiperf/workers/worker.py` around lines 546 - 552, Add a regression test
for the worker flow around _finalize_session_response that makes it raise, then
verifies _send_inference_result_message is still awaited with
parsed_responses=None. Preserve the exception behavior while asserting the
record is sent through the finally safety path.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Nitpick comments:
In `@src/aiperf/workers/worker.py`:
- Around line 546-552: Add a regression test for the worker flow around
_finalize_session_response that makes it raise, then verifies
_send_inference_result_message is still awaited with parsed_responses=None.
Preserve the exception behavior while asserting the record is sent through the
finally safety path.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: CHILL

Plan: Enterprise

Run ID: 2742a21d-94de-41de-af4c-fbaf179b2335

📥 Commits

Reviewing files that changed from the base of the PR and between 94ece0f and b68b3fe.

📒 Files selected for processing (9)
  • docs/architecture.md
  • src/aiperf/common/messages/inference_messages.py
  • src/aiperf/records/inference_result_parser.py
  • src/aiperf/records/record_processor_service.py
  • src/aiperf/workers/worker.py
  • tests/unit/common/messages/test_inference_messages.py
  • tests/unit/records/test_inference_result_parser.py
  • tests/unit/records/test_record_processor_service.py
  • tests/unit/workers/test_worker.py

@jthomson04 jthomson04 left a comment

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Reviewed the parsed-response IPC change end-to-end (including a standalone round-trip run under -W error::UserWarning). The correctness-critical invariants all hold:

  • last_response_perf_ns's gt=0 validator can never trip, because RequestRecord.valid already requires all(0 < response.perf_ns ...) and the worker only sets it on validated records.
  • record.responses = None is dropped via exclude_none=True and rehydrated to [] by default_factory, so the untyped-None-into-list[...] field round-trips cleanly.
  • Worker parsing and RecordProcessor parsing are the same call (endpoint.extract_response_data; no process_responses override exists anywhere), so the reused results are equivalent.
  • The compaction guard export_level != RAW exactly matches the RecordProcessor's own raw-freeing guard, and the only post-parse reader of raw request.responses is the RAW writer (raw_record_writer_processor.py:108, RAW-gated). Accuracy grading and every metric read the parsed parsed_record.responses, and token counting reads parsed responses + payload_bytes. Compaction is observationally equivalent to existing behavior.
  • Skipping create_error_from_invalid under responses_validated is necessary (the raw list is now empty) and safe (the downstream ParsedResponseRecord.create_error_from_invalid() still catches content-empty parsed records).
  • The try/finally reorder guarantees the message is always sent and falls back to the raw re-parse path when _finalize_session_response raises.

Overall: approve with minor nits. The change is well-fenced — every non-happy path (RAW export, errors, custom subclasses, legacy messages, finalize failure) falls back to the pre-existing raw path. Inline comments below are readability/robustness polish, not blockers.

Two test gaps worth a quick add: an explicit worker-level ExportLevel.RECORDS case (shares the != RAW branch with SUMMARY), and a case for the try/finally fallback where _finalize_session_response raises and the raw path is taken with parsed_responses=None.

Comment thread src/aiperf/records/record_processor_service.py Outdated
Comment thread src/aiperf/common/messages/inference_messages.py Outdated
Comment thread src/aiperf/common/messages/inference_messages.py Outdated
Comment thread src/aiperf/workers/worker.py
Signed-off-by: jthomson04 <jwillthomson19@gmail.com>
@ajcasagrande

Copy link
Copy Markdown
Contributor

Not sure I'm sold on this one.

Signed-off-by: jthomson04 <jwillthomson19@gmail.com>

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🧹 Nitpick comments (2)
tests/unit/common/messages/test_inference_messages.py (2)

23-72: 🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick win

Cover the atomic encoding fallback.

Add a case containing a custom or unsupported response and assert that encode_parsed_responses returns None; otherwise a regression could incorrectly compact responses that must remain on the raw path.

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@tests/unit/common/messages/test_inference_messages.py` around lines 23 - 72,
Extend test_inference_results_parsed_responses_round_trip_builtin_types with a
custom or unsupported ParsedResponse case and verify encode_parsed_responses
returns None for that input. Keep the existing round-trip assertions unchanged,
ensuring unsupported responses cannot be compacted and remain on the raw path.

80-89: 🗄️ Data Integrity & Integration | 🔵 Trivial | ⚡ Quick win

Assert that the legacy payload omits the new fields.

The restored defaults can match even if the fields were serialized explicitly. Verify that all four keys are absent before calling from_json:

Proposed assertion
     legacy_data = message.model_dump(
         mode="json",
         exclude={
             "parsed_responses",
             "last_response_perf_ns",
             "raw_response_count",
             "responses_compacted",
         },
     )
+    assert {
+        "parsed_responses",
+        "last_response_perf_ns",
+        "raw_response_count",
+        "responses_compacted",
+    }.isdisjoint(legacy_data)
     restored = InferenceResultsMessage.from_json(legacy_data)
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@tests/unit/common/messages/test_inference_messages.py` around lines 80 - 89,
Update the test around InferenceResultsMessage serialization to assert that
legacy_data excludes parsed_responses, last_response_perf_ns,
raw_response_count, and responses_compacted before calling from_json. Keep the
existing restoration assertions unchanged.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Nitpick comments:
In `@tests/unit/common/messages/test_inference_messages.py`:
- Around line 23-72: Extend
test_inference_results_parsed_responses_round_trip_builtin_types with a custom
or unsupported ParsedResponse case and verify encode_parsed_responses returns
None for that input. Keep the existing round-trip assertions unchanged, ensuring
unsupported responses cannot be compacted and remain on the raw path.
- Around line 80-89: Update the test around InferenceResultsMessage
serialization to assert that legacy_data excludes parsed_responses,
last_response_perf_ns, raw_response_count, and responses_compacted before
calling from_json. Keep the existing restoration assertions unchanged.

ℹ️ Review info
⚙️ Run configuration

Configuration used: Path: .coderabbit.yaml

Review profile: CHILL

Plan: Enterprise

Run ID: 4ecf6ea4-4f9f-4c42-81d8-f22acec2957d

📥 Commits

Reviewing files that changed from the base of the PR and between 8b92586 and 9c256b7.

📒 Files selected for processing (2)
  • tests/unit/common/messages/test_inference_messages.py
  • tests/unit/workers/test_worker.py
🚧 Files skipped from review as they are similar to previous changes (1)
  • tests/unit/workers/test_worker.py

@jthomson04
jthomson04 marked this pull request as ready for review July 24, 2026 04:10
@github-actions

Copy link
Copy Markdown

This PR is stale because it has been open 30 days with no activity. Remove stale label or comment or this will be closed in 5 days.

@github-actions github-actions Bot added the stale label Sep 11, 2026
@github-actions

Copy link
Copy Markdown

This PR has been closed due to inactivity. If you believe this PR is still relevant, please feel free to reopen it with additional context or information.

@github-actions github-actions Bot closed this Sep 16, 2026
@github-actions
github-actions Bot deleted the jthomson04/parsed-response-ipc branch September 16, 2026 09:37
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants