Skip to content
Open
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
44 changes: 44 additions & 0 deletions experimental/aisim/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
# experimental/aisim — offline rollout replay adapter

> **Where it fits:** an optional, default-off offline tool. It observes a real UniRL rollout
> window, converts it into an existing AISimulate / Dynamo Replay input, and compares the
> replay against the matched real execution. It does not schedule, score, train, or feed the
> optimizer, and normal training must run with this package absent.

## Layout

| File | Role |
| --- | --- |
| `schema.py` | Versioned trace contract: manifest, workload, measured outcomes. Stdlib only. |
| `convert.py` | Workload envelope → upstream replay input. Reads no measured outcomes. |
| `runner.py` | Not present. A runner is the only file allowed to import the simulator API, and it lands once a pinned simulator revision is confirmed. |

There is deliberately no `requirements.txt`: the adapter has no dependency of its own, and the
simulator belongs to a separate environment (see the repository rule against adding
non-additive requirements under `experimental/`).

## Boundaries that are load-bearing

- **Three inputs, one direction.** `manifest.json` and `workload.jsonl` are predictor input;
`measured_requests.jsonl` is report input only. `convert.py` cannot read a measured result
because no function accepts one, and `assert_no_outcome_leakage` re-checks the converted rows.
- **Observed arrivals are preserved, not reconstructed.** The open-loop contract keeps the
real arrival offsets (rebased so the first submission is 0) and rejects a window that also
carries dependency arrivals. It never adds a `session_id` or a synthesised tool wait.
- **A missing wait is not a zero wait.** `external_delay_ms` is required for a dependency
arrival; `0.0` has to be written explicitly.
- **Unsupported lifecycles are refused, not trimmed.** `failed` / `aborted` / `partial` /
`retry` terminations raise, and `capture_status="incomplete"` cannot produce a fidelity run.
Dropping the failures and calling the remainder a complete rollout is the failure mode this
check exists to prevent.
- **A GRPO group is not a barrier, and a trajectory slot is not an in-flight request.** Neither
is modelled here; `to_dependency` labels both `not_modelled` in its sidecar rather than
inventing a representation.

## Status of this package

The UniRL-side contract, converters and their CPU tests are implemented. Everything that needs
the simulator itself — locating the pinned trace parser, running the CLI, request-level result
export, trajectory-slot admission and the collector barrier — is **NOT_RUN**: no AISimulate /
Dynamo package is installed in the target environment, and the docs were only read. See the
evidence directory's capability report for the exact probe results.
1 change: 1 addition & 0 deletions experimental/aisim/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
"""Optional offline AISimulate/Dynamo Replay integration (see README)."""
186 changes: 186 additions & 0 deletions experimental/aisim/convert.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,186 @@
"""UniRL workload envelope → upstream replay input, with no outcome leakage (see the package README)."""

from __future__ import annotations

import json
from dataclasses import dataclass, field
from typing import Any, Dict, Mapping, Sequence, Tuple

from experimental.aisim.schema import (
SCHEMA_VERSION,
Manifest,
WorkloadRequest,
validate_workload,
)

# Upstream-facing field names. Kept in one place because the pinned simulator revision owns
# their meaning; this adapter never renames them locally.
REQUEST_ID = "request_id"
INPUT_TOKENS = "input_tokens"
OUTPUT_TOKENS = "output_tokens"
ARRIVAL_MS = "arrival_ms"
WAIT_FOR = "wait_for"
TOOL_WAIT_MS = "tool_wait_ms"


@dataclass(frozen=True)
class ConvertedWorkload:
"""One converted window: predictor input plus the UniRL-only sidecar."""

contract: str
requests: Tuple[Dict[str, Any], ...]
sidecar: Dict[str, Any] = field(default_factory=dict)
notes: Tuple[str, ...] = ()

def to_json(self) -> str:
"""Serialise the converted workload and its sidecar."""
return json.dumps(
{
"contract": self.contract,
"requests": list(self.requests),
"sidecar": self.sidecar,
"notes": list(self.notes),
},
indent=1,
sort_keys=True,
)


def normalize_observed_offsets(requests: Sequence[WorkloadRequest]) -> Tuple[Tuple[float, ...], float]:
"""Return observed offsets rebased so the first submission is 0, plus the origin."""
observed = [request for request in requests if request.arrival.kind == "observed"]
if not observed:
raise ValueError("normalize_observed_offsets: the window has no observed arrival to anchor on")
origin = min(float(request.arrival.offset_ms) for request in observed)
return tuple(float(request.arrival.offset_ms) - origin for request in observed), origin


def to_open_loop(manifest: Manifest, requests: Sequence[WorkloadRequest]) -> ConvertedWorkload:
"""Convert observed-arrival requests into the open-loop contract, adding no session/wait field."""
manifest.validate()
validate_workload(requests)
if manifest.replay_contract != "open_loop_observed_arrival":
raise ValueError(
f"to_open_loop: manifest.replay_contract is {manifest.replay_contract!r}; "
"an observed-arrival conversion must not be run under a dependency contract"
)
dependency = [request.request_id for request in requests if request.arrival.kind != "observed"]
if dependency:
raise ValueError(
f"to_open_loop: request(s) {dependency[:5]} carry dependency arrivals; convert them with to_dependency"
)
if manifest.capture_status == "incomplete":
raise ValueError(
"to_open_loop: capture_status='incomplete' — an incomplete trace cannot produce a fidelity run"
)

offsets, origin = normalize_observed_offsets(requests)
rows = []
for request, offset in zip([r for r in requests if r.arrival.kind == "observed"], offsets):
rows.append(
{
REQUEST_ID: request.request_id,
INPUT_TOKENS: int(request.input_tokens),
OUTPUT_TOKENS: int(request.output_tokens),
ARRIVAL_MS: float(offset),
}
)
sidecar = {
"unirl_schema_version": SCHEMA_VERSION,
"run_id": manifest.run_id,
"arrival_origin_ms": origin,
"requests": [
{
"request_id": request.request_id,
"rollout_batch_id": request.rollout_batch_id,
"group_id": request.group_id,
"trajectory_id": request.trajectory_id,
"turn_index": request.turn_index,
"attempt_id": request.attempt_id,
"engine_id": request.engine_id,
"worker_id": request.worker_id,
"policy_identity": request.policy_identity,
}
for request in requests
],
}
return ConvertedWorkload(
contract="open_loop_observed_arrival",
requests=tuple(rows),
sidecar=sidecar,
notes=("observed arrivals are preserved verbatim; no dependency or session field was added",),
)


def to_dependency(manifest: Manifest, requests: Sequence[WorkloadRequest]) -> ConvertedWorkload:
"""Convert predecessor-driven requests, keeping measured service time out of the predictor input."""
manifest.validate()
validate_workload(requests)
observed = [request.request_id for request in requests if request.arrival.kind == "observed"]
if observed:
raise ValueError(
f"to_dependency: request(s) {observed[:5]} still carry observed arrivals; a dependency workload must not keep "
"an arrival the simulator would also honour"
)
rows = []
for request in requests:
row: Dict[str, Any] = {
REQUEST_ID: request.request_id,
INPUT_TOKENS: int(request.input_tokens),
OUTPUT_TOKENS: int(request.output_tokens),
}
if request.arrival.wait_for:
row[WAIT_FOR] = list(request.arrival.wait_for)
row[TOOL_WAIT_MS] = float(request.arrival.external_delay_ms)
else:
# First turn: the window releases it at an absolute offset instead of via a predecessor.
row[ARRIVAL_MS] = float(request.arrival.offset_ms)
rows.append(row)
sidecar = {
"unirl_schema_version": SCHEMA_VERSION,
"run_id": manifest.run_id,
"slot_semantics": "not_modelled",
"barrier_semantics": "not_modelled",
}
return ConvertedWorkload(
contract="dependency_arrival",
requests=tuple(rows),
sidecar=sidecar,
notes=(
"external_delay_ms is a measured external wait, never a difference of neighbouring submit timestamps",
"trajectory slots and collector barriers are not represented in this contract",
),
)


def result_sidecar(results: Sequence[Mapping[str, Any]]) -> Dict[str, Any]:
"""Wrap measured outcomes for the report only; never merge this into a converter input."""
return {
"purpose": "report_only",
"measured": [dict(row) for row in results],
}


def assert_no_outcome_leakage(converted: ConvertedWorkload) -> None:
"""Fail if any measuring-outcome key reached the predictor-facing rows."""
forbidden = {"submit_ms", "complete_ms", "first_token_ms", "latency_ms", "service_time_ms", "status"}
for row in converted.requests:
leaked = sorted(forbidden & set(row))
if leaked:
raise ValueError(f"predictor input carries measured-outcome key(s) {leaked} for {row.get(REQUEST_ID)!r}")


def describe(converted: ConvertedWorkload) -> str:
"""Short human summary for logs and PR evidence."""
return f"{converted.contract}: {len(converted.requests)} request(s); notes={list(converted.notes)}"


__all__ = [
"ConvertedWorkload",
"assert_no_outcome_leakage",
"describe",
"normalize_observed_offsets",
"result_sidecar",
"to_dependency",
"to_open_loop",
]
Loading
Loading