Skip to content

Records completion barrier waits indefinitely when a credit completes without emitting a record #1273

Description

@janbernloehr

Summary

RecordsTracker.check_and_set_all_records_received() gates phase completion on

all_records_received = self._final_requests_completed is not None and (
    self._success_records + self._error_records >= self._final_requests_completed
)

There is no timeout, deadline, or watchdog anywhere on this path (on main today, src/aiperf/records/records_tracker.py contains no wait_for and no timeout). The predicate is therefore monotonically unsatisfiable if even one credit is counted into final_requests_completed but never produces a success or error record: the phase never completes, the process never exits, and the run has to be killed externally. In our case the benchmark sat completely idle for ~63 minutes with no further output until the job scheduler terminated the step at its time limit.

The correctness of the whole run therefore rests on every producer being individually perfect about emitting exactly one record per counted credit. That invariant is only enforced by convention, and the source acknowledges the hazard in at least four places — e.g. credit/structs.py:

record_emitted: bool = False
"True once an inference record has been pushed for this credit. Used to keep the records-side count in lockstep with the credit-side count: a completed (non-cancelled) credit with no record would hang the RecordsManager completion barrier."

Every time this invariant has been broken so far, the fix has been another producer-side patch (e.g. #1079, and the _emit_credit_failure_record backstop added for v0.12.0) while the barrier itself stayed unbounded. This issue asks for the barrier to be bounded, so the next producer-side gap degrades into a diagnosable failure instead of an indefinite hang.

Field occurrence

Seen on aiperf 0.11.0, driving a multimodal (image) benchmark against an OpenAI-compatible inference server on a single-GPU aarch64 host, at concurrency 128.

The triggering producer-side gap was the shared-mmap read race in MemoryMapDatasetClient.get_conversation (the open fix for that race is #1245). Chain:

  1. Concurrent get_conversation calls are dispatched to a thread pool via run_in_executor(None, self._client.get_conversation, conversation_id), all sharing one mmap object.
  2. get_conversation uses a position-dependent self.data_mmap.seek(offset) followed by self.data_mmap.read(size). A competing thread's seek() can land between them, so read() returns another record's bytes. (mmap_read_method in CPython does not release the GIL, but the executor hop and the interpreter's thread switching are more than enough to interleave the two calls.)
  3. Deserialization of the foreign slice raises MemoryMapSerializationError: Failed to decode conversation data.
  4. That exception propagates out of _retrieve_conversation (which has no handler) into the broad except Exception as e: in the credit-processing path, which records credit_context.error, logs, and returns — without ever calling _send_inference_result_message.
  5. The credit is still counted toward final_requests_completed. The barrier predicate becomes permanently false. The run hangs forever.

Concretely, our accounting showed success + error = 527 against final_requests_completed = 640 — a fixed 113-record deficit that could never close.

Reproduction

The producer-side trigger reproduces deterministically and needs no GPU. Drive MemoryMapDatasetClient.get_conversation from several threads over a single client instance, with records large enough (multimodal base64 image payloads are hundreds of KB) that read() stays in flight long enough to be interleaved. Shrinking the interpreter's thread-switch interval makes the window easy to hit:

import sys, asyncio, threading
from concurrent.futures import ThreadPoolExecutor

sys.setswitchinterval(1e-6)  # default 5ms hides the window

# ... build a MemoryMapDatasetBackingStore with ~24 conversations of ~200KB-4MB each,
# finalize it, then open one MemoryMapDatasetClient over it and call
# client.get_conversation(f"conv-{i}") from 8 threads, 400 lookups each.

On aiperf 0.11.0 this yields, per 3200 reads, a handful of failures such as:

MemoryMapSerializationError: Failed to decode conversation data: 1 validation error for Conversation
  Invalid JSON: EOF while parsing a string at line 1 column 1400525
MemoryMapSerializationError: Failed to decode conversation data: ... trailing characters at line 1 column 800526

and asking for conv-6 can return bytes that begin {"session_id":"conv-14"....

The barrier hang itself follows deterministically from source once any such producer gap exists. A minimal, targeted reproduction (which we have reasoned out from the source but not executed end-to-end) is to inject a single exception into conversation retrieval for one credit on 0.11.0 and observe that the run never terminates.

Expected behavior

A record-accounting deficit should surface as a bounded, diagnosable failure: the run should log the deficit and fail (or complete the phase as failed) rather than blocking forever with no output.

Actual behavior

The phase never completes. The process emits nothing further and must be killed externally. There is no log line reporting the outstanding delta, so from the outside the hang is indistinguishable from a slow benchmark.

Environment

  • aiperf 0.11.0 (installed via uv tool install aiperf==0.11.0 --with tiktoken --with transformers); the unbounded barrier is also present on current main and in v0.12.0.
  • Python 3.12, Linux aarch64, single GPU, concurrency 128, multimodal (image) input.
  • The aiperf --version / environment-collection command could not be run inside the original CI artifact context; the version above is the exact pinned install from the container build, and all source claims in this report were verified directly against the published 0.11.0 wheel, the v0.12.0 tag, and refs/heads/main.

Suggested fix

Bound the barrier instead of relying solely on per-producer discipline:

  1. After the credit phase has finished sending, apply a configurable inactivity deadline to the wait for outstanding records. On expiry, log final_requests_completed vs success + error (and, if cheaply available, which credit IDs are outstanding) and terminate the phase as failed.
  2. Emit a periodic progress line for the outstanding delta while waiting, so an in-progress stall is diagnosable without attaching a debugger.
  3. Optionally assert the lockstep invariant in debug builds, so a producer that completes a credit without emitting a record fails loudly at the source of the bug rather than at the barrier.

This is complementary to, not a substitute for, the producer-side fixes: #1245 removes the race that triggered our incident, and the _emit_credit_failure_record backstop in v0.12.0 closes the specific swallow path. Neither prevents the next unaccounted credit from hanging a run indefinitely.

Validation evidence for the related producer-side fix

Applying #1245's core change — replacing the seek() + read() pair with a position-free slice, self.data_mmap[offset : offset + size], matching what the sibling get_payload_bytes / get_payload_turn methods already do — to the deployed 0.11.0 module eliminated the corruption completely:

Variant Runs Decode errors per 3200 reads
0.11.0 as shipped (seek + read) 4 6, 8, 7, 9
With #1245's position-free slice 4 0, 0, 0, 0

This issue was drafted with assistance from the opus AI model.

Activity

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

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions