Skip to content

Pipeline as a durable workflow CorpusBuild

Moduledata.09 · build · Go · Pass 8 · 3 to 4 h
You buildgo/workflows/corpus_build.go: CorpusBuild(rt, in), ConfigSHA256, CorpusWorkflowID, StageOptions, the activity names corpus.stage and corpus.cleanup; then the activity form of your {corpus} CLI (run --stage <s> and cleanup, entry-point territory)
Contractcourse/contracts/formats/corpus-config.schema.json (the stage spec) and course/contracts/spec/subprocess-activity.md (how a stage runs)
Testscourse/tests/go/data_09/ (CorpusBuild through the replaying simulator, over fake stages that keep the subprocess contract; section 4)
Needsdur.08 (the Runtime seam and Saga); reading: data.08 (the pipeline you wrap, data.01 to data.08), dur.09 (the subprocess runner the stages run through)
Used bydur.11 (TrainRun builds its corpus first); later C1
MilestoneMS-durable ({ctl} data build survives kills with the same manifest hash)
Optional depthdatatrove: pipelines and executors (free); Temporal: idempotency of activities (free)
  • The workflow orders, keys, and undoes; the stages do the work. Each stage is one subprocess activity at an on-disk boundary of the pipeline (fetch, shard, tokenize).
  • The workflow id carries the config’s content hash and each stage’s activity id is its name, so the idempotency key corpus/<dataset>/<version>/<hash12>/<stage> is stable across attempts and different for every config.
  • The hash is of the config’s content, not its spelling: decoded and re-encoded with sorted keys.
  • A kill at any point gives the same result and files as an uninterrupted build; a finished stage is replayed from history, a stage in flight reruns as a no-op or a whole rewrite.
  • Cancel or failure deletes the build’s partial outputs (a compensation registered before the shard stage), never the shared raw-document cache.
Terminal window
ol start data.09 # stubs go/workflows/corpus_build.go
ol tests data.09 # read the test catalog first
ol check data.09
ol diff data.09

Then, in your worker’s composition root, register the workflow (adapting workflow.Context to Runtime) and two subprocess activities over {corpus}: corpus.stage runs run --stage <in.Stage> with in.Config as the spec, corpus.cleanup runs cleanup with its input as the spec. Add {ctl} data build --config <toml>: parse the TOML, start CorpusBuild with workflow id CorpusWorkflowID(config).


Your corpus pipeline (data.01 to data.08) is one command: {corpus} run --config small.toml. On the real TinyStories corpus it runs for an hour, downloads gigabytes, and dies with the laptop’s battery or a flaky mirror; a rerun starts from scratch, and a run you stop halfway leaves shards that the next training run reads as a complete corpus. You now own a durable engine (dur.01 to dur.08) and a way to run Python as activities (dur.09). CorpusBuild puts them together: each stage becomes an activity with a stable idempotency key, a kill anywhere resumes, and a cancel cleans up after itself.

A durable workflow pays a round trip through the server (history events, a task, a poll) for every activity, and only work that writes its result somewhere durable can be resumed. The pipeline’s stages pass documents to each other in memory (compose of generator stages, data.02), so CorpusBuild cuts it where the outputs are files:

Activity corpus.stageRunsWrites
fetchrun --stage fetchcorpus/raw/, corpus/LEDGER.jsonl (a cache shared by every build)
shardrun --stage shard (filters, dedup, PII, shards, ledger check)corpus/<dataset>/<version>/ with _MANIFEST.json
tokenizerun --stage tokenize (from the manifest)tokens/<tokenizer>/<dataset>/
SymbolMeaning
ccthe corpus config (JSON)
canon(c)\text{canon}(c)cc decoded and re-encoded: object keys sorted, no whitespace
h=sha256(canon(c))h = \text{sha256}(\text{canon}(c))the config hash, 64 hex digits
ww = corpus/<dataset>/<version>/ + h0..11h_{0..11}the workflow id
ks=wk_s = w + / + ssstage ss‘s idempotency key

StartWorkflow is idempotent on the workflow id (dur.02): starting the same config again finds the run that exists, and a changed config, even one changed value, is a different run with different keys, so it can never be mistaken for a rerun of the old one. Hashing the canonical form means a config with other whitespace or key order is the same build.

A stage that finished is in the run’s history: a replay answers it without running anything. A stage that was running when the worker died is redelivered with the same key and runs again in the same work directory, where DONE.json (dur.09) makes it a no-op, or, without one, the stage rewrites its outputs whole (every file is published by rename). Either way the build’s result, and every file it wrote, equals an uninterrupted build’s.

The compensation corpus.cleanup deletes corpus/<dataset>/<version>/ and tokens/<tokenizer>/<dataset>/. It is registered with the saga before the shard stage (a shard stage interrupted halfway may have left files) and not before fetch (a fetch that failed left nothing of this build, and deleting would destroy a previous good build of the same version). It runs once, on the detached runtime, whether the build was cancelled (ErrCanceled) or a stage failed for good (an unlicensed source exits 65).

Config {"version": "v1", "dataset": "tinystories", "tokenizer": {"id": "bytes"}}. Its canonical form is the 67 bytes

{"dataset":"tinystories","tokenizer":{"id":"bytes"},"version":"v1"}

whose sha256 is a03711ec959b2a8fb6cf51d9445bee90fcb69a103cfa033483b52f836269c0a6. So the workflow id is corpus/tinystories/v1/a03711ec959b, and the three activities run in order with keys .../a03711ec959b/fetch, /shard, /tokenize, each with the stage name as activity id, a 2-minute heartbeat timeout, and (for fetch) more attempts. The result names the manifest corpus/tinystories/v1/_MANIFEST.json and the tokens tokens/bytes/tinystories. A cancel that arrives while tokenize waits runs corpus.cleanup with key .../cleanup, deletes the shards and tokens, keeps corpus/raw/, and ends the run canceled.

This is TestHandExampleKeysAndOrder.

const (
ActCorpusStage = "corpus.stage"
ActCorpusCleanup = "corpus.cleanup"
)
var CorpusStages = []string{"fetch", "shard", "tokenize"}
type CorpusBuildInput struct{ Config json.RawMessage }
type StageInput struct{ Stage, ConfigSHA256 string; Config json.RawMessage }
type CleanupInput struct{ Dataset, Version, TokenizerID, ConfigSHA256 string }
type StageDone struct{ Outputs []string }
type CorpusBuildResult struct {
Dataset, Version, ConfigSHA256, Manifest string
Tokens []string
Stages map[string][]string
}
func ConfigSHA256(cfg json.RawMessage) (string, error)
func CorpusWorkflowID(cfg json.RawMessage) (string, error)
func StageOptions(stage string) StepOptions
func CorpusBuild(rt Runtime, in CorpusBuildInput) (CorpusBuildResult, error)
TestKINDChecksWhy it matters downstream
TestHandExampleKeysAndOrderunitsection 3: the hash, the workflow id, the keys, the stage order, the resultMS-durable compares manifests by these keys
TestConfigHashIsContentNotSpellingpropertywhitespace and key order do not change the id; a changed value doesone config, one build
TestBadConfigFailsBeforeAnyActivityboundarymissing dataset, version, or tokenizer fails non-retryable, no activityno gigabytes fetched for nothing
TestKillLoopSameResultfaulta crash at every position, during or between stages: same result, same files, each written once, no cleanupC1 trains on a corpus that survived kills
TestCancelDeletesPartialShardsfaultcancel during shard or tokenize: one cleanup, no shards or tokens left, raw cache kept, canceleda cancelled build is never read as a corpus
TestFailedStageCompensatesfaulta non-retryable stage failure cleans up and reports the stage’s failurean unlicensed source leaves nothing behind
TestFetchFailureNeedsNoCleanupboundarya failed fetch runs no cleanup and keeps an earlier buildthe previous good corpus survives
TestStageOptionsunitids are stage names, heartbeat timeouts set, fetch retries moredead workers are noticed in minutes
PitfallSymptomCaught by
1. stages out of ordertokenize reads a manifest that does not exist yetTestHandExampleKeysAndOrder (mutant s01)
2. activity ids by call orderthe key no longer names the stage; a code change shifts every keyTestHandExampleKeysAndOrder (mutant s02)
3. hashing the config’s bytes, or leaving the hash out of the idthe same config twice is two builds; a changed config reuses the old runTestHandExampleKeysAndOrder, TestConfigHashIsContentNotSpelling (mutants s03, s04)
4. no config validationa broken config fetches before it failsTestBadConfigFailsBeforeAnyActivity (mutant s05)
5. cleanup registered after the shard stagea cancel during shard leaves partial shardsTestCancelDeletesPartialShards (mutant s07)
6. compensating only on cancel, or only on failurehalf a corpus left behind by the other pathTestCancelDeletesPartialShards, TestFailedStageCompensates (mutants s08, s09)
7. cleanup registered before fetcha failed download deletes yesterday’s good buildTestFetchFailureNeedsNoCleanup (mutant s10)
8. no heartbeat timeouta dead worker is noticed after the 2 h start-to-closeTestStageOptions (mutant s11)
DirectionModuleHow it uses this
Backdur.08Runtime, Saga, ErrCanceled
Backdata.08the pipeline the stages run, ledger checks included
Backdur.09the subprocess runner the stages run through
Forwarddur.11TrainRun builds its corpus with CorpusBuild before training
ForwardC1the capstone corpus is a CorpusBuild
Your pieceProduction equivalentWhat it addsWhere to look
three stage activitiesdatatrove executors (local, Slurm)many tasks per stage, each over a slice of the input, with completion markers per taskdatatrove/executor/base.py
one shard stageper-shard fan-outeach shard its own activity, run in parallel; the manifest written by a final reduceTemporal “batch processing” samples
content-hashed workflow idsDolma’s versioned configs, DVC pipelinesoutputs addressed by the hash of their inputs and codedvc.lock