perf: send parsed responses to record processors - #1180
jthomson04 wants to merge 3 commits into
Conversation
Signed-off-by: jthomson04 <jwillthomson19@gmail.com>
Try out this PRQuick install: pip install --upgrade --force-reinstall git+https://github.com/ai-dynamo/aiperf.git@9c256b765a6ab35ac770189d9158b782c74e05d2Recommended 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@9c256b765a6ab35ac770189d9158b782c74e05d2Last updated for commit: |
Codecov Report❌ Patch coverage is
📢 Thoughts on this report? Let us know! |
WalkthroughThe 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. ChangesParsed inference IPC flow
Estimated code review effort: 3 (Moderate) | ~25 minutes Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
Comment |
There was a problem hiding this comment.
🧹 Nitpick comments (1)
src/aiperf/workers/worker.py (1)
546-552: 🩺 Stability & Availability | 🔵 Trivial | ⚡ Quick winAdd a regression test for the try/finally safety net.
This try/finally is a meaningful correctness fix: if
_finalize_session_responseraises,_send_inference_result_messagestill fires (withparsed_responses=None, falling back to raw responses) instead of silently dropping the record. That's exactly the lockstep hazardrecord_processor_service.py's_on_inference_resultsdocstring 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_responseraise and asserts_send_inference_result_messageis still awaited (withparsed_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
📒 Files selected for processing (9)
docs/architecture.mdsrc/aiperf/common/messages/inference_messages.pysrc/aiperf/records/inference_result_parser.pysrc/aiperf/records/record_processor_service.pysrc/aiperf/workers/worker.pytests/unit/common/messages/test_inference_messages.pytests/unit/records/test_inference_result_parser.pytests/unit/records/test_record_processor_service.pytests/unit/workers/test_worker.py
jthomson04
left a comment
There was a problem hiding this comment.
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'sgt=0validator can never trip, becauseRequestRecord.validalready requiresall(0 < response.perf_ns ...)and the worker only sets it on validated records.record.responses = Noneis dropped viaexclude_none=Trueand rehydrated to[]bydefault_factory, so the untyped-None-into-list[...]field round-trips cleanly.- Worker parsing and RecordProcessor parsing are the same call (
endpoint.extract_response_data; noprocess_responsesoverride exists anywhere), so the reused results are equivalent. - The compaction guard
export_level != RAWexactly matches the RecordProcessor's own raw-freeing guard, and the only post-parse reader of rawrequest.responsesis the RAW writer (raw_record_writer_processor.py:108, RAW-gated). Accuracy grading and every metric read the parsedparsed_record.responses, and token counting reads parsed responses +payload_bytes. Compaction is observationally equivalent to existing behavior. - Skipping
create_error_from_invalidunderresponses_validatedis necessary (the raw list is now empty) and safe (the downstreamParsedResponseRecord.create_error_from_invalid()still catches content-empty parsed records). - The
try/finallyreorder guarantees the message is always sent and falls back to the raw re-parse path when_finalize_session_responseraises.
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.
Signed-off-by: jthomson04 <jwillthomson19@gmail.com>
|
Not sure I'm sold on this one. |
Signed-off-by: jthomson04 <jwillthomson19@gmail.com>
There was a problem hiding this comment.
🧹 Nitpick comments (2)
tests/unit/common/messages/test_inference_messages.py (2)
23-72: 🎯 Functional Correctness | 🔵 Trivial | ⚡ Quick winCover the atomic encoding fallback.
Add a case containing a custom or unsupported response and assert that
encode_parsed_responsesreturnsNone; 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 winAssert 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
📒 Files selected for processing (2)
tests/unit/common/messages/test_inference_messages.pytests/unit/workers/test_worker.py
🚧 Files skipped from review as they are similar to previous changes (1)
- tests/unit/workers/test_worker.py
|
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. |
|
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. |
Summary
InferenceResultsMessageWhy
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.
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 passedin focused worker, message, endpoint/parser, and RecordProcessor tests59 passedin property tests15457 passed, 94 skipped, 1 xfailed; the two remaining TTY console failures reproduce unchanged on the pinned main controlSummary by CodeRabbit