Skip to content

Subprocess activity runner (Go) and the Python activity helper

Moduledur.09 · build · Go and Python · Pass 8 · 5 to 7 h
You buildgo/activities/subprocess.go: Subprocess.Run, Task, Heartbeat, Failure, WorkDir, Tail, Classify; python/tinyllm/io/activity.py: Activity, run, the exit codes; python/tinyllm/io/telemetry.py: Tracer, parse_traceparent, OTLP/HTTP JSON export
Contractcourse/contracts/spec/subprocess-activity.md (the interface between the two halves), course/contracts/py/tinyllm/io/activity.pyi, course/contracts/py/tinyllm/io/telemetry.pyi, and the spans of otel/semconv.md
Testscourse/tests/go/dur_09/ (the runner, against a scripted fake child) and course/tests/dur.09/ (the Python helper and the exporter); what they check: section 4
NeedsL0.6 atomic checkpoints (verify_step_dir decides when a checkpoint may be reported); reading: lang.02 exit codes and signals, the worker and activity SDK of dur.04 and dur.05 (the runner is plugged into them by your worker’s composition root)
Used bydata.09 (every corpus stage runs through the runner) · dur.11 (train and eval activities) · dep.06 (the worker image carries these units) · obs.05 (runs telemetry.py as the Python end of the control-plane trace); later dur.12 (export) and C1 (the capstone run is a TrainRun)
MilestoneMS-durable (<system> train --spec specs/tiny.json survives a worker kill)
Optional depthTemporal: activity heartbeats and cancellation (free); W3C Trace Context (free); OTLP specification, JSON encoding (free); W. Richard Stevens, Advanced Programming in the UNIX Environment, ch. 9 and 10 (process groups, signals)
  • Python never speaks gRPC to the platform (D9). A Go worker runs it as a child process and the whole interface is files, environment variables, signals, and an exit code.
  • The exit code is the verdict: 65 never retries, 75 and crashes retry, 130 is a cancellation. A worker drain is not a cancellation, so the runner reports it as retryable even when the child exits 130.
  • A checkpoint becomes a heartbeat detail the moment it is complete, and the next attempt gets it as --resume. Until a newer one exists, every heartbeat keeps sending the old one.
  • Every attempt shares one work directory named by the idempotency key. Outputs are published by rename, DONE.json last, and an exclusive lock keeps an orphaned child from writing beside its successor.
  • The Python spans continue the activity’s trace from TRACEPARENT, honor its sampled flag, and a dead collector never fails training.
Terminal window
ol start dur.09 # stubs go/activities/subprocess.go and python/tinyllm/io/{activity,telemetry}.py
ol tests dur.09 # read the test catalog first
ol check dur.09 # Go tests (fake child) and Python tests; the exit code is the verdict
ol diff dur.09 # after passing: your code against the reference

Then wire it in your worker’s composition root (go/cmd/worker, learner territory). For each Python verb, register an activity that adapts the delivery to the runner:

r := activities.Subprocess{Entry: tinyllmEntry, Artifacts: "/artifacts", ServiceName: name + "-python",
OTLPEndpoint: os.Getenv("OTEL_EXPORTER_OTLP_ENDPOINT"),
Canceled: func(cause error) bool { return errors.Is(cause, worker.ErrCancelRequested) }}
w.RegisterActivity("train", func(ctx context.Context, spec []byte) ([]byte, error) {
inf, _ := worker.InfoFrom(ctx)
task := activities.Task{IdempotencyKey: inf.IdempotencyKey, Attempt: inf.Attempt,
LastHeartbeat: inf.HeartbeatDetails, TraceContext: inf.TraceContext, HeartbeatTimeout: hbTimeout}
return r.Run(ctx, task, []string{"train"}, spec, func(_ context.Context, d []byte) (bool, error) {
return false, worker.Heartbeat(ctx, d) // a cancel arrives as ctx's cause, read by r.Canceled
})
})

In your Python entries, wrap train, eval, export, and corpus run --stage in activity.run.


Your durable engine can run Go activities with retries, heartbeats, and fenced leases (dur.03 to dur.05), but everything that matters to the platform’s users is Python: the corpus stages from Pass 3 and the trainer from Pass 2. Today {tinyllm} train is a command you type. If the laptop sleeps at step 9,000, the run is gone, and nothing records which checkpoint was last good. If you start it from a Go activity with exec.Command and wait, three things break at once: the server sees no heartbeats and fences the activity mid-run, a worker restart kills the child without a checkpoint (or worse, leaves it running), and a retry starts again from step 0 in a fresh directory. This module is the bridge: a Go runner and a Python helper that keep one written contract, so CorpusBuild (data.09) and TrainRun (dur.11) can run Python work durably.

A subprocess activity is one attempt at one activity, run as a child process of the worker. The runner gives the child an argv, an environment, and a work directory; the child gives back a progress file, maybe some outputs, and an exit code. Nothing else crosses: no shared memory, no RPC. That is why the Python side stays testable without a cluster (you can run {tinyllm} train --spec ... --progress ... by hand) and why one durable protocol (gRPC, Go only) is enough.

SymbolMeaningType
kkthe idempotency key <workflow_id>/<activity_id>string
aathe attempt number, starting at 1int
ThbT_{hb}the activity’s heartbeat timeoutduration
h=Thb/3h = T_{hb}/3the longest gap allowed between heartbeatsduration
ggthe grace period between SIGTERM and SIGKILL, 30 s by defaultduration

2.2 One work directory per activity, shared by every attempt

Section titled “2.2 One work directory per activity, shared by every attempt”

The work directory is <TL_ARTIFACTS>/activities/<k>/, each /-separated segment of kk a directory (WorkDir). A segment must match [A-Za-z0-9._-]+ and must not be . or ..: workflow ids are caller-chosen text, and a key like ../../etc would otherwise escape the artifact root. Every attempt a=1,2,…a = 1, 2, \dots of one activity lands in the same directory, which is what makes a retry a resume rather than a restart.

Three rules make the directory safe to share:

  1. Publish by rename. An output is written to name.tmp, flushed and fsynced, then renamed to name, and the directory is fsynced. rename is atomic within one filesystem: a reader sees the old file or the whole new one. The runner writes spec.json the same way.
  2. DONE.json last. The last act of a successful run publishes DONE.json ({"outputs": [...]}), then emits the done event, then exits 0. A run that starts and finds DONE.json re-emits done and exits 0 without doing anything: a duplicate delivery has one effect.
  3. One attempt at a time. The helper takes an exclusive, non-blocking flock on <dir>/.lock before any work. If a previous attempt still holds it (a child orphaned by a SIGKILLed worker), the new attempt exits 75 and is retried later.

The child appends one JSON object per line to --progress and flushes after each line: the runner reads the file while it grows, so an event stuck in Python’s buffer is invisible exactly when it matters (the worker is killed one second later). The runner’s Tail keeps a line back until its newline arrives: half a line is not an event.

A ckpt event names a checkpoint step directory, relative to TL_ARTIFACTS. The helper emits it only after L0.6’s verify_step_dir says the directory is complete (its MANIFEST.json lists every file with the right size and hash). The runner treats each new ckpt path as the heartbeat details and heartbeats it at once; it also heartbeats at least every hh so the server does not fence the activity between checkpoints (training checkpoints every few minutes; heartbeat timeouts are tens of seconds). On a retry the server hands back the last details, and the runner passes them as --resume <ckpt>. Until the child reports a newer checkpoint, every heartbeat repeats the old one: an empty heartbeat would erase the resume point on the server.

CodeMeaningFailure
0 and DONE.json presentsuccessnone: the result is the bytes of DONE.json (at most 2 MiB)
0 without DONE.jsona broken entryMissingDone, retryable
65 (EX_DATAERR)the input is wrong: invalid spec, unlicensed sourceExitCode65, non-retryable
75 (EX_TEMPFAIL)transient: a flaky download, an OOM the next attempt may avoidExitCode75, retryable
130cancelledCanceled, non-retryable
any other code, or a signala crashExitCode<n> or Signal<n>, retryable

Failure.message is the last 4 KiB of the child’s stderr, where a traceback ends. On the Python side run(main) maps exceptions to the same codes: SpecError 65, RetryableError 75, Cancelled 130, anything else 1.

2.5 Stopping a child: process groups, SIGTERM, SIGKILL

Section titled “2.5 Stopping a child: process groups, SIGTERM, SIGKILL”

The runner starts the child with Setpgid, so the child leads a new process group whose id is its pid. Signals go to the whole group (kill(-pgid, sig)), which reaches anything the child started (a tokenizer pool, a data loader worker), and a signal meant for the worker’s own group does not reach the child behind the runner’s back.

There are two reasons to stop a child, and they are different results:

WhyHow the runner learns itResult
the activity is cancelled (a workflow cancel)a heartbeat answer with cancel_requestedCanceled, non-retryable
the worker is going away (SIGTERM drain on a deploy, or a lost lease)ctx endsWorkerShutdown, retryable: another worker resumes from the checkpoint

The worker SDK of dur.04 reports a cancel_requested heartbeat answer by cancelling the activity’s context with its own cause; Subprocess.Canceled tells the runner which causes mean “cancelled” (everything else is a shutdown). Either way the runner sends SIGTERM to the group, waits up to gg, then sends SIGKILL. The Python handler only sets a flag; the training loop notices it at its next step, saves a checkpoint, emits its ckpt event, and raises Cancelled (exit 130). The runner heartbeats that last checkpoint even after ctx has ended (on a fresh context with a short timeout) and returns it in Failure.Details for FailActivityRequest.last_heartbeat_details.

A worker that dies outright cannot stop its child. The child notices instead: its parent pid changes (it is reparented), and Activity.cancelled becomes true.

2.6 Telemetry: one trace from the CLI to the training step

Section titled “2.6 Telemetry: one trace from the CLI to the training step”

The worker passes the activity span’s W3C context in TRACEPARENT: 00-<trace-id>-<parent-id>-<flags>, lowercase hex, 32 and 16 digits, neither all zeros, flags bit 0 meaning sampled. The exporter makes train.run (or corpus.stage <s>) a child of that span in the same trace, train.step a child of train.run every 50th step, and posts them as OTLP/HTTP JSON to OTEL_EXPORTER_OTLP_ENDPOINT/v1/traces (gauges such as tl.train.loss to /v1/metrics: a subprocess cannot be scraped). Three rules:

  • An invalid TRACEPARENT starts a new root trace; it is never parsed leniently.
  • An unsampled parent means export nothing: the decision belongs to the root of the trace.
  • Export failures return False and drop the batch. Telemetry never raises into the work it observes.

Activity train-1/3 (workflow train-1, activity id 3) with TL_ARTIFACTS=/artifacts, heartbeat timeout 150 ms, so hh = 50 ms.

Attempt 1. The runner computes the work directory /artifacts/activities/train-1/3/, publishes spec.json there, and execs

<entry> train --spec /artifacts/activities/train-1/3/spec.json --progress /artifacts/activities/train-1/3/progress.jsonl

with TL_IDEMPOTENCY_KEY=train-1/3, TL_ATTEMPT=1. No --resume: the task has no heartbeat details. The child writes

{"ts":...,"kind":"step","step":500,"loss":2.31,"lr":0.0009,"tokens":8192000}
{"ts":...,"kind":"ckpt","step":500,"ckpt":"runs/train-1/ckpt/step-000500"}

The tail reads the ckpt line within one poll, and the runner heartbeats with details runs/train-1/ckpt/step-000500. Then the child runs out of memory and exits 75 with CUDA out of memory on stderr. The result is Failure{Type: "ExitCode75", NonRetryable: false, Message: "...CUDA out of memory..."}: retry.

Attempt 2. The server redelivers with attempt = 2 and last_heartbeat_details = runs/train-1/ckpt/step-000500. The runner uses the same directory and execs the same argv plus --resume runs/train-1/ckpt/step-000500, with TL_ATTEMPT=2. The helper (clock fixed at 1760000000.0) writes exactly

{"ts":1760000000.0,"kind":"step","step":600,"loss":2.25,"lr":0.0009,"tokens":9830400}
{"ts":1760000000.0,"kind":"ckpt","step":1000,"ckpt":"runs/train-1/ckpt/step-001000"}
{"ts":1760000000.0,"kind":"done","outputs":["runs/train-1/ckpt/step-001000"]}

and DONE.json is the 45 bytes {"outputs":["runs/train-1/ckpt/step-001000"]}, which become the activity’s result. A third delivery (the worker died before CompleteActivityTask) finds DONE.json, re-emits done, and exits 0 without training.

Telemetry. With TRACEPARENT=00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01 (the W3C example) and span ids drawn from a counter, train.run gets span id 0000000000000001 and parent 00f067aa0ba902b7; the sampled step 50 (step 49 is not sampled) gets 0000000000000002 with parent 0000000000000001; both carry trace id 4bf92f3577b34da6a3ce929d0e0e4736. In the OTLP JSON body, tl.train.step = 50 is {"intValue": "50"}: a string, because JSON numbers lose precision past 2532^{53}.

These numbers are TestHandExampleResumeFlow, test_hand_example_progress_lines, and test_hand_example_span_tree.

go/activities/subprocess.go
type Task struct {
IdempotencyKey string
Attempt int
LastHeartbeat []byte
TraceContext map[string]string
HeartbeatTimeout time.Duration
}
type Heartbeat func(ctx context.Context, details []byte) (cancelRequested bool, err error)
type Subprocess struct {
Entry []string; Artifacts, OTLPEndpoint, ServiceName string
Grace, Poll time.Duration; Env []string; Dir string
Canceled func(cause error) bool // is this end of ctx a cancel (not a shutdown)?
}
func (s Subprocess) Run(ctx context.Context, t Task, args []string, spec []byte, hb Heartbeat) ([]byte, error) // *Failure
type Failure struct { Type, Message string; NoRetry bool; ExitCode, Signal int; Details []byte }
func (f *Failure) FailureType() string // what the worker SDK reports as Failure.type
func (f *Failure) NonRetryable() bool // and as Failure.non_retryable
func WorkDir(artifacts, key string) (string, error)
type Tail struct{ Path string } // Read() ([]Progress, error): complete lines only
func Classify(state *os.ProcessState, stderr []byte) *Failure
python/tinyllm/io/activity.py
class Activity: # Activity(argv, env=None, clock=time.time)
def load_spec(self, validate=None) -> dict: ...
def step(self, step, loss, lr, tokens) -> None: ...
def checkpoint(self, step, path) -> str: ...
def publish(self, name, data) -> str: ...
def done(self, outputs) -> bytes: ...
@property
def cancelled(self) -> bool: ...
def run(main, argv=None, env=None) -> int: ...
# python/tinyllm/io/telemetry.py
class Tracer: # Tracer.from_env(); with tracer.span("train.run", {...}) as s: ...
def flush(self) -> bool: ...
def parse_traceparent(value) -> SpanContext | None: ...
def should_sample_step(step, every=50) -> bool: ...

The runner takes a Task and a Heartbeat function instead of importing the worker SDK, so it is tested here against a scripted fake child, and your worker’s composition root adapts each delivery to it (a dozen lines per verb).

TestKINDChecksWhy it matters downstream
TestHandExampleResumeFlowunitsection 3: attempt 1 heartbeats its checkpoint and fails retryable; attempt 2 gets --resume, TL_ATTEMPT=2, the same directory, and returns DONE.jsonTrainRun resumes from the last checkpoint after a kill (dur.11)
TestWorkDirFromKeyboundarykey segments become directories; .., ., empty, and odd characters are rejected, non-retryablea workflow id can never write outside /artifacts/activities
TestEnvironmentContractunitTL_*, TRACEPARENT, OTEL variables set per task; the worker’s own copies never leakPython spans land in the right trace (obs.05)
TestExitCodeTableunitevery row of the table in section 2.4, signals includeda poison input fails fast; a transient one retries
TestSuccessNeedsDoneFileboundaryexit 0 without DONE.json is a retryable MissingDone; over 2 MiB is non-retryable; exactly 2 MiB is fineresults fit in a gRPC message and a history event
TestTailKeepsPartialLinesunithalf a line waits for its newline; garbage lines are skippedheartbeats never carry half a path
TestHeartbeatCarriesNewCheckpointsuniteach new checkpoint is heartbeated while the child runs, even a line written in two piecesa kill one second after a checkpoint loses nothing
TestResumePointSurvivesUntilANewCheckpointfaulta retry’s heartbeats repeat the old resume point until a newer one existsa kill before the retry’s first checkpoint does not restart from step 0
TestHeartbeatEveryThirdOfTimeoutunitat least one heartbeat per Thb/3T_{hb}/3the server does not fence a long step
TestCancelRequestedStopsTheChildfaultcancel_requested sends SIGTERM; the run ends Canceled with the exit checkpoint heartbeateddur.08 cancel reaches Python
TestGraceThenSIGKILLfaulta child ignoring SIGTERM is SIGKILLed after the grace, grandchild includeda stuck child never holds a worker slot
TestOwnProcessGroupunitthe child leads its own groupgroup signals reach grandchildren only
TestWorkerShutdownIsRetryablefaultctx ending gives a retryable WorkerShutdown with Details = the exit checkpointa deploy does not end a training run for good
TestCancelCauseFromContextfaultwith Subprocess.Canceled set, a ctx cancelled with the SDK’s cancel cause ends Canceled; any other cause stays WorkerShutdownthe worker SDK’s cancel reaches the child as a cancel
TestSpecWrittenAtomicallyunitspec.json replaced by rename, never rewritten in placea retry never shows half a spec
TestStderrTailIsBoundedboundarythe last 4 KiB of stderr, traceback includedfailures stay small in the history
test_hand_example_progress_linesunitsection 3’s three progress lines and DONE.json, byte for bytethe runner and the schema read these bytes
test_flags_and_environmentunit--spec, --progress, --resume (both spellings), TL_* defaults; missing flags are SpecErrorthe entry and the runner agree on argv
test_each_event_is_one_flushed_lineuniteach event is on disk as one line when emit returnsthe tail sees checkpoints at once
test_rejects_bad_eventsboundaryunknown kinds and NaN or inf values raise, nothing writtena NaN loss never corrupts the progress file
test_checkpoint_only_when_completeunitno event for an incomplete directory; paths relative to TL_ARTIFACTS; outside paths rejectedthe resume point is always loadable (L0.6)
test_publish_never_exposes_a_partial_filefaultwhen the rename never happens, the final name does not existDONE.json exists only for finished work
test_outputs_stay_in_the_work_dirboundary.., absolute, and empty names rejectedone activity cannot overwrite another’s files
test_run_exit_codesunitSpecError 65, RetryableError 75, Cancelled 130, crash 1, success 0 with DONE.jsonthe Go side’s table has the right inputs
test_spec_validation_is_exit_65unitvalidation errors and non-object or bad JSON specs are 65, before any outputinvalid specs do not burn retries
test_rerun_after_done_is_a_noopregressiona second run finds DONE.json, skips main, re-emits donea duplicate delivery has one effect
test_sigterm_checkpoints_and_exits_130faultSIGTERM, then a checkpoint as the last event, then exit 130 within 5 scancel and drain keep the work done so far
test_second_attempt_is_locked_outfaultwhile one attempt holds the lock, another exits 75 without runningno two children write one directory
test_orphaned_child_stopsfaulta child whose parent died sees cancelledan orphan does not train on beside its successor
test_hand_example_span_treeconformancesection 3’s ids, parents, kinds, and JSON encoding at a real HTTP receiverobs.05 finds train.step under activity train
test_parse_traceparentboundaryvalid, unsampled, future-version, and nine invalid headersinvalid input starts a new trace
test_format_roundtripunitformat_traceparent and parse_traceparent agreecontexts can be handed on
test_unsampled_parent_exports_nothingboundaryflags 00: no spans, no gauges sentno orphan spans
test_no_traceparent_starts_a_rootboundaryno or bad TRACEPARENT: a fresh trace id and no parenthand runs still trace
test_from_envunitendpoint, service name, resource attributes, parent from the environmentthe runner’s variables are the only input
test_dead_collector_never_fails_the_workfaultconnection refused: flush returns False within its timeout, queue droppedtraining survives a collector outage
test_collector_error_status_is_falsefaulta 500 answer is Falsedrops are visible
test_error_marks_the_spanunitan exception sets status ERROR and error.type, and propagatesfailed stages show in the trace
test_step_samplingunitevery 50th stepspan volume stays bounded
test_gauges_are_pushedunitthe latest gauge value per name and attributes goes to /v1/metricstl.train.loss reaches Prometheus
test_shutdown_stops_exportunitnothing is queued after shutdownclean exit
test_body_shape_without_parentunita root span has no parentSpanId; attribute value typesthe collector accepts the body
PitfallSymptomCaught by
1. checkpoints never reach a heartbeatafter a kill, the retry starts from step 0TestHandExampleResumeFlow, TestHeartbeatCarriesNewCheckpoints (mutant s01)
2. the key’s segments are not checkeda workflow id ../../x writes outside /artifactsTestWorkDirFromKey (mutant s02)
3. --resume is not passedretries ignore the checkpointTestHandExampleResumeFlow (mutant s03)
4. the worker’s TRACEPARENT is inheritedPython spans join an unrelated traceTestEnvironmentContract (mutant s04)
5. a wrong row in the exit-code tablea poison spec retries until the dead-letter queue; a cancel restartsTestExitCodeTable, test_run_exit_codes (mutants s05, s06, s24)
6. exit 0 trusted without DONE.jsona broken entry “succeeds” with no outputsTestSuccessNeedsDoneFile (mutant s07)
7. half a progress line is parseda checkpoint event is lost foreverTestTailKeepsPartialLines (mutant s08)
8. a retry heartbeats empty detailsthe resume point is erased before the first new checkpointTestResumePointSurvivesUntilANewCheckpoint (mutant s09)
9. heartbeats only on checkpointsthe server fences the activity mid-stepTestHeartbeatEveryThirdOfTimeout (mutant s10)
10. cancel_requested ignoreda cancelled workflow keeps trainingTestCancelRequestedStopsTheChild (mutant s11)
11. no process group, or no SIGKILLgrandchildren and stuck children outlive the activityTestGraceThenSIGKILL, TestOwnProcessGroup (mutants s12, s13)
12. a drain reported as a cancel, or a cancel as a drainevery deploy ends a training run for good; a cancelled run is retriedTestWorkerShutdownIsRetryable, TestCancelCauseFromContext (mutants s14, s40)
13. spec.json written in placea reader sees a truncated specTestSpecWrittenAtomically (mutant s15)
14. all of stderr kepta chatty child blows the 4 MiB message limitTestStderrTailIsBounded (mutant s16)
15. progress not flushedthe runner sees nothing until exittest_each_event_is_one_flushed_line (mutant s17)
16. absolute paths in events or DONE.jsona path that is right on one pod is wrong on the nexttest_hand_example_progress_lines, test_checkpoint_only_when_complete (mutants s18, s21)
17. a NaN loss written as NaNthe progress line is not JSONtest_rejects_bad_events (mutant s19)
18. a ckpt event before the checkpoint is completethe retry cannot load its resume pointtest_checkpoint_only_when_complete (mutant s20)
19. outputs written under their final namea kill leaves a truncated DONE.json that looks finishedtest_publish_never_exposes_a_partial_file (mutant s22)
20. output names not confinedone activity overwrites another’s filestest_outputs_stay_in_the_work_dir (mutant s23)
21. spec validation errors ignoreda bad spec trains on garbage, or retriestest_spec_validation_is_exit_65 (mutant s25)
22. no DONE.json check on starta duplicate delivery trains againtest_rerun_after_done_is_a_noop (mutant s26)
23. the SIGTERM handler exits at oncethe work since the last checkpoint is losttest_sigterm_checkpoints_and_exits_130 (mutant s27)
24. no work-dir lockan orphan and its successor write one directorytest_second_attempt_is_locked_out (mutant s28)
25. orphans not detecteda child trains on after its worker diedtest_orphaned_child_stops (mutant s29)
26. the remote parent ignoredPython spans form a separate tracetest_hand_example_span_tree (mutant s30)
27. intValue as a JSON numbercollectors reject or round the attributetest_hand_example_span_tree, test_body_shape_without_parent (mutant s31)
28. lenient traceparent parsingall-zero ids glue unrelated runs togethertest_parse_traceparent (mutant s32)
29. the sampled flag ignoredorphan spans whose parents were never recordedtest_unsampled_parent_exports_nothing (mutant s33)
30. OTEL_SERVICE_NAME ignoredevery system’s Python is tinyllm-pythontest_from_env (mutant s34)
31. export errors raiseda collector outage kills trainingtest_dead_collector_never_fails_the_work, test_collector_error_status_is_false (mutant s35)
32. exceptions swallowed by the spana failing stage looks successfultest_error_marks_the_span (mutant s36)
33. sampling off by onestep 0 never traced, step 1 alwaystest_step_sampling (mutant s37)
34. the first gauge value keptdashboards show the loss of step 1test_gauges_are_pushed (mutant s38)
35. spans queued after shutdownmemory grows at exit; spans half senttest_shutdown_stops_export (mutant s39)
DirectionModuleHow it uses this
BackL0.6verify_step_dir gates every ckpt event; --resume loads with load_checkpoint
Backdur.04, dur.05the worker and activity SDK your composition root adapts to Task and Heartbeat
Forwarddata.09CorpusBuild runs each corpus stage as {corpus} run --stage <s> through the runner
Forwarddur.11TrainRun and EvalSuite run {tinyllm} train and eval; kills resume from the heartbeated checkpoint
Forwardobs.05the control-plane trace reaches train.step through this exporter
Forwarddep.06the worker image carries the activity helper and the exporter
Forwarddur.12, C1export and the capstone run use the same contract

If you skip this module, ol check data.09 fails with needs dur.09: build it, or pass --ref-deps.

Your pieceProduction equivalentWhat it addsWhere to look
Subprocess.RunTemporal activity heartbeats with detailsheartbeat throttling, details as typed payloads, async completion for work finished elsewhereTemporal Go SDK activity.RecordHeartbeat, activity.GetHeartbeatDetails
the exit-code tableKubernetes Job podFailurePolicyexit-code rules that fail the Job, ignore the failure, or count itKubernetes docs, “Handling retriable and non-retriable pod failures”
process groups and gracesystemd KillMode=control-group, Kubernetes terminationGracePeriodSecondscgroups instead of process groups: no child can escape by calling setsidsystemd.kill(5)
orphan detection by parent pidLinux PR_SET_PDEATHSIGthe kernel signals the child when the parent thread diesprctl(2)
telemetry.pythe OpenTelemetry Python SDK with the OTLP/HTTP exporterbatching span processor, retries with backoff, protobuf encodingopentelemetry-exporter-otlp-proto-http