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
File renamed without changes.
1 change: 1 addition & 0 deletions docs/design.md
100644 → 100755
Original file line number Diff line number Diff line change
Expand Up @@ -315,6 +315,7 @@ revenant status # print each stage's offset / running-or-idl
- **Adding new input mid-run.** Explicitly deferred — input is currently treated as a fixed, fully-assembled batch known before the pipeline starts (§1, §7). If this changes later, the fix is a single `input.closed.json` marker (`{"final_seq": N, "closed_at": "..."}`) written explicitly whenever input is done being appended to, which becomes the new base case for §7's recursion; every other stage boundary is unaffected. Two operating modes would fall out of this: batch mode (close input immediately, run to drain) vs. long-running/streaming mode (leave input open, keep enqueuing, only close it to deliberately shut the pipeline down for good).
- **Snapshot frequency.** Confirmed: checkpoint after *every* completed item (not batched every N items), since tasks may be long-running and the cost of reprocessing more than one item on crash was judged not worth the savings from less-frequent snapshotting.
- **Log rotation / compaction.** Not yet designed. `A.jsonl` grows unboundedly for the life of a run; if a stage's output file needs periodic compaction (e.g. dropping already-fully-consumed lines), note that `seq` numbers — not line-position-in-file — are the resumption unit specifically so this remains possible without breaking downstream consumers.
- **Blob storage garbage collection.** Deferred, same bucket as the log rotation/compaction item above — see `docs/design-blob-storage.md` for the full blob-storage design and the one-hop rule that keeps a future GC pass tractable.
- **Retry backoff policy for `RetryableError`.** Max attempts / backoff schedule not yet implemented; current behavior is an unconditional retry with the fixed poll-interval delay already used by the runner.
- **Docker bind-mount atomicity verification.** Confirm the actual host filesystem backing the bind mount before relying on `write()`-append and `rename()` atomicity (§6's caveat) — this is load-bearing for the entire crash-safety story.

Expand Down
Empty file.
47 changes: 47 additions & 0 deletions examples/blob_split_merge/make_input.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
"""
Seed state/input.jsonl for the blob_split_merge example.

This script creates a tiny input stream for the blob split/merge
pipeline by writing one logical record that points to a sample file.
Run from the repo root:
python examples/blob_split_merge/make_input.py
"""

from __future__ import annotations

from pathlib import Path

from revenant.io_utils import atomic_append_lines, make_checkpoint_line, make_record_line

LINES = [
"README.md",
]

def main() -> None:
state_dir = Path("state")
state_dir.mkdir(exist_ok=True)
input_path = state_dir / "input.jsonl"
if input_path.exists():
input_path.unlink()

# Create one logical input record for the sample file.
records = [
make_record_line(
seq=i,
src_seq=i,
parent_seq=i,
payload={"input": text},
)
for i, text in enumerate(LINES, start=1)
]
# Treat the generated input as already committed so downstream
# stages can consume it immediately. The checkpoint line marks the
# latest visible state for tooling that reads the last committed
# record (see docs/design.md, section 5.1).
records.append(make_checkpoint_line(last_consumed_seq=len(LINES), last_emitted_seq=len(LINES)))
atomic_append_lines(input_path, records)
print(f"Wrote {len(LINES)} input records to {input_path}")


if __name__ == "__main__":
main()
64 changes: 64 additions & 0 deletions examples/blob_split_merge/pipeline.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
"""
Example pipeline: blob split/merge.

Demonstrates:
- a stateless stage that fans out one input file into many blob-backed
parts (SplitFile: one file -> many part records)
- a stateful stage that reassembles those parts once all fragments
are present (MergeFile: collect parts until the full file is ready)

This pipeline is a reference for the Step API shape (docs/design.md,
section 8) and the blob-store integration pattern, not a performance
example.
"""

from __future__ import annotations

from revenant.config import StageConfig
from revenant.step import Step


class SplitFile(Step):
"""Stateless: split one input file into blob-backed parts."""

def process(self, payload, state):
input = payload.get("input", "")
with open(input) as f:
lines = [l for l in f]
total = len(lines)
for l in lines:
# Each line becomes a separate blob-backed part that can be
# merged later once the full input has arrived.
part = self.blob_store.write(bytes(l, 'utf8'))
yield {'input': input, 'part': part, 'total': total}
return state

class MergeFile(Step):
"""Stateful: collect blob parts until the file can be reassembled."""

def process(self, payload, state):
input = payload.get("input")
part = payload.get("part")
total = payload.get("total")
state = state or {"parts": {}}
parts = state['parts'].setdefault(input, [])
parts.append(part)
state["parts"][input] = parts
if len(parts) < total:
# Keep gathering fragments for this input until all expected
# parts have arrived.
return state
# Once every fragment is present, read them back from blob storage
# and emit the merged blob for the original input.
texts = [self.blob_store.resolve(p).read_text() for p in parts]
text = "".join(texts)
result = self.blob_store.write(bytes(text, 'utf8'))
yield {"input": input, "total":total, "result":result}
del(state["parts"][input])
return state


PIPELINE = [
StageConfig(name="split", step_class=SplitFile, upstream="input"),
StageConfig(name="merge", step_class=MergeFile, upstream="split"),
]
1 change: 1 addition & 0 deletions pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -28,3 +28,4 @@ where = ["src"]

[tool.pytest.ini_options]
testpaths = ["tests"]
pythonpath = ["."]
88 changes: 88 additions & 0 deletions src/revenant/blob_store.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
"""
Content-addressed blob storage for stage outputs that shouldn't be
embedded directly in .jsonl payloads (e.g. images, media files). See
docs/design-blob-storage.md for the full rationale.

One BlobStore per stage process, constructed once at stage-runner
startup and attached to the Step instance as `step.blob_store`.
"""

from __future__ import annotations

import hashlib
import os
import uuid
from pathlib import Path


class BlobStore:
def __init__(self, blobs_dir: Path, state_dir: Path):
self.blobs_dir = blobs_dir
self.state_dir = state_dir
self.scratch_dir = blobs_dir / ".scratch"
self.blobs_dir.mkdir(parents=True, exist_ok=True)
self.scratch_dir.mkdir(parents=True, exist_ok=True)

def reserve_path(self, suffix: str = "") -> Path:
"""Return a scratch path an external tool can write to directly."""
return self.scratch_dir / f"{uuid.uuid4().hex}{suffix}"

def reserve_dir(self) -> Path:
"""Same idea, for tools that produce many output files at once
(e.g. ffmpeg segmenting one input into N chunk files)."""
d = self.scratch_dir / uuid.uuid4().hex
d.mkdir(parents=True, exist_ok=True)
return d

def write(self, data: bytes, suffix: str = "") -> str:
"""Convenience wrapper: write in-memory bytes and commit in one call."""
scratch_path = self.reserve_path(suffix=suffix)
fd = os.open(scratch_path, os.O_WRONLY | os.O_CREAT, 0o644)
try:
os.write(fd, data)
os.fsync(fd)
finally:
os.close(fd)
return self.commit(scratch_path, suffix=suffix)

def commit(self, scratch_path: Path, suffix: str | None = None) -> str:
"""Hash a completed scratch file and move it into content-addressed
storage. Returns a path relative to state_dir, suitable for storing
directly in a payload dict.

`scratch_path` must already be fully written and fsync'd before
calling this -- commit() does not fsync the source file itself,
only the containing-directory fsync implied by os.replace on most
Linux filesystems. Callers writing via write() get this for free;
callers using reserve_path()/reserve_dir() directly (e.g. an
external subprocess) are responsible for the tool having finished
and flushed before commit() is called.
"""
digest = _sha256_file(scratch_path)
actual_suffix = suffix if suffix is not None else scratch_path.suffix
shard = digest[:2]
name = f"{digest}{actual_suffix}"
target_dir = self.blobs_dir / shard
target_dir.mkdir(parents=True, exist_ok=True)
final_path = target_dir / name

if final_path.exists():
# Identical content already stored (e.g. a deterministic retry
# after a crash) -- discard the duplicate scratch file.
scratch_path.unlink(missing_ok=True)
else:
os.replace(scratch_path, final_path) # atomic: same filesystem

return str(final_path.relative_to(self.state_dir))

def resolve(self, relative_path: str) -> Path:
"""Turn a payload's stored relative path back into a real path."""
return self.state_dir / relative_path


def _sha256_file(path: Path) -> str:
h = hashlib.sha256()
with open(path, "rb") as f:
for chunk in iter(lambda: f.read(1024 * 1024), b""):
h.update(chunk)
return h.hexdigest()
3 changes: 3 additions & 0 deletions src/revenant/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,3 +69,6 @@ def upstream_output_path(self, state_dir: Path) -> Path:
if self.upstream == "input":
return state_dir / "input.jsonl"
return state_dir / f"{self.upstream}.jsonl"

def blobs_dir(self, state_dir: Path) -> Path:
return state_dir / f"{self.name}.blobs"
3 changes: 3 additions & 0 deletions src/revenant/stage_runner.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@
read_last_checkpoint_line,
)
from revenant.step import RetryableError, SkipItem
from revenant.blob_store import BlobStore

POLL_INTERVAL_SECONDS = 1.0

Expand Down Expand Up @@ -137,6 +138,8 @@ def run_stage(
last_consumed_seq, last_emitted_seq, saved_state = load_resume_point(stage, state_dir)

step = stage.step_class()
step.blob_store = BlobStore(stage.blobs_dir(state_dir), state_dir)
step.state_dir = state_dir
state = step.load(saved_state)

input_final_seq = read_input_final_seq(state_dir)
Expand Down
19 changes: 19 additions & 0 deletions src/revenant/step.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,25 @@ class Step:
untouched, and the stage process exits non-zero. On restart, the
same item is retried from scratch. This is intentional — see
docs/design.md section 8.

Attributes set by stage runner:
blob_store: BlobStore -- for writing/reading binary artifacts
too large or unsuitable to embed in
.jsonl payloads. See
docs/design-blob-storage.md.
state_dir: Path -- the pipeline's state directory, for
resolving blob paths read from an
upstream payload.

IMPORTANT: a payload field containing a blob path (as returned by
self.blob_store.write()/commit()) must not be forwarded unchanged into
this step's own yield. If downstream needs the referenced content to
persist past this stage, re-write it via self.blob_store into this
stage's own blob directory. Blob paths are one-hop only -- see
docs/design-blob-storage.md section 1. This is not enforced by the
framework; violating it will not fail loudly today, but will make a
blob unsafe to reclaim whenever garbage collection is eventually
implemented.
"""

def load(self, saved_state: Any | None) -> Any:
Expand Down
76 changes: 76 additions & 0 deletions tests/examples/test_blob_split_merge.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,76 @@
import json

from revenant.config import StageConfig
from revenant.io_utils import make_checkpoint_line, make_record_line
from revenant.stage_runner import run_stage

from examples.blob_split_merge.pipeline import SplitFile, MergeFile


def _write_input_jsonl(path, payloads):
path.parent.mkdir(parents=True, exist_ok=True)
lines = [
json.dumps(make_record_line(seq=i, src_seq=i, parent_seq=i, payload=p))
for i, p in enumerate(payloads, start=1)
]
lines.append(json.dumps(make_checkpoint_line(last_consumed_seq=len(payloads), last_emitted_seq=len(payloads))))
path.write_text("\n".join(lines) + "\n")


def _run_split_then_merge(tmp_path, source_file):
state_dir = tmp_path / "state"
state_dir.mkdir()

split_stage = StageConfig(name="split", step_class=SplitFile, upstream="input")
merge_stage = StageConfig(name="merge", step_class=MergeFile, upstream="split")
pipeline = [split_stage, merge_stage]

_write_input_jsonl(split_stage.upstream_output_path(state_dir), [{"input": str(source_file)}])

run_stage(split_stage, state_dir, pipeline=pipeline)
run_stage(merge_stage, state_dir, pipeline=pipeline)

output_path = merge_stage.output_path(state_dir)
if not output_path.exists():
return state_dir, []

output_lines = [
json.loads(line)
for line in output_path.read_text().splitlines()
if line.strip()
]
return state_dir, [line for line in output_lines if line.get("type") == "record"]


def test_empty_file_produces_no_merge_output(tmp_path):
source_file = tmp_path / "empty.txt"
source_file.write_text("")

state_dir, records = _run_split_then_merge(tmp_path, source_file)

# SplitFile has no lines to yield for an empty file, so nothing ever
# reaches MergeFile for this input -- zero parts is not the same as
# "one part with total=0", and no output record should be produced.
assert records == []


def test_single_line_file_merges_to_one_result(tmp_path):
source_file = tmp_path / "single_line.txt"
source_file.write_text("only one line\n")

state_dir, records = _run_split_then_merge(tmp_path, source_file)

assert len(records) == 1
resolved = state_dir / records[0]["payload"]["result"]
assert resolved.read_text() == "only one line\n"


def test_multi_line_file_merges_all_parts_in_order(tmp_path):
source_file = tmp_path / "multi_line.txt"
source_file.write_text("line one\nline two\nline three\n")

state_dir, records = _run_split_then_merge(tmp_path, source_file)

assert len(records) == 1
resolved = state_dir / records[0]["payload"]["result"]
assert resolved.read_text() == "line one\nline two\nline three\n"
Loading
Loading