Skip to content
Merged
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: 43 additions & 1 deletion skycap/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -241,11 +241,53 @@ and each token's offset in it), routed experts and sampling masks. The format
is specified in [`docs/format.md`](docs/format.md), which is what any reader,
such as the viewer, implements.

`result.record` says where it is: `path` on the server's disk and, with a
mirror, its `mirror` URI (`result.record.uri` is the mirror URI when there is
one, else the path). `result.record.files` names the record's files, sidecars
then document, all beside the document: the ones the mirror holds when there is
one (without the kinds its `exclude` leaves out), else the ones on disk.

### Mirror the record to remote storage

`--record-mirror URL` (`record_mirror=` for `CaptureService`) copies each
record, once written to `--record-dir`, to an fsspec URL in the background.
Install the `remote` extra and the store's fsspec implementation yourself:

```bash
uv sync --extra remote && uv pip install s3fs # or gcsfs for gs://
uv run skycap serve ... --record-dir ./record --record-mirror s3://bucket/run-7
```

The mirror never fails a trajectory. Its queue is bounded (a record that finds
it full is dropped), each file's copy is abandoned after a timeout rather than
retried, errors that may pass (connection resets, timeouts) are retried up to
3 times with backoff, and a graceful shutdown waits up to a deadline for the
queue and drops the rest. Each loss is logged and counted: `/healthz` reports
`record_mirror` with `mirrored`, `failed`, `timed_out`, `dropped`, `retried`
and `pending`. Credentials come from the store's usual environment (e.g.
`AWS_*` for s3fs).

`--record-mirror-config JSON` (`record_mirror_config=` for `CaptureService`)
sets the mirror's options: `exclude`, `workers`, `queue_size`, `timeout`,
`attempts`, `backoff`, `shutdown_timeout` and `storage_options`. `exclude`
leaves sidecar kinds out of the remote copy, e.g. routed experts when the
mirror is for reading rather than retraining:

```bash
uv run skycap serve ... --record-dir ./record --record-mirror s3://bucket/run-7 \
--record-mirror-config '{"exclude": ["experts", "sampling_mask"], "timeout": 120}'
```

The document is copied unchanged, so in the mirror it lists sidecars the mirror
doesn't hold; readers treat those as absent (see [`docs/format.md`](docs/format.md)).
The local record is untouched. Excluding `tokens` is allowed, but a viewer of
the mirror then shows message text only.

## Develop

```bash
cd skycap
uv sync --extra tokens
uv sync --extra tokens --extra remote
uv run pytest
```

Expand Down
74 changes: 66 additions & 8 deletions skycap/docs/format.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,16 @@ no Python and no tokenizer.

A record directory holds, per trajectory `{id}`:

| File | Always Outputted | Holds |
| File | Present | Holds |
| --- | --- | --- |
| `{id}.json.zst` | yes | the document (graph) |
| `{id}.tokens.zst` | token mode, when any node has tokens | token ids, logprobs, the text the tokens decode to, and each token's byte offset in it |
| `{id}.experts.zst` | when routed experts were captured | routed experts (R3) |
| `{id}.sampling_mask.zst` | when a captured sampling mask has at least one row | per sampled token, the ids it could have been drawn from |
| `{id}.json.zst` | always | the document (graph) |
| `{id}.tokens.zst` | optional: token mode only, when any node has tokens; never in text mode | token ids, logprobs, the text the tokens decode to, and each token's byte offset in it |
| `{id}.experts.zst` | optional: when routed experts were captured (R3) | routed experts |
| `{id}.sampling_mask.zst` | optional: when a captured sampling mask has at least one row | per sampled token, the ids it could have been drawn from |

Every sidecar is optional. The document's `sidecars` manifest lists the ones
the trajectory captured, and a kind missing from it is one the trajectory
doesn't have.

Every file is exactly one zstd frame: compress the whole payload in one call
and write it once. Never add to a file that already exists, whether by
Expand All @@ -27,11 +31,65 @@ A trajectory is not partially written when it is live, it is only written when i
Each document has a `sidecars` field naming the sidecar files that belong to
it. The writer writes those sidecars before the document, and writes every file
under a temporary name and then renames it, so no file is ever seen
half-written. If a document exists, every sidecar named in its `sidecars` field
exists too. A crash can leave sidecars with no document, but never a document
with a missing sidecar.
half-written. In a record directory, if a document exists, every sidecar named
in its `sidecars` field exists too. A crash can leave sidecars with no
document, but never a document with a missing sidecar.

A copy of a record directory, such as a mirror that leaves sidecar kinds out,
may lack sidecar files its manifest lists. A reader treats a listed sidecar
whose file is missing as absent, as if the trajectory hadn't captured it, and
doesn't tell it apart from one that went missing: without `tokens` it has
message text only and should say so; without `experts` or `sampling_mask` it
loses nothing a reader shows, since only training uses them, and training
takes its samples from `finish`.
A reader lists trajectories by listing `*.json.zst`.

### A mirror

A server started with a record mirror (`--record-mirror URL`, an fsspec URL
such as `s3://bucket/run-7`) also copies each trajectory's files, under the
same names, to `{URL}/{name}`, after writing them to the record directory.
The copy is made in the background, sidecars before the document, so a
mirrored document also has its sidecars beside it. The copy fails open: a
record the store never got is missing from the mirror (the server's
`/healthz` counts them under `record_mirror`), and the record directory is
always complete. A reader reads a mirror exactly as it reads a record
directory.

A mirror can leave sidecar kinds out (`--record-mirror-config '{"exclude":
["experts"]}'`), e.g. when the remote copy is for reading rather than
retraining. It still copies the document unchanged, so the document lists
sidecars the mirror doesn't hold, which a reader reads as absent (above). The
record directory's copy is untouched. Leaving out `tokens` is allowed, but a
reader of the mirror then has message text only.

### Where a trajectory's record is

`finish` answers with the document's location as `record`, or `null` when the
server has no record directory or couldn't write the record:

```json
{"record": {"host": "10.0.0.5",
"path": "/data/record/tr_ab12.json.zst",
"mirror": "s3://bucket/run-7/tr_ab12.json.zst",
"files": ["tr_ab12.tokens.zst", "tr_ab12.json.zst"]}}
```

`path` is on the server's own disk, and `host` is that machine: its IP address
(or hostname), so another node can reach a record that exists only there, e.g.
`scp 10.0.0.5:/data/record/tr_ab12.json.zst .` (an IPv6 host is bracketed
there, `[fe80::1]:/data/...`). The server reports the address it is reached
at, or its primary IP when it has none to report; `--record-host` sets it. `mirror` is null without a mirror, and
otherwise where the copy is going: it may not be there yet, or at all. The
sidecars are beside the document in both places.

`files` names the record's files, sidecars first and the document last. With a
mirror they are the files the mirror holds (or will hold), so the kinds its
`exclude` leaves out are not listed; without one they are the files in the
record directory. A sidecar the trajectory didn't capture is never listed.
Each file is beside the document: `{dirname(mirror)}/{name}`, or
`{dirname(path)}/{name}` without a mirror.

## The document

The decompressed document is a UTF-8 JSON object:
Expand Down
6 changes: 6 additions & 0 deletions skycap/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,11 @@ tokens = [
"renderers>=0.1.9.dev9",
"transformers>=4.50",
]
# The record mirror (--record-mirror). A store fsspec doesn't ship needs its own
# implementation as well, e.g. s3fs for s3:// or gcsfs for gs://.
remote = [
"fsspec>=2024.2",
]

[project.scripts]
skycap = "skycap.cli:main"
Expand All @@ -27,6 +32,7 @@ Repository = "https://github.com/NovaSky-AI/SkyRL"

[dependency-groups]
dev = [
"fsspec>=2024.2",
"openai>=1.40",
"pytest>=8.2",
"pytest-asyncio>=0.23",
Expand Down
2 changes: 2 additions & 0 deletions skycap/src/skycap/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
CapturePool,
FinishResult,
PathRuleError,
RecordLocation,
Trajectory,
)
from skycap.samples import Sample # noqa: E402
Expand All @@ -18,6 +19,7 @@
"CaptureService",
"FinishResult",
"PathRuleError",
"RecordLocation",
"Sample",
"Trajectory",
"__version__",
Expand Down
38 changes: 37 additions & 1 deletion skycap/src/skycap/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,29 @@ def build_parser() -> argparse.ArgumentParser:
default=None,
help="where ended trajectories are written; without it they stay in memory (development only)",
)
serve.add_argument(
"--record-mirror",
default=None,
metavar="URL",
help="also copy each written record to this fsspec URL (s3://, gs://, ...) in the background; "
"needs --record-dir and skycap[remote], plus the store's fsspec implementation (s3fs, gcsfs)",
)
serve.add_argument(
"--record-mirror-config",
type=json.loads,
default=None,
metavar="JSON",
help='the mirror\'s options as a JSON object, e.g. \'{"exclude": ["experts", "sampling_mask"]}\'; '
"keys: exclude, workers, queue_size, timeout, attempts, backoff, shutdown_timeout, storage_options. "
"Excluding tokens leaves viewers of the mirror with message text only",
)
serve.add_argument(
"--record-host",
default=None,
metavar="ADDRESS",
help="the address other machines reach this one at, reported as finish's record.host so they can "
"fetch a record from this node's --record-dir (default: this machine's primary IP)",
)
serve.add_argument(
"--ttl",
type=float,
Expand Down Expand Up @@ -112,6 +135,12 @@ def build_parser() -> argparse.ArgumentParser:
def build_server(args: argparse.Namespace) -> CaptureServer:
if args.mode == "tokens" and not args.tokenizer:
raise SystemExit("--mode tokens needs --tokenizer")
if args.record_mirror and not args.record_dir:
raise SystemExit("--record-mirror needs --record-dir")
if args.record_mirror_config is not None and not args.record_mirror:
raise SystemExit("--record-mirror-config needs --record-mirror")
if args.record_mirror_config is not None and not isinstance(args.record_mirror_config, dict):
raise SystemExit("--record-mirror-config must be a JSON object")
backend = build_backend(
args.upstream_url,
mode=args.mode,
Expand All @@ -130,7 +159,14 @@ def build_server(args: argparse.Namespace) -> CaptureServer:
name, _, rule = spec.rpartition("=")
rules[name or rule] = rule
return CaptureServer(
backend, record_dir=args.record_dir, ttl=args.ttl, path_rules=rules, require_api_key=args.require_api_key
backend,
record_dir=args.record_dir,
record_mirror=args.record_mirror,
record_mirror_config=args.record_mirror_config,
record_host=args.record_host,
ttl=args.ttl,
path_rules=rules,
require_api_key=args.require_api_key,
)


Expand Down
45 changes: 45 additions & 0 deletions skycap/src/skycap/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,48 @@ class PathRuleError(CaptureError):
"""``finish``'s path rule raised on the server. The trajectory is ended, without samples."""


@dataclass(frozen=True, slots=True)
class RecordLocation:
"""Where a finished trajectory's record is. See ``docs/format.md``."""

#: The document's path on the capture server's own disk, in its ``record_dir``.
path: str
#: The document's URI in the server's record mirror, or None without one. The copy is made in the
#: background and fails open, so it may not be there yet, or at all.
mirror: str | None = None
#: The record's file names, sidecars then the document (e.g. ``tr_ab12.tokens.zst``, ``tr_ab12.json.zst``):
#: as the mirror holds them when there is one (so without the kinds its ``exclude`` leaves out), else as
#: the record directory does. They sit beside the document, in the mirror and on disk.
files: tuple[str, ...] = ()
#: The machine ``path`` is on: the capture server's node, as an IP address (or hostname). None from a
#: server that predates it.
host: str | None = None

@property
def local(self) -> str:
"""``host:path``, the way scp and rsync over ssh name a remote file; an IPv6 host is bracketed."""
if self.host is None:
return self.path
host = f"[{self.host}]" if ":" in self.host else self.host
return f"{host}:{self.path}"

@property
def uri(self) -> str:
"""The mirror URI when there is one, else ``local``."""
return self.mirror or self.local

@classmethod
def from_json(cls, body: dict[str, Any] | None) -> RecordLocation | None:
if body is None:
return None
return cls(
path=body["path"],
mirror=body.get("mirror"),
files=tuple(body.get("files", ())),
host=body.get("host"),
)


@dataclass(slots=True)
class FinishResult:
id: str
Expand All @@ -52,6 +94,8 @@ class FinishResult:
#: Token-mode calls whose prompt had to be rendered rather than extended (``CallInfo.bridged``).
#: Zero for a harness that keeps its history append-only.
unbridged_calls: int = 0
#: Where the record is, or None when the server has no ``record_dir`` (or couldn't write it).
record: RecordLocation | None = None


class Trajectory:
Expand Down Expand Up @@ -98,6 +142,7 @@ async def finish(self, annotations: dict[str, Any] | None = None, *, paths: str
status=body["status"],
samples=[Sample.from_json(s) for s in body["samples"]],
unbridged_calls=body.get("unbridged_calls", 0),
record=RecordLocation.from_json(body.get("record")),
)
return self.result

Expand Down
Loading
Loading