# `vllm_mlx.output_collector`

Output collector for streaming with low-latency optimizations.

[View the complete module source at #L1-L212](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L1-L212).

## API details

Each callable below includes its exact signature, type annotations, inputs, defaults, return contract, documented exceptions, implementation source, and parsed docstring sections when the source provides them.

::: vllm_mlx.output_collector
    options:
      members:
        - RequestOutputCollector
        - RequestStreamState
      filters: []
      show_if_no_docstring: true

## Complete contract reference

Expand any definition for its exact inputs, annotations, defaults, return contract, directly raised exceptions, source-grounded behavior, and immutable line link. This section includes private and nested definitions that ordinary API generators omit.

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestOutputCollector" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestOutputCollector</code> · class</summary>

```python
vllm_mlx.output_collector.RequestOutputCollector(aggregate: bool = True)
```

Per-request output collector with smart buffering.

**Parameters**

| Name | Type | Required | Default | Description |
| --- | --- | --- | --- | --- |
| `aggregate` | `bool` | `no` | `True` | If True, merge outputs when producer gets ahead. This prevents buffer explosion under load. |

**Returns**

- Constructs: `vllm_mlx.output_collector.RequestOutputCollector`

**Exceptions and behavior**

Class `RequestOutputCollector` declares 7 direct member(s).
No direct `raise` statement appears in this definition.

[View source #L17-L170](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L17-L170).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestOutputCollector.__init__" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestOutputCollector.__init__</code> · method</summary>

```python
vllm_mlx.output_collector.RequestOutputCollector.__init__(aggregate: bool = True) -> not annotated
```

Initialize the collector.

**Parameters**

| Name | Type | Required | Default | Description |
| --- | --- | --- | --- | --- |
| `aggregate` | `bool` | `no` | `True` | If True, merge outputs when producer gets ahead. This prevents buffer explosion under load. |

**Returns**

- Type: `not annotated`

**Exceptions and behavior**

Method `RequestOutputCollector.__init__` updates `self.output`, `self.ready`, `self.aggregate`, `self._is_waiting`; calls `asyncio.Event`.
No direct `raise` statement appears in this definition.

[View source #L42-L53](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L42-L53).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestOutputCollector.put" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestOutputCollector.put</code> · method</summary>

```python
vllm_mlx.output_collector.RequestOutputCollector.put(output: RequestOutput) -> None
```

Put an output into the collector (non-blocking).

**Parameters**

| Name | Type | Required | Default | Description |
| --- | --- | --- | --- | --- |
| `output` | `RequestOutput` | `yes` | `none` | The RequestOutput to store |

**Returns**

- Type: `None`

**Exceptions and behavior**

Method `RequestOutputCollector.put` updates `self.output`; calls `self._merge_outputs`, `self.ready.set`.
No direct `raise` statement appears in this definition.

[View source #L55-L73](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L55-L73).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestOutputCollector.get_nowait" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestOutputCollector.get_nowait</code> · method</summary>

```python
vllm_mlx.output_collector.RequestOutputCollector.get_nowait() -> Optional[RequestOutput]
```

Get output without blocking.

**Parameters**

This callable has no explicit inputs.

**Returns**

- Type: `Optional[RequestOutput]`
- Direct return expressions: `output`

**Exceptions and behavior**

Method `RequestOutputCollector.get_nowait` updates `self.output`; calls `self.ready.clear`; returns `output`.
No direct `raise` statement appears in this definition.

[View source #L75-L89](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L75-L89).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestOutputCollector.get" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestOutputCollector.get</code> · method</summary>

```python
async vllm_mlx.output_collector.RequestOutputCollector.get() -> RequestOutput
```

Get output, blocking only if none available.

**Parameters**

This callable has no explicit inputs.

**Returns**

- Type: `RequestOutput`
- Direct return expressions: `output`

**Exceptions and behavior**

Method `RequestOutputCollector.get` updates `self._is_waiting`; calls `self.ready.wait`, `self.get_nowait`; awaits asynchronous work; returns `output`.
No direct `raise` statement appears in this definition.

[View source #L91-L118](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L91-L118).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestOutputCollector._merge_outputs" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestOutputCollector._merge_outputs</code> · method</summary>

```python
vllm_mlx.output_collector.RequestOutputCollector._merge_outputs(existing: RequestOutput, new: RequestOutput) -> RequestOutput
```

Merge two outputs when producer gets ahead of consumer.

**Parameters**

| Name | Type | Required | Default | Description |
| --- | --- | --- | --- | --- |
| `existing` | `RequestOutput` | `yes` | `none` | The existing output in the buffer |
| `new` | `RequestOutput` | `yes` | `none` | The new output to merge |

**Returns**

- Type: `RequestOutput`
- Direct return expressions: `RequestOutput(request_id=new.request_id, new_token_ids=merged_new_token_ids, new_text=merged_new_text, output_token_ids…`

**Exceptions and behavior**

Method `RequestOutputCollector._merge_outputs` calls `RequestOutput`; returns `RequestOutput(request_id=new.request_id, new_token_ids=merged_new_token_ids, new_text=merged_new_text, output_token_ids…`.
No direct `raise` statement appears in this definition.

[View source #L120-L152](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L120-L152).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestOutputCollector.clear" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestOutputCollector.clear</code> · method</summary>

```python
vllm_mlx.output_collector.RequestOutputCollector.clear() -> None
```

Clear any pending output.

**Parameters**

This callable has no explicit inputs.

**Returns**

- Type: `None`

**Exceptions and behavior**

Method `RequestOutputCollector.clear` updates `self.output`, `self._is_waiting`; calls `self.ready.clear`.
No direct `raise` statement appears in this definition.

[View source #L154-L161](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L154-L161).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestOutputCollector.has_waiting_consumers" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestOutputCollector.has_waiting_consumers</code> · method</summary>

```python
vllm_mlx.output_collector.RequestOutputCollector.has_waiting_consumers() -> bool
```

Check if any collector has waiting consumers.

**Parameters**

This callable has no explicit inputs.

**Returns**

- Type: `bool`
- Direct return expressions: `cls._waiting_consumers > 0`

**Exceptions and behavior**

Method `RequestOutputCollector.has_waiting_consumers` returns `cls._waiting_consumers > 0`.
No direct `raise` statement appears in this definition.

[View source #L164-L170](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L164-L170).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestStreamState" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestStreamState</code> · class</summary>

```python
vllm_mlx.output_collector.RequestStreamState(stream_interval: int = 1, sent_tokens: int = 0)
```

Tracks streaming state for a request.

**Parameters**

| Name | Type | Required | Default | Description |
| --- | --- | --- | --- | --- |
| `stream_interval` | `int` | `no` | `1` | Optional constructor field; defaults to `1`. |
| `sent_tokens` | `int` | `no` | `0` | Optional constructor field; defaults to `0`. |

**Returns**

- Constructs: `vllm_mlx.output_collector.RequestStreamState`

**Exceptions and behavior**

Class `RequestStreamState` declares 2 direct member(s).
No direct `raise` statement appears in this definition.

[View source #L174-L212](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L174-L212).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestStreamState.should_send" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestStreamState.should_send</code> · method</summary>

```python
vllm_mlx.output_collector.RequestStreamState.should_send(total_tokens: int, finished: bool) -> bool
```

Determine if output should be sent based on stream_interval.

**Parameters**

| Name | Type | Required | Default | Description |
| --- | --- | --- | --- | --- |
| `total_tokens` | `int` | `yes` | `none` | Total tokens generated so far |
| `finished` | `bool` | `yes` | `none` | Whether generation is complete |

**Returns**

- Type: `bool`
- Direct return expressions: `True`; `total_tokens - self.sent_tokens >= self.stream_interval`

**Exceptions and behavior**

Method `RequestStreamState.should_send` has 2 explicit return paths.
No direct `raise` statement appears in this definition.

[View source #L185-L203](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L185-L203).

</details>

<details class="api-contract" id="contract-vllm_mlx.output_collector.RequestStreamState.mark_sent" markdown="1">
<summary><code>vllm_mlx.output_collector.RequestStreamState.mark_sent</code> · method</summary>

```python
vllm_mlx.output_collector.RequestStreamState.mark_sent(total_tokens: int) -> None
```

Update state after sending output.

**Parameters**

| Name | Type | Required | Default | Description |
| --- | --- | --- | --- | --- |
| `total_tokens` | `int` | `yes` | `none` | Total tokens at time of send |

**Returns**

- Type: `None`

**Exceptions and behavior**

Method `RequestStreamState.mark_sent` updates `self.sent_tokens`.
No direct `raise` statement appears in this definition.

[View source #L205-L212](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L205-L212).

</details>

## Complete symbol map

This map also includes private definitions and nested helpers. The signature column exposes every explicit input even when an internal helper has no dedicated parameter prose.

| Symbol | Kind | Signature and inputs | What it does | Source |
| --- | --- | --- | --- | --- |
| [`RequestOutputCollector`](#contract-vllm_mlx.output_collector.RequestOutputCollector) | class | `RequestOutputCollector(aggregate: bool = True)` | Per-request output collector with smart buffering. | [#L17-L170](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L17-L170) |
| [`RequestOutputCollector.__init__`](#contract-vllm_mlx.output_collector.RequestOutputCollector.__init__) | method | `RequestOutputCollector.__init__(aggregate: bool = True) -> not annotated` | Initialize the collector. | [#L42-L53](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L42-L53) |
| [`RequestOutputCollector.put`](#contract-vllm_mlx.output_collector.RequestOutputCollector.put) | method | `RequestOutputCollector.put(output: RequestOutput) -> None` | Put an output into the collector (non-blocking). | [#L55-L73](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L55-L73) |
| [`RequestOutputCollector.get_nowait`](#contract-vllm_mlx.output_collector.RequestOutputCollector.get_nowait) | method | `RequestOutputCollector.get_nowait() -> Optional[RequestOutput]` | Get output without blocking. | [#L75-L89](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L75-L89) |
| [`RequestOutputCollector.get`](#contract-vllm_mlx.output_collector.RequestOutputCollector.get) | method | `async RequestOutputCollector.get() -> RequestOutput` | Get output, blocking only if none available. | [#L91-L118](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L91-L118) |
| [`RequestOutputCollector._merge_outputs`](#contract-vllm_mlx.output_collector.RequestOutputCollector._merge_outputs) | method | `RequestOutputCollector._merge_outputs(existing: RequestOutput, new: RequestOutput) -> RequestOutput` | Merge two outputs when producer gets ahead of consumer. | [#L120-L152](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L120-L152) |
| [`RequestOutputCollector.clear`](#contract-vllm_mlx.output_collector.RequestOutputCollector.clear) | method | `RequestOutputCollector.clear() -> None` | Clear any pending output. | [#L154-L161](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L154-L161) |
| [`RequestOutputCollector.has_waiting_consumers`](#contract-vllm_mlx.output_collector.RequestOutputCollector.has_waiting_consumers) | method | `RequestOutputCollector.has_waiting_consumers() -> bool` | Check if any collector has waiting consumers. | [#L164-L170](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L164-L170) |
| [`RequestStreamState`](#contract-vllm_mlx.output_collector.RequestStreamState) | class | `RequestStreamState(stream_interval: int = 1, sent_tokens: int = 0)` | Tracks streaming state for a request. | [#L174-L212](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L174-L212) |
| [`RequestStreamState.should_send`](#contract-vllm_mlx.output_collector.RequestStreamState.should_send) | method | `RequestStreamState.should_send(total_tokens: int, finished: bool) -> bool` | Determine if output should be sent based on stream_interval. | [#L185-L203](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L185-L203) |
| [`RequestStreamState.mark_sent`](#contract-vllm_mlx.output_collector.RequestStreamState.mark_sent) | method | `RequestStreamState.mark_sent(total_tokens: int) -> None` | Update state after sending output. | [#L205-L212](https://github.com/waybarrios/vllm-mlx/blob/a69d47912bcb21d8fe04d48f75fa896b620ffcfa/vllm_mlx/output_collector.py#L205-L212) |
