Persistent rollout with straw#
straw provides filesystem-based durable
queues and packed tensor storage for vime. With --rollout-data-transport straw,
prompt tasks, partial rollouts, accepted results and training batches have
persistent identities. Multiple generation and training processes access their
payloads directly through a shared filesystem; Ray schedules actors and carries
control messages and small references.
The current shared-filesystem target is JuiceFS. Each node must mount the same pool at the same absolute path. Other network filesystems, including NFS, need separate qualification of cross-client locks, visibility and durability.
What changes in the architecture#
Concern |
Ray object-store transport (default) |
straw transport |
|---|---|---|
Data source |
|
|
Payload exchange |
Ray object references |
|
Worker assignment |
In-memory producer state |
Leases, attempts, durable continuations and accepted receipts |
Rollout-to-training conversion |
|
The same conversion, with a persisted selection plan and ready batch |
Data lifetime |
Process/object-store ownership |
Explicit task, consumer, reader and checkpoint ownership |
Cleanup |
Object-store lifetime |
Optional online GC of sealed packs with no remaining owners |
There are two logical stages: rollout work/results, and ready training batches. They share the same straw storage pool, so a tensor retained by either stage or a checkpoint remains live. These stages do not require two copies of each tensor. The vime controller coordinates both stages; straw implements the storage, journal, task acceptance and ownership protocol in Rust.
flowchart LR
D[Dataset producer] -->|cursor + task submission| Q[Durable prompt tasks]
Q -->|lease + input reference| W[Distributed generation workers]
W -->|durable partial continuation| Q
W -->|accepted group reference| R[Rollout results]
R --> B[BatchBuilder: selection plan and DP shards]
B -->|ready batch reference| T[Training ranks]
T -->|training completed| C[Queue controller: progress and ownership]
C --> G[straw online GC]
P[(Shared JuiceFS pack files)] -. payload reads and writes .-> W
P -. shared tensor references .-> B
P -. rank-local reads .-> T
One non-restarting Ray actor owns the coordinator for the run. It serializes queue transitions and dataset allocation; payload writes happen in the producing processes. This permits parallel I/O across machines while retaining one owner for each queue journal. There is no automatic coordinator election or failover.
End-to-end data flow#
Submit prompts. The dataset cursor and prompt tasks commit together. Workers claim leases and load prompt groups from the shared pool. A lost worker does not permanently abandon a preallocated range of dataset indices.
Generate and continue. Distributed fully async generation uses one worker per eligible Ray node. Workers generate, score and filter complete groups; partial groups are saved as durable continuations. R3 routes and SC rows must cover the saved token prefix. An incomplete server capture retries from the previous consistent input.
Accept results. Payloads become durable before their references are accepted. Receipts identify accepted groups, allowing recovery after an RPC reply is lost. Reassigned tasks receive new attempts; stale attempts cannot accept a replacement result.
Build the training batch.
BatchBuilderpersists the selected inputs and conversion configuration before calling reward/conversion hooks. It publishes all DP shards before marking one batch ready. Training readers check batch, plan and rank identity. The manager still materializes the selected samples for conversion, so its host memory remains a capacity consideration.Acknowledge completion. After all rank training calls return, the manager reports
training_completed. This advances runtime consumption and frees ready-batch capacity. It is distinct from a durable model/optimizer checkpoint.
Ray remains responsible for scheduling, RPC and failures of execution processes. vime owns Sample schemas, rewards, filters, batching and training. straw owns bytes, references, persistence and storage lifetime. Custom rollout functions can still return Sample lists; the manager publishes a compatibility collection. See customization.
Queue scheduling#
QueueReader is a worker’s queue client: it acquires and returns groups and renews
active leases. QueueDataSource adds job-level reader creation, consumer lifecycle
and source checkpoints. It inherits the same sample-access methods; it does not
implement another buffer. With --rollout-data-transport straw, the source is
selected automatically. Explicit --data-source-path values should use
vime.data.queue_data_source.QueueDataSource.
QueueReader has no local continuation buffer. add_samples() publishes valid
continuations, then yield_tasks() atomically replaces their inputs and returns
them to pending in one WAL transaction per bounded batch. The old leases end;
any worker can acquire the saved inputs with a new lease. Readers keep only
active leases and executing samples locally. Closing a reader does not rewrite
returned groups.
The persistent order is completed delivery → partial → fresh prompt. Within
a stage, the oldest numeric value in a group’s weight_versions comes first
(higher staleness), then FIFO submission/return order. Groups without numeric versions sort last
within the stage. Comparing stored versions avoids rewriting queue priorities
on every model update. This ordering does not drop stale samples automatically.
Staleness is current serving weight version - oldest generated-token weight version: fresh samples have staleness 0.
There is one indexed queue, not three independently managed queues.
--buffer-filter-path is rejected for straw. --buffer-sort-by-staleness remains
an option for the in-memory data source; straw always applies the order above.
Reward and sample-selection hooks remain available. Already accepted groups
explicitly returned to the source get separate delivery tasks; accepted history
is immutable, and consuming/discarding a delivery also acknowledges its earlier
accepted versions. The fully async scheduler still keeps a bounded ready-results
window, with its own checkpoint state and durable accepted-result references.
Source snapshots now store the shared pending tasks, their exact immutable input
references and scheduling metadata once; worker snapshots store reader metadata.
The reader snapshot format is version 3; old reader-local v1/v2 buffer snapshots
require migration and are rejected explicitly. This requires straw-queue>=0.1.2
with yield_tasks() and persistent scheduling keys.
Sample encoding and integrity#
The adapter stores samples as a vime.v1 JSON record {version: 1, tree: ...}
with explicit tags for containers, Sample fields, NumPy arrays, images and
rollout references. It preserves nested groups, supported dynamic fields,
multimodal inputs, R3, SC top-k and ragged top-p. Unsupported Python types and
cycles fail explicitly. The payload protocol uses no pickle and never imports
classes named by stored data.
Tensor nodes carry shape, dtype, semantic kind, lazy/validated flags and a
{dependency: index, ordinal: index} pointer into the enclosing manifest’s
dependencies. Physical locations belong to those manifests, so publishing
another sample or training batch can share the same immutable tensor records.
All tensor dependencies must be published before the enclosing sample is accepted.
straw’s tensor.v1 stores contiguous, row-major, little-endian typed bytes.
Lazy readers verify metadata and the checksum of each chunk covering requested
rows; they do not verify untouched chunks. Full inspection verifies the whole
segment. Tensor descriptors are restored in batches to reuse shared pack indices.
Enable it#
requirements.txt includes straw-queue>=0.1.2, so the standard Docker build and
normal vime installation install it automatically. For an existing environment,
install it on every generation and training node, using the same version throughout the job:
pip install 'straw-queue>=0.1.2'
Without straw installed, the default Ray object-store transport still works
and startup logs the installation command. Selecting --rollout-data-transport straw
without the package fails at startup with the same command.
The Python import name is straw. Add these arguments to your existing training
command, using a fresh shared directory:
--rollout-data-transport straw \
--rollout-data-dir /shared/juicefs/jobs/my-run/rollout_data \
--rollout-queue-run-id my-run \
--rollout-storage-profile juicefs \
--rollout-storage-declaration /shared/juicefs/deployment.json
The declaration records direct_mount: true, writeback: false,
open_cache: 0, readdir_cache: false, the deployed client_version and a
durability_description. It declares verified deployment settings; it does
not discover or configure the mount. See straw’s
filesystem contract.
The default local storage profile represents the POSIX development contract;
using it on a shared mount does not verify that service’s durability.
Storage and execution mode are separate choices. The default transport remains
object-store; the default rollout entrypoint remains synchronous. To select
distributed fully async execution with straw, also add:
--rollout-function-path vime.rollout.fully_async_rollout.generate_rollout_fully_async
If --rollout-data-dir is omitted, straw mode uses <save>/rollout_data; without
--save, an explicit shared directory is required. --rollout-io-concurrency
bounds off-event-loop serialization and I/O (default 4). The adapter uses worker
processes and bounded threads; straw does not create an I/O process pool.
During weight synchronization, distributed producers pause new admissions and allow in-flight work to finish and persist. They resume if another training rollout remains. Weight synchronization itself remains a separate subsystem. On shutdown, all generation workers share a five-minute deadline to finish durable writes. Exceeding it reports an error rather than successful shutdown.
Online GC and ownership#
Online GC is disabled by default. Enable --rollout-queue-online-gc to let straw
reclaim unused sealed packs during a job. vime provides completion and discard
signals through straw APIs. A pack is deletable only after all task, publication,
queue, reader and checkpoint owners release it. A lease timeout alone does not
prove that a reader has stopped. A live record retains its entire pack; active
writers, retained checkpoints and WAL history still consume space.
GC failures stop the background loop, propagate the original cause through later coordinator operations and surface at shutdown. They do not stop work already running on remote workers. Retained data stays available for inspection. Disk exhaustion requires freeing unneeded retention or expanding storage; do not reset a live queue or unlink pack files to make room. Offline removal requires all participants to stop. Storage quotas and application retention policies remain necessary.
Processed positions and filter decisions are recorded explicitly, so gaps in acceptance order do not retain already-consumed data.
Recovery and checkpoints#
Save model and rollout state together through normal training checkpoints. They include the dataset cursor, sample/group counters, pending and partial inputs, ready groups, and training progress. Existing straw payloads are referenced rather than copied; retaining a checkpoint keeps its data available.
Stop the previous job, including remote workers, before restarting. The save-directory lock rejects concurrent coordinators but does not stop orphaned readers; there is no automatic coordinator failover.
Resume or select a step#
Resume the latest completed checkpoint using the same logical directories:
--rollout-data-transport straw \
--load /shared/checkpoints/run \
--save /shared/checkpoints/run \
--save-interval 1
Megatron requires a positive --save-interval when --save is set. To restore
the state saved after rollout 7, add --ckpt-step 7; the next rollout is 8.
Keep the dataset, model/tokenizer configuration, straw run ID, storage profile,
and fully async worker topology consistent. Training recovery requires saved
optimizer and RNG state.
Restore creates an isolated queue that shares the checkpoint’s immutable
payloads. Pending/partial inputs, ready-group order, and dataset cursor come
from that checkpoint; later samples from the source run are excluded. Further
writes do not modify the source checkpoint. If --save already contains a run,
outputs go to a unique branches/<id> directory tracked by
rollout/current.json. A later restart can use the same logical --load and
--save; omitting --load resumes the active branch. Use a fresh path for each
immutable .straw.json debug archive.
Selection |
Arguments |
|---|---|
Step in current branch history |
|
A specific branch |
|
Exact checkpoint |
|
Separate output directory |
|
Editing latest_checkpointed_iteration.txt in the logical save directory or
active branch also selects an earlier step. Explicit --ckpt-step or an exact
commit file takes precedence. Automatic selection uses completed model-and-queue
checkpoints, follows branch ancestry only to each fork point, and ignores newer
incomplete model saves.
If a model checkpoint has no queue snapshot, recovery starts an empty queue and
restores rollout/global_dataset_state_dict_<step>.pt when present; otherwise
the dataset restarts at offset 0 with a warning. Missing models, incomplete or
corrupt snapshots, and missing payloads raise instead of silently falling back.
Before the first model checkpoint, the original run can recover persisted work
only before its first training batch is planned. Recovery uses the original
straw pool; copying the checkpoint directory alone does not copy payloads. The
boundary is a completed rollout batch, not an optimizer microstep. GPU KV
caches and generation RNG are not restored. When extending --num-rollout,
use --use-checkpoint-opt-param-scheduler to retain the saved optimizer schedule.
Debug archives and sample lookup#
.pt and indexed .straw.json archives both work with the existing debug flags:
--save-debug-rollout-data '/shared/debug/rollout_{rollout_id}.straw.json'
# In a separate train-only job, without vLLM:
--load-debug-rollout-data '/shared/debug/rollout_{rollout_id}.straw.json'
An archive retains its data independently of queue consumption and GC. With
straw transport it reuses existing tensors; otherwise it creates a straw-data
pool beside the index. Each rollout has one immutable index, with samples
stored in chunks. Evaluation uses eval_<id>. Copying only the JSON index does
not copy payloads. Train-only straw replay uses the archive’s writable storage
pool and run, overriding --rollout-data-dir and --rollout-queue-run-id.
Each replay creates an isolated queue without changing the source queue. Full
replay without subsampling reuses the archived Sample and tensor records; all
archives in one replay must belong to the same pool and run.
--load-debug-rollout-data-subsample also applies to archives.
from vime.data.archive import RolloutArchive
from vime.observability.rollout_data_utils import load_debug_rollout_data
with RolloutArchive('/shared/debug/rollout_7.straw.json') as archive:
print(archive.keys())
samples = archive.load_samples(sample_key='sample:42')
group = archive.load_samples(task_key='prompt:21')
archive.export_pt('/shared/debug/rollout_7.pt')
samples = load_debug_rollout_data('/shared/debug/rollout_7.pt', rollout_id=7)
RolloutArchive.save('/shared/debug/imported_7.straw.json', samples, rollout_id=7)
Lookups return lists because compact trajectories can share a sample index.
Samples without an index use their archive position; samples without queue
provenance have no task key. Closing an archive closes its reader but retains
the data. After all readers finish, call archive.release(); other checkpoints,
queues, and archives keep their own references. An exported .pt remains
readable after straw reclaims the archive payloads.
Validation and remaining limits#
CPU coverage includes codecs, R3/SC transport, lazy CP/TP reads, worker failures,
continuations, checkpoint views and GC ownership. straw separately checks its
Rust/Python core, bounded ownership model and native crash boundaries. Physical
multi-client filesystem checks and real GPU training complement these tests.
tests/test_straw_fully_async_recovery.py starts two local Ray nodes, interrupts
the entire driver/coordinator/generation process tree with SIGKILL, and starts a
new job against the same straw pool without a graceful checkpoint. It checks
accepted-result replay, unchanged token/logprob/R3/SC prefixes, new leases and
single reward evaluation, with online GC both disabled and enabled. Inference
and rewards are CPU fixtures; this test does not claim optimizer recovery.
It runs automatically in cpu-unittest; see CI setup.
Run the distributed rollout and interruption recovery tests locally with:
PYTHONPATH=. python tests/test_distributed_rollout.py
PYTHONPATH=. python tests/test_straw_fully_async_recovery.py
The main deployment limits are shared-storage bandwidth, durability latency, coordinator throughput and manager conversion memory. Live-pack compaction, journal compaction and automatic coordinator failover are not implemented. See usage for the rollout controls and straw’s verification guide for the storage checks and their scope.