Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
630 changes: 564 additions & 66 deletions miles/rollout/fully_async_rollout.py

Large diffs are not rendered by default.

24 changes: 22 additions & 2 deletions miles/rollout/inference_rollout/fully_async.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
import asyncio
from collections.abc import Iterator
from collections.abc import Callable, Iterator
from copy import deepcopy
from typing import TypeVar, cast

Expand Down Expand Up @@ -157,8 +157,14 @@ async def wait_terminal(self) -> FullyAsyncExecutionOutcome:
class InferenceFullyAsyncExecutor(FullyAsyncExecutor):
"""Execute receipt-bound inference groups on the caller's event loop."""

def __init__(self, state: GenerateState) -> None:
def __init__(
self,
state: GenerateState,
*,
sample_done_callback: Callable[[], None] | None,
) -> None:
self._state = state
self._sample_done_callback = sample_done_callback
self._cancellation = _InferenceCancellationCoordinator(state)
self._tasks: set[asyncio.Task[list[Sample | list[Sample]]]] = set()
self._closed = False
Expand Down Expand Up @@ -187,6 +193,7 @@ def submit(
_execute_group(
self._state,
deepcopy(list(reservation.samples)),
sample_done_callback=self._sample_done_callback,
)
)
self._tasks.add(task)
Expand Down Expand Up @@ -225,14 +232,27 @@ async def close(self) -> None:
async def _execute_group(
state: GenerateState,
samples: list[Sample],
*,
sample_done_callback: Callable[[], None] | None,
) -> list[Sample | list[Sample]]:
if sample_done_callback is None:
return cast(
list[Sample | list[Sample]],
await generate_and_rm_group(
state,
samples,
sampling_params=state.sampling_params.copy(),
evaluation=False,
),
)
return cast(
list[Sample | list[Sample]],
await generate_and_rm_group(
state,
samples,
sampling_params=state.sampling_params.copy(),
evaluation=False,
sample_done_callback=sample_done_callback,
),
)

Expand Down
25 changes: 22 additions & 3 deletions miles/rollout/inference_rollout/inference_rollout_common.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
from argparse import Namespace
from collections.abc import Callable
from copy import deepcopy
from typing import Any
from typing import Any, cast

from miles.rollout.base_types import (
GenerateFnInput,
Expand Down Expand Up @@ -199,6 +199,9 @@ async def generate_and_rm_group(
args = state.args

if state.aborted:
if sample_done_callback is not None:
for _ in group:
sample_done_callback()
return group

if policy_uses_routing_key(args):
Expand All @@ -208,7 +211,7 @@ async def generate_and_rm_group(

log_prefix = f"[group indices={[getattr(s, 'index', '?') for s in group]}]"
logger.debug(f"{log_prefix} Starting group with {len(group)} samples")
tasks = []
tasks: list[asyncio.Task[Sample | list[Sample]]] = []
for idx, sample in enumerate(group):
current_sampling_params = sampling_params.copy()
if getattr(args, "sglang_enable_deterministic_inference", False):
Expand All @@ -221,7 +224,23 @@ async def generate_and_rm_group(
task.add_done_callback(lambda _task: sample_done_callback())
tasks.append(task)

group = await asyncio.gather(*tasks)
terminal_wait = asyncio.gather(*tasks, return_exceptions=True)
cancellation: asyncio.CancelledError | None = None
while not terminal_wait.done():
try:
await asyncio.shield(terminal_wait)
except asyncio.CancelledError as error:
if cancellation is None:
cancellation = error
results = terminal_wait.result()
errors = [result for result in results if isinstance(result, BaseException)]
if cancellation is not None:
if errors:
raise cancellation from errors[0]
raise cancellation
if errors:
raise errors[0]
group = cast(list[Sample], [task.result() for task in tasks])
logger.debug(f"{log_prefix} [group] All {len(group)} samples completed")
if state.aborted:
return group
Expand Down
107 changes: 102 additions & 5 deletions tests/fast/rollout/inference_rollout/test_fully_async.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ async def generate(input: GenerateFnInput) -> GenerateFnOutput:


async def test_executor_returns_receipt_bound_success_without_mutating_reservation() -> None:
executor = InferenceFullyAsyncExecutor(make_generate_state())
executor = InferenceFullyAsyncExecutor(make_generate_state(), sample_done_callback=None)
reservation = SourceReservation(
reservation_id=SourceReservationId("source-0"),
samples=(Sample(group_index=0, index=0, prompt="prompt"),),
Expand Down Expand Up @@ -87,7 +87,7 @@ async def generate(input: GenerateFnInput) -> GenerateFnOutput:
return GenerateFnOutput(samples=sample)

state.generate_function = generate
executor = InferenceFullyAsyncExecutor(state)
executor = InferenceFullyAsyncExecutor(state, sample_done_callback=None)
reservation = SourceReservation(
reservation_id=SourceReservationId("source-1"),
samples=(Sample(group_index=1, index=10, prompt="prompt"),),
Expand Down Expand Up @@ -119,7 +119,7 @@ async def generate(input: GenerateFnInput) -> GenerateFnOutput:
return GenerateFnOutput(samples=sample)

state.generate_function = generate
executor = InferenceFullyAsyncExecutor(state)
executor = InferenceFullyAsyncExecutor(state, sample_done_callback=None)
reservation = SourceReservation(
reservation_id=SourceReservationId("source-invalid"),
samples=(Sample(group_index=2, index=20, prompt="prompt"),),
Expand Down Expand Up @@ -153,7 +153,7 @@ async def generate_and_rm_group(state, samples, sampling_params, evaluation=Fals
return generated_samples

monkeypatch.setattr(fully_async_module, "generate_and_rm_group", generate_and_rm_group)
executor = InferenceFullyAsyncExecutor(make_generate_state())
executor = InferenceFullyAsyncExecutor(make_generate_state(), sample_done_callback=None)
reservation = SourceReservation(
reservation_id=SourceReservationId("source-malformed"),
samples=(Sample(group_index=3, index=30, prompt="prompt"),),
Expand All @@ -171,6 +171,103 @@ async def generate_and_rm_group(state, samples, sampling_params, evaluation=Fals
await executor.close()


async def test_cancellation_requests_abort_and_waits_for_terminal_generation(
monkeypatch: pytest.MonkeyPatch,
) -> None:
generation_started = asyncio.Event()
release_generation = asyncio.Event()
abort_requested = asyncio.Event()
release_abort = asyncio.Event()

async def generate(input: GenerateFnInput) -> GenerateFnOutput:
generation_started.set()
await release_generation.wait()
sample = cast(Sample, input.sample)
sample.status = Sample.Status.COMPLETED
sample.reward = 1.0
return GenerateFnOutput(samples=sample)

async def request_abort(args: Namespace) -> None:
abort_requested.set()
await release_abort.wait()

state = make_generate_state()
state.generate_function = generate
monkeypatch.setattr(fully_async_module, "request_abort", request_abort, raising=False)
executor = InferenceFullyAsyncExecutor(state, sample_done_callback=None)
reservation = SourceReservation(
reservation_id=SourceReservationId("source-2"),
samples=(Sample(group_index=2, index=20, prompt="prompt"),),
)
executor_receipt = cast(ReservationExecutorReceipt, object())
execution = executor.submit(reservation, executor_receipt)
terminal_wait = asyncio.create_task(execution.wait_terminal())

await generation_started.wait()
execution.request_cancellation()

try:
await asyncio.wait_for(abort_requested.wait(), timeout=0.01)
release_abort.set()
await asyncio.sleep(0)
assert not terminal_wait.done()
finally:
release_abort.set()
release_generation.set()
outcome = await terminal_wait
await executor.close()

assert outcome == FullyAsyncExecutionRetry(
executor_receipt=executor_receipt,
reason=FullyAsyncRetryReason.CANCELLATION_REQUESTED,
)
assert state.aborted


async def test_cancellation_preserves_terminal_generation_failure(
monkeypatch: pytest.MonkeyPatch,
) -> None:
generation_started = asyncio.Event()
release_generation = asyncio.Event()
abort_requested = asyncio.Event()
generation_error = RuntimeError("generation failed after cancellation")

async def generate(input: GenerateFnInput) -> GenerateFnOutput:
generation_started.set()
await release_generation.wait()
raise generation_error

async def request_abort(args: Namespace) -> None:
abort_requested.set()

state = make_generate_state()
state.generate_function = generate
monkeypatch.setattr(fully_async_module, "request_abort", request_abort)
executor = InferenceFullyAsyncExecutor(state, sample_done_callback=None)
reservation = SourceReservation(
reservation_id=SourceReservationId("source-3"),
samples=(Sample(group_index=3, index=30, prompt="prompt"),),
)
executor_receipt = cast(ReservationExecutorReceipt, object())
execution = executor.submit(reservation, executor_receipt)
terminal_wait = asyncio.create_task(execution.wait_terminal())

await generation_started.wait()
execution.request_cancellation()
await abort_requested.wait()
assert not terminal_wait.done()

release_generation.set()
outcome = await terminal_wait

assert outcome == FullyAsyncExecutionFailure(
executor_receipt=executor_receipt,
error=generation_error,
)

await executor.close()


async def test_executor_close_settles_siblings_before_raising_failure(
monkeypatch: pytest.MonkeyPatch,
) -> None:
Expand All @@ -196,7 +293,7 @@ async def generate_and_rm_group(state, samples, sampling_params, evaluation=Fals
return samples

monkeypatch.setattr(fully_async_module, "generate_and_rm_group", generate_and_rm_group)
executor = InferenceFullyAsyncExecutor(make_generate_state())
executor = InferenceFullyAsyncExecutor(make_generate_state(), sample_done_callback=None)
first_execution = executor.submit(
SourceReservation(
reservation_id=SourceReservationId("source-4"),
Expand Down
Loading
Loading