Repository navigation
[skycap] Mirror records to remote storage, and say where each is at finish - #2350
Conversation
There was a problem hiding this comment.
Code Review
This pull request introduces a background record mirroring feature using fsspec to copy written trajectory records to remote storage. A critical issue was identified in the shutdown logic of the RecordMirror where worker threads can leak and block indefinitely on self._queue.get() if they do not join before the shutdown deadline. It is recommended to implement a timeout on the queue fetch and check the stopping event to ensure clean thread termination.
|
…inish
A capture server can copy every record it writes to a remote store (any
fsspec URL: s3://, gs://, ...) in the background. The local record_dir stays
the source of truth; the copy fails open.
- `--record-mirror URL` and `--record-mirror-config JSON` (`record_mirror`,
`record_mirror_config` on CaptureServer/CaptureService). The config holds the
mirror's options: exclude, workers, queue_size, timeout, attempts, backoff,
shutdown_timeout, storage_options. fsspec comes with the optional `remote`
extra and is imported only when a mirror is made.
- The mirror (`skycap.mirror.RecordMirror`): a bounded queue, worker threads,
a per-file timeout after which a copy is abandoned rather than retried,
bounded retries for transient errors, and a deadline at shutdown. Every loss
is logged and counted on /healthz.
- `exclude` leaves sidecar kinds out of the copy; the mirrored document names
them in `omitted_sidecars` and reads as a whole record.
- finish returns `record: {path, mirror, files}` (`FinishResult.record`,
`RecordLocation`).
- format.md: sidecars are optional and the manifest lists exactly what a copy
has; the mirror, `omitted_sidecars`, finish's `record`, and an optional run
index (`index/<phase>/step-<N>.json`).
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com>
b1bb140 to
dc3a553
Compare
There was a problem hiding this comment.
Cursor Bugbot has reviewed your changes and found 1 potential issue.
❌ Bugbot Autofix is OFF. To automatically fix reported issues with cloud agents, have a team admin enable autofix in the Cursor dashboard.
Reviewed by Cursor Bugbot for commit dc3a553. Configure here.
…omitted_sidecars - finish's `record` gains `host`: the IP address (or hostname) of the machine `path` is on, so another node can reach a record that exists only there. `RecordLocation.local` is `host:path` (an IPv6 host bracketed), and `uri` is the mirror's URI or else `local`. A server reports `record_host` when given (`--record-host` on the CLI), a service its `advertise_host` unless that's loopback, and otherwise the machine's primary IP. - A mirror with `exclude` now copies the document unchanged and skips the excluded sidecar files; nothing is rewritten and there is no `omitted_sidecars`. A copy may lack sidecar files its manifest lists, and a reader treats those as absent (format.md says so); `record.load` does, with a warning, instead of asserting. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com>
…hout omitted_sidecars - Each index row's `record` includes `host`, the machine `path` is on, as finish now reports it, so a head node can reach a local-only record. - The mirror test checks that `exclude` leaves the sidecar out and copies the document unchanged (#2350 no longer writes `omitted_sidecars`). - Tests pass `records=` by keyword: main added `train_paths` before it. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com>
kouroshHakha
left a comment
There was a problem hiding this comment.
leaving comments
… and an exact pending count - `submit` checks `_closed` and queues under the lock `close` sets it with, so a record is either queued before close (and handled or counted by it) or refused and counted as dropped; never accepted and then lost. - Workers wait on the queue with a timeout and exit once the mirror is stopping, so one still copying at close's deadline (whose stop signal close drained) finishes its copy and exits instead of blocking forever. - `pending` is a counter raised when a record is queued and lowered when it is settled (mirrored, failed or dropped), together with that outcome's count, so it never reads 0 while a record is between the queue and its copy. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com>
… Harbor integration skycap writes records; which trajectories made up a training step is the trainer integration's knowledge, so its index file is specified there (#2351), not in skycap's format.md or README. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com>
…hout omitted_sidecars - Each index row's `record` includes `host`, the machine `path` is on, as finish now reports it, so a head node can reach a local-only record. - The mirror test checks that `exclude` leaves the sidecar out and copies the document unchanged (#2350 no longer writes `omitted_sidecars`). - Tests pass `records=` by keyword: main added `train_paths` before it. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com>
…_index.md It moved out of skycap's format.md (#2350 review): skycap writes records, and which trajectories made up a step is this integration's knowledge. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com>
**TLDR:** each training step's skycap records are indexed in W&B: one artifact version per step listing which trajectories made up the step, and how each one ended up. Mirrored records (#2350) appear as W&B *references*, so no record bytes are uploaded. It replaces #2333's upload path, and its design (config shape, per-phase log, eval opt-in) comes from that PR. --------- Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com> Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
**TLDR:** `record_index pull` turns a run's W&B record index (#2351) back into a local record directory: the step's records, read from the mirror (#2350), plus the step index. Stacked on #2351. --------- Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com> Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
…ull writes a directory per phase (#2441) **TLDR:** lands #2351 (the per-step W&B record index) and #2436 (`record_index pull`) on `main`. Both were reviewed and merged, but into their stacked base branches after #2350 had already landed, so neither reached `main`. On top of that, `pull` writes each phase as a record directory of its own. --------- Signed-off-by: Kourosh Hakhamaneshi <kourosh@anyscale.com> Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>

What does this PR do?
TLDR: a capture server can copy each record it writes to remote storage (any fsspec URL:
s3://,gs://, ...) in the background, andfinishsays where each record is. The localrecord_dirstays the source of truth; the copy fails open. Splits the storage half out of #2333; the W&B half is #2351.How it works
Options:
--record-mirror URLand--record-mirror-config JSON, orrecord_mirrorandrecord_mirror_configonCaptureServerandCaptureService. Both needrecord_dir.skycap serve ... --record-dir ./record --record-mirror s3://bucket/run-7 \ --record-mirror-config '{"exclude": ["experts", "sampling_mask"], "timeout": 120}'The config keys are
exclude,workers,queue_size,timeout,attempts,backoff,shutdown_timeoutandstorage_options. fsspec comes from the new optionalremoteextra and is imported only when a mirror is made.RecordMirror(skycap/mirror.py): worker threads copy a record's sidecars, then its document. It fails open:Every loss is logged and counted on
/healthzunderrecord_mirror.excludeskips sidecar kinds when copying; the document is copied unchanged. A copy may therefore lack sidecar files its manifest lists, and readers treat those as absent: withouttokensthey show message text only, and they need nothing else, since only training usesexpertsandsampling_mask.finishreturnsrecord: {host, path, mirror, files}, also available asFinishResult.record(aRecordLocation).hostis the machinepathis on, so another node can reach a record that exists only there.RecordLocation.localgiveshost:path, the scp form.hostcomes from--record-hostif set, otherwise from a service'sadvertise_hostunless that's loopback, otherwise from the machine's primary IP.fileslists the record's files as the mirror holds them, or as the record directory does without a mirror.format.md:record.Test plan
skycap/tests/test_mirror.py: copies are byte-for-byte, sidecars first; a failing or slow store fails open; transient errors are retried and permanent ones aren't; a sidecar failure keeps the document out; a timeout is abandoned without a duplicate; a full queue drops; shutdown respects its deadline;exclude(an unchanged document, and missing sidecars reading as absent);record.host; the config over the CLI and the API;finish'srecord.files.--extra tokens --extra remote(230) and on a clean install (224 passed, 1 skipped).test_imports.pychecks that fsspec isn't loaded without a mirror.uv lock --checkand pre-commit are clean.🤖 Generated with Claude Code
Note
Medium Risk
New async upload path and changed finish API/record read semantics; failures are isolated from trajectories but remote copies may be incomplete or missing sidecars by design.
Overview
Adds background mirroring of finished trajectory records to any fsspec URL (
s3://,gs://, …) via--record-mirror/CaptureService(record_mirror=…)and the new optionalremoteextra (fsspec). Local--record-dirstays authoritative;RecordMirrorcopies sidecars then the document on worker threads with a bounded queue, per-file timeouts, transient retries, and fail-open behavior (drops/timeouts never failfinish)./healthzexposesrecord_mirrorcounters; shutdown drains the queue up to a deadline.finishnow returnsrecord(RecordLocation:host,path, optionalmirror,files) so trainers can fetch from the node or remote URI;--record-hostoverrides the reported address. Mirrorexcludeskips sidecar kinds in the upload while leaving the JSON manifest unchanged;record.loadandformat.mdtreat missing listed sidecars as absent (e.g. text-only whentokensis excluded).CLI validation ties mirror to
record_dir; extensive coverage intest_mirror.py.Reviewed by Cursor Bugbot for commit 1a29349. Bugbot is set up for automated code reviews on this repo. Configure here.