Documentation
¶
Overview ¶
Package durable provides embeddable durable task execution for a single OS process. Register typed tasks, run memoised steps, and resume from a local journal after a crash. A step may return ErrStepPending and complete later via CompleteStep.
e, err := durable.NewEngine(ctx, "./data")
if err != nil { ... }
defer e.Close()
_ = durable.RegisterTask(e, "process-order", durable.Func(
func(ctx context.Context, s *durable.StepRunner, in OrderInput) (OrderOutput, error) {
charged, err := durable.RunStep(ctx, s, "charge", in, func(ctx context.Context, in OrderInput) (string, error) {
return chargeCard(in)
}).Get(ctx)
if err != nil {
return OrderOutput{}, err
}
return OrderOutput{Result: charged}, nil
},
))
run := durable.RunTask[OrderInput, OrderOutput](ctx, e, "process-order", "", input)
output, err := run.Get(ctx)
Index ¶
- Constants
- Variables
- func CompleteStep[O any](ctx context.Context, e *Engine, token string, result O) error
- func RegisterTask[I, O any](e *Engine, taskID string, task Task[I, O], opts ...TaskOption) error
- type AESGCMKey
- type Engine
- func (e *Engine) CancelRun(ctx context.Context, taskID, runID string) error
- func (e *Engine) Close() error
- func (e *Engine) DeleteTask(ctx context.Context, taskID string) error
- func (e *Engine) DeleteTaskRun(ctx context.Context, taskID, runID string) error
- func (e *Engine) GetStep(ctx context.Context, taskID, runID, stepID string) (StepRecord, bool, error)
- func (e *Engine) GetTask(ctx context.Context, taskID, runID string) (TaskInfo, bool, error)
- func (e *Engine) ListTasks(ctx context.Context, statuses ...TaskStatus) ([]TaskInfo, error)
- func (e *Engine) ListTasksPage(ctx context.Context, offset, limit int, statuses ...TaskStatus) (TaskPage, error)
- func (e *Engine) LoadInput(ctx context.Context, taskID, runID string) ([]byte, bool, error)
- func (e *Engine) LoadSteps(ctx context.Context, taskID, runID string) ([]StepRecord, error)
- func (e *Engine) WatchSteps(ctx context.Context, taskID, runID string, fromOffset int, ...) (<-chan StepEvent, error)
- type Option
- func WithAutoPurge(age time.Duration, interval ...time.Duration) Option
- func WithAutoPurgeMaxBytes(n int64) Option
- func WithAutoPurgeMaxRuns(n int) Option
- func WithDefaultStepTokenTTL(d time.Duration) Option
- func WithJournalMACKey(key []byte) Option
- func WithLockTimeout(d time.Duration) Option
- func WithLogger(l *slog.Logger) Option
- func WithMaxRetries(n int) Option
- func WithPayloadCodec(c PayloadCodec) Option
- func WithStepTokenKey(key []byte) Option
- func WithTimeout(d time.Duration) Option
- func WithUnsignedStepTokens() Option
- type PayloadCodec
- type ReadOnlyEngine
- func (r *ReadOnlyEngine) Close() error
- func (r *ReadOnlyEngine) GetStep(ctx context.Context, taskID, runID, stepID string) (StepRecord, bool, error)
- func (r *ReadOnlyEngine) GetTask(ctx context.Context, taskID, runID string) (TaskInfo, bool, error)
- func (r *ReadOnlyEngine) ListTasks(ctx context.Context, statuses ...TaskStatus) ([]TaskInfo, error)
- func (r *ReadOnlyEngine) ListTasksPage(ctx context.Context, offset, limit int, statuses ...TaskStatus) (TaskPage, error)
- func (r *ReadOnlyEngine) LoadInput(ctx context.Context, taskID, runID string) ([]byte, bool, error)
- func (r *ReadOnlyEngine) LoadSteps(ctx context.Context, taskID, runID string) ([]StepRecord, error)
- type RunOption
- type StepEvent
- type StepOption
- type StepRecord
- type StepRun
- type StepRunner
- type StepStatus
- type Task
- type TaskFunc
- type TaskInfo
- type TaskOption
- type TaskPage
- type TaskRun
- type TaskStatus
Constants ¶
const ( // AESGCMRekeyAfter is the recommended maximum number of Encode calls // on one 96-bit-nonce key. Random-nonce birthday collision risk // becomes material around 2^32 seals; rotate sooner for high volume. AESGCMRekeyAfter uint64 = 1 << 32 )
Variables ¶
var ( // ErrTaskNotRegistered is returned by RunTask.Get when the taskID has not // been registered on this Engine. The registry is in-memory only and must // be rebuilt after every NewEngine call. ErrTaskNotRegistered = errors.New("durable: task not registered") // ErrTaskAlreadyRegistered is returned by RegisterTask when the same // taskID is registered twice on one Engine. ErrTaskAlreadyRegistered = errors.New("durable: task already registered") // ErrRunAlreadyFinished is returned by CompleteStep when the target step // is already completed or the run is in a terminal state (completed or failed). ErrRunAlreadyFinished = errors.New("durable: run already completed or failed") // ErrRunActive is returned by DeleteTaskRun and DeleteTask when the target run // (or any run under the task) is currently executing, including while // blocked in StatusWaiting. ErrRunActive = errors.New("durable: run is currently executing") // ErrInvalidRunID is returned when a runID contains path separators or // parent-directory references that would escape the data directory. ErrInvalidRunID = errors.New("durable: invalid run ID") // ErrEngineLocked is returned by NewEngine and NewReadOnlyEngine when // dataDir is already held by another engine instance (same process or // another OS process) and the lock cannot be acquired before the timeout. ErrEngineLocked = errors.New("durable: dataDir locked by another engine") // ErrStepPending is returned from a step function to suspend that step // until CompleteStep delivers a result for it. The engine writes // StepStatusWaiting and blocks that step's Get — not the whole task. // Sibling steps started before this one keep running; call Get on them // independently or select on their Done channels. ErrStepPending = errors.New("durable: step pending external completion") // ErrInvalidToken is returned by CompleteStep when the token cannot be // decoded into taskID, runID, and stepID, or when an HMAC token has a // missing or wrong MAC (including unsigned tokens while a secret is set). ErrInvalidToken = errors.New("durable: invalid step token") // ErrTokenExpired is returned by CompleteStep when an HMAC step token is // past its TTL. ErrTokenExpired = errors.New("durable: step token expired") // ErrPayloadTooLarge is returned when a step's persisted record does not // fit in one journal frame (32 MiB of encoded input plus result). The step // fails rather than writing a frame that replay and compaction could not // read back. Keep large blobs out of step results — store a handle and // fetch the payload inside fn. ErrPayloadTooLarge = errors.New("durable: payload too large for one journal frame") // ErrEngineClosed is returned by RunTask when the Engine is already // closing or closed, so no new run is started after the exclusive flock // on dataDir has been released. ErrEngineClosed = errors.New("durable: engine is closed") // ErrRunCancelled is the error stored on a run (TaskInfo.Error, as its // string form) and returned by RunStep/Get after CancelRun. It marks // the run as StatusFailed for the same reason engine Close and task/run // timeouts already do — but with a distinct message so callers can // tell a deliberate CancelRun apart from a generic ctx cancellation or // deadline. Because TaskInfo.Error is a plain string, matching after a // resume requires comparing against ErrRunCancelled.Error(), not // errors.Is. ErrRunCancelled = errors.New("durable: run was cancelled") )
Functions ¶
func CompleteStep ¶ added in v0.1.2
CompleteStep delivers result to a suspended step identified by token. The payload is appended as a SignalEntry (durable) before the in-process waiter is signalled, so a crash after this call still resumes on the next RunTask. A second call with the same token is a no-op once the SignalEntry is on disk. Returns ErrRunAlreadyFinished if the step or run is already terminal, ErrInvalidToken if the token cannot be decoded or authenticated, and ErrTokenExpired if an HMAC token is past its TTL.
The run's status transition back to StatusRunning is owned by the waiting step itself (see leaveWaiting), not by CompleteStep — with several steps possibly waiting at once, only the step that actually resumes knows whether a sibling is still pending.
Concurrent CompleteStep calls for the same token are serialised on a per-(taskID,runID,stepID) lock distinct from the run-execution lock (which is held for the entire run, including while blocked waiting — locking it here would deadlock). A CompleteStep that loses a race with the run reaching a terminal state may append an orphan SignalEntry that is never read; this is harmless and does not corrupt the journal.
func RegisterTask ¶ added in v0.1.2
RegisterTask stores taskID → (closure + config) in the in-memory registry. Must be called before RunTask. I is the single task input type (one struct of fields, or struct{} if the task has no payload). Re-registering the same taskID returns ErrTaskAlreadyRegistered. The registry is not persisted; call this again after every NewEngine. taskID must not contain path separators or ':'.
Types ¶
type AESGCMKey ¶ added in v0.1.5
AESGCMKey is one AES-GCM key plus the identifier written next to each ciphertext when using the v2 wire format. ID 0 is also how v1 blobs (no key-ID byte) are decoded.
type Engine ¶ added in v0.1.2
type Engine struct {
// contains filtered or unexported fields
}
Engine is the process-level entry point for durable task execution. Create one per dataDir via NewEngine. Safe for concurrent use.
Always call Close to cancel in-flight runs, release the exclusive flock, stop the purger, and close journal file handles.
func NewEngine ¶ added in v0.1.2
NewEngine opens or creates dataDir, acquires an exclusive flock on <dataDir>/.lock, and initialises the in-memory task registry. Fails with ErrEngineLocked if another Engine or ReadOnlyEngine holds the directory.
func (*Engine) CancelRun ¶ added in v0.1.3
CancelRun requests cancellation of a run. It persists a durable cancel signal — reusing the SignalEntry/CompleteStep journal plumbing via a reserved signal ID — so the intent survives a crash: on the next RunTask for this taskID/runID, the run's ctx is cancelled before Task.Exec is invoked, and every RunStep call fails fast with ErrRunCancelled instead of re-running fn. If the run is currently executing in this process, its ctx is also cancelled immediately.
Cancelling ctx only signals a step function; it cannot forcibly stop one. RunStep.Get and RunTask.Get both return promptly regardless — they select on ctx.Done() independently of whether the step's goroutine has exited. But the engine still waits for that goroutine to actually return before the run reaches a terminal state or Close/compactJournal can safely proceed (see Close). If the step function never checks ctx and never returns, that goroutine runs to completion in the background and the run stays non-terminal until it does — Close called afterward would block on it too. Write step functions so they select on ctx instead of running unconditionally to completion.
The run ends up StatusFailed (the same terminal status used for engine Close and task/run timeouts) with TaskInfo.Error set to ErrRunCancelled.Error(), so callers can distinguish a deliberate cancel from another failure by comparing that string.
Returns ErrRunAlreadyFinished if the run is already StatusCompleted or StatusFailed, or an error if taskID/runID is invalid or the run does not exist. Idempotent: calling it more than once on the same run is a no-op after the first call's signal is durably written.
func (*Engine) Close ¶ added in v0.1.2
Close cancels in-flight runs, waits for them and the auto-purger to finish writing, then releases the exclusive flock and journal handles. Safe to call more than once.
Cancelling ctx only signals — it does not forcibly stop a step function. Close waits for every in-flight RunStep goroutine to actually return (see StepRunner.waitInFlight) so a step's own journal append can never race compactJournal or a second Engine opening the same dataDir. If a step function never checks ctx and never returns, Close blocks forever and the exclusive flock on dataDir is never released. Write step functions so they select on ctx (or pass it to ctx-aware calls like http.NewRequestWithContext) instead of running unconditionally to completion.
func (*Engine) DeleteTask ¶ added in v0.1.2
DeleteTask removes all runs under taskID. Destructive — use for a full wipe only. Returns ErrRunActive if any run under taskID is currently executing (Running or Waiting in-process).
func (*Engine) DeleteTaskRun ¶ added in v0.1.2
DeleteTaskRun removes one run directory and its journal. No-op if not found. Returns ErrRunActive if the run is currently executing, including while blocked in StatusWaiting — runLocks is held for that entire duration.
func (*Engine) GetStep ¶ added in v0.1.3
func (e *Engine) GetStep(ctx context.Context, taskID, runID, stepID string) (StepRecord, bool, error)
GetStep returns the latest StepRecord for one stepID in O(1) after one journal load — (zero, false, nil) if the run has no record for that stepID yet (including one stuck mid-execution: an unfinished STARTED step never appears here — see loadJournal).
func (*Engine) GetTask ¶ added in v0.1.2
GetTask returns a single run's metadata. (zero, false, nil) if not found.
func (*Engine) ListTasks ¶ added in v0.1.2
ListTasks returns TaskInfo records. Zero args returns every status. Pass explicit statuses to filter, e.g. ListTasks(ctx, StatusRunning, StatusWaiting) for recovery after a restart.
func (*Engine) ListTasksPage ¶ added in v0.1.5
func (e *Engine) ListTasksPage(ctx context.Context, offset, limit int, statuses ...TaskStatus) (TaskPage, error)
ListTasksPage is ListTasks with a 0-based offset into the sorted result. limit <= 0 means the rest of the list. Use it when dataDir holds more runs than you want to load at once (recovery, dashboards, purge).
func (*Engine) LoadInput ¶ added in v0.1.4
LoadInput returns the JSON-encoded task input written on first RunTask (one I value). (nil, false, nil) if input.json is missing.
func (*Engine) LoadSteps ¶ added in v0.1.2
LoadSteps returns all StepRecords for a run. Used to inspect progress. Includes waiting, completed, and failed steps. Order is not sorted or otherwise guaranteed — it reflects map iteration order internally. Use WatchSteps if you need append order or history.
func (*Engine) WatchSteps ¶ added in v0.1.3
func (e *Engine) WatchSteps(ctx context.Context, taskID, runID string, fromOffset int, fromByteOffset int64) (<-chan StepEvent, error)
WatchSteps streams step lifecycle events for a run: STARTED, WAITING, COMPLETED, FAILED — every append, in journal order. fromOffset (event count) or fromByteOffset (file position) skip already-seen history; fromByteOffset takes priority when non-zero (pass 0 for fromOffset if you only have a byte offset). Already-written events after that point are sent first, then each new event as it is persisted. The channel closes when ctx is cancelled or the engine closes; cancelling the watch does not stop the run. A slow consumer may miss live events — a warning is logged and delivery continues; reconnect with the last Offset/ByteOffset seen to catch up from the journal.
type Option ¶
type Option func(*engineConfig, *readOnlyConfig)
Option configures NewEngine and NewReadOnlyEngine. Writer-only options (WithAutoPurge, WithMaxRetries, WithTimeout, WithStepTokenKey, WithDefaultStepTokenTTL) are ignored by NewReadOnlyEngine.
func WithAutoPurge ¶
WithAutoPurge starts a background goroutine that deletes Completed and Failed runs whose UpdatedAt is older than age. One pass runs at NewEngine (so a process that restarts more often than interval still purges). The optional interval controls later ticks; it defaults to one hour. Running and Waiting runs are never purged. Combine with WithAutoPurgeMaxRuns / WithAutoPurgeMaxBytes to cap disk. Ignored by NewReadOnlyEngine.
func WithAutoPurgeMaxBytes ¶ added in v0.1.5
WithAutoPurgeMaxBytes deletes the oldest Completed/Failed runs when terminal-run directories exceed n bytes on disk. Zero (default) means no size cap. Starts the purger even when WithAutoPurge is omitted. Ignored by NewReadOnlyEngine.
func WithAutoPurgeMaxRuns ¶ added in v0.1.5
WithAutoPurgeMaxRuns deletes the oldest Completed/Failed runs when the number of those terminal runs exceeds n. Running and Waiting runs are never counted or removed. Zero (default) means no count cap. Starts the purger even when WithAutoPurge is omitted. Ignored by NewReadOnlyEngine.
func WithDefaultStepTokenTTL ¶ added in v0.1.5
WithDefaultStepTokenTTL sets how long HMAC StepToken values remain valid. Ignored when no step-token key is configured. Zero means no expiry. Omit this option to use the default of 24 hours when a key is set. Overridden per step by WithStepTokenTTL. Ignored by NewReadOnlyEngine.
func WithJournalMACKey ¶ added in v0.1.5
WithJournalMACKey signs journal.log frames and the run sidecar files (input.json, output.json, meta.json) with HMAC-SHA256. Frame MACs are bound to taskID, runID, and the 1-based frame index (journal v2) so a journal.log cannot be copied between runs or have its frames reordered. Sidecar MACs are bound to taskID/runID. Omit it (or pass nil/empty) to keep CRC32 journal trailers and unsigned JSON sidecars, the default. Enabling a MAC on an existing CRC tree is not a migrate. The key is copied and never written under dataDir. Required on NewReadOnlyEngine to read a MAC-signed tree.
func WithLockTimeout ¶ added in v0.1.2
WithLockTimeout sets how long NewEngine waits for the exclusive flock and NewReadOnlyEngine waits for the shared flock. Default is 2 seconds.
func WithLogger ¶
WithLogger sets the slog.Logger for NewEngine and NewReadOnlyEngine. If nil or omitted, a discard logger is used.
func WithMaxRetries ¶
WithMaxRetries sets the engine-wide default for task-level retries (re-invoking the task closure). Default is 0 — retries are opt-in so non-idempotent work is not silently repeated. Overridden by WithTaskMaxRetries and WithRunMaxRetries. Ignored by NewReadOnlyEngine.
func WithPayloadCodec ¶ added in v0.1.5
func WithPayloadCodec(c PayloadCodec) Option
WithPayloadCodec sets the codec for task input/output, step input/result, and CompleteStep signal payloads on NewEngine and NewReadOnlyEngine. nil is identity (plaintext).
func WithStepTokenKey ¶ added in v0.1.5
WithStepTokenKey sets the key used to sign StepToken values that CompleteStep accepts. Pass one when tokens must stay valid across a restart: without it NewEngine signs with a per-process random key, so tokens issued before a restart are rejected afterwards. The key is copied and never written under dataDir. Ignored by NewReadOnlyEngine.
func WithTimeout ¶
WithTimeout sets the engine-wide default task deadline. Zero (default) means no timeout. Overridden by WithTaskTimeout and WithRunTimeout. Ignored by NewReadOnlyEngine.
func WithUnsignedStepTokens ¶ added in v0.1.5
func WithUnsignedStepTokens() Option
WithUnsignedStepTokens opts out of authenticated step tokens: StepToken returns a plain taskID:runID:stepID triple with no signature and no expiry. Anyone who can reach CompleteStep can mint one for any run and step, so use this only where the CompleteStep caller is already authenticated by other means and tokens must survive a restart. WithStepTokenKey gives you both properties and takes precedence over this option. Ignored by NewReadOnlyEngine.
type PayloadCodec ¶ added in v0.1.5
type PayloadCodec interface {
Encode(plaintext, aad []byte) ([]byte, error)
Decode(ciphertext, aad []byte) ([]byte, error)
}
PayloadCodec transforms task and step payload bytes at the persistence boundary. Encode runs after JSON marshal and before the bytes are written; Decode runs after a read and before JSON unmarshal. Implementations must be reversible — hashing or redaction here would break resume.
aad binds the blob to its location (kind, taskID, runID, stepID) so a ciphertext cannot be copied between fields or steps. Omit WithPayloadCodec (or pass nil) to persist JSON plaintext, the default.
func NewAESGCMCodec ¶ added in v0.1.5
func NewAESGCMCodec(key []byte) (PayloadCodec, error)
NewAESGCMCodec returns a PayloadCodec that wraps payloads as version | nonce | ciphertext+tag (AES-GCM v1). key must be 16, 24, or 32 bytes (AES-128/192/256). The caller owns key storage (env / secret manager); this value is copied and never persisted by the engine.
Prefer NewAESGCMCodecWithKeys when you need a key ID on the wire and a dual-key read window for rotation.
func NewAESGCMCodecWithKeys ¶ added in v0.1.5
func NewAESGCMCodecWithKeys(current AESGCMKey, previous ...AESGCMKey) (PayloadCodec, error)
NewAESGCMCodecWithKeys encodes with current (v2: version | keyID | nonce | sealed) and can decode current, any previous key, and v1 blobs (treated as key ID 0). Pass the retiring key as previous so in-flight journals keep reading during a rotation. Rotate before AESGCMRekeyAfter seals on one key.
type ReadOnlyEngine ¶ added in v0.1.2
type ReadOnlyEngine struct {
// contains filtered or unexported fields
}
ReadOnlyEngine is a compile-time-restricted view of a dataDir. It acquires a shared flock so multiple readers can coexist. It cannot register or run tasks. Safe to open from a separate CLI process with no task registration.
func NewReadOnlyEngine ¶ added in v0.1.2
func NewReadOnlyEngine(dataDir string, opts ...Option) (*ReadOnlyEngine, error)
NewReadOnlyEngine acquires a shared flock on <dataDir>/.lock. Multiple readers coexist. Returns ErrEngineLocked after the lock timeout if a writer holds the exclusive lock.
func (*ReadOnlyEngine) Close ¶ added in v0.1.2
func (r *ReadOnlyEngine) Close() error
Close releases the shared flock. Safe to call more than once.
func (*ReadOnlyEngine) GetStep ¶ added in v0.1.3
func (r *ReadOnlyEngine) GetStep(ctx context.Context, taskID, runID, stepID string) (StepRecord, bool, error)
GetStep returns the latest StepRecord for one stepID. Same semantics as Engine.GetStep.
func (*ReadOnlyEngine) GetTask ¶ added in v0.1.2
GetTask returns a single run's metadata. (zero, false, nil) if not found.
func (*ReadOnlyEngine) ListTasks ¶ added in v0.1.2
func (r *ReadOnlyEngine) ListTasks(ctx context.Context, statuses ...TaskStatus) ([]TaskInfo, error)
ListTasks returns TaskInfo records with the same filter semantics as Engine.ListTasks.
func (*ReadOnlyEngine) ListTasksPage ¶ added in v0.1.5
func (r *ReadOnlyEngine) ListTasksPage(ctx context.Context, offset, limit int, statuses ...TaskStatus) (TaskPage, error)
ListTasksPage is ListTasksPage for a read-only engine.
func (*ReadOnlyEngine) LoadInput ¶ added in v0.1.4
LoadInput returns the JSON-encoded task input written on first RunTask (one I value). Same semantics as Engine.LoadInput.
func (*ReadOnlyEngine) LoadSteps ¶ added in v0.1.2
func (r *ReadOnlyEngine) LoadSteps(ctx context.Context, taskID, runID string) ([]StepRecord, error)
LoadSteps returns all StepRecords for a run. Same semantics as Engine.LoadSteps.
type RunOption ¶ added in v0.1.2
type RunOption func(*runConfig)
RunOption configures a single RunTask call.
func WithRunMaxRetries ¶ added in v0.1.2
WithRunMaxRetries overrides the task and engine task-level retry count for this run. An explicit 0 disables a non-zero parent default.
func WithRunTimeout ¶ added in v0.1.2
WithRunTimeout overrides the task and engine timeout for this run. An explicit 0 disables a non-zero parent timeout.
type StepEvent ¶ added in v0.1.3
type StepEvent struct {
StepRecord
Offset int
ByteOffset int64
// Generation is TaskInfo.JournalGeneration when this event was
// observed. compactJournal increments it and rewrites offsets from 1,
// so a WatchSteps cursor is only valid while Generation is unchanged.
Generation int
}
StepEvent is one journal entry delivered by WatchSteps: STARTED (running), WAITING, COMPLETED, or FAILED — every append, in journal order (never last-write-wins). Offset is the 1-based position of this event across the run's whole journal (steps and signals share one counter); ByteOffset is the file position immediately after it. Reconnect with either value as fromOffset/fromByteOffset to resume a watch without missing or repeating events; byteOffset skips a full rescan, offset does not.
type StepOption ¶ added in v0.1.2
type StepOption func(*stepConfig)
StepOption configures a single RunStep call.
func WithStepMaxRetries ¶ added in v0.1.2
func WithStepMaxRetries(n int) StepOption
WithStepMaxRetries sets how many times this step's function is re-invoked on a non-panic, non-ErrStepPending error. Default is 0 (one attempt). Does not inherit from task or engine retry settings.
func WithStepTimeout ¶ added in v0.1.2
func WithStepTimeout(d time.Duration) StepOption
WithStepTimeout is an inner bound: the step deadline is min(task deadline, step timeout). It cannot extend past the task deadline. The deadline only cancels the step's ctx — it does not forcibly stop fn. A fn that never checks ctx keeps running past the deadline in the background, and the engine (Close, a later CancelRun's drain, etc.) still waits for it to actually return.
func WithStepTokenTTL ¶ added in v0.1.5
func WithStepTokenTTL(d time.Duration) StepOption
WithStepTokenTTL overrides the engine StepToken TTL for this step. Ignored when no step-token key is configured. Zero means no expiry for this step.
func WithStepVersion ¶ added in v0.1.4
func WithStepVersion(v string) StepOption
WithStepVersion tags this step so a later deploy can opt into re-running it. Resume still returns the cached result when the stored version equals v. If this step's code or params change and in-flight runs must re-execute it, pass a new v ("1" → "2") or rename the stepID. Empty v is a no-op (cache by stepID only). Side effects on re-run are the caller's problem.
type StepRecord ¶
type StepRecord struct {
StepID string
Version string
Input []byte
Result []byte
Error string
PanicTrace string
Status StepStatus
StartedAt time.Time
CompletedAt time.Time
}
StepRecord is the latest persistent checkpoint for one memoised step.
type StepRun ¶ added in v0.1.2
type StepRun[O any] struct { // contains filtered or unexported fields }
StepRun is the handle returned by RunStep. RunStep starts work and returns immediately; Get waits for the result, Done reports readiness without blocking so callers can select across several handles (first-of-N).
func RunStep ¶ added in v0.1.2
func RunStep[I, O any](ctx context.Context, s *StepRunner, stepID string, in I, fn func(ctx context.Context, in I) (O, error), opts ...StepOption) *StepRun[O]
RunStep starts fn as a memoised checkpoint and returns immediately. Get waits for the result. On a completed or failed cache hit, fn is not called. On a miss, fn runs on a goroutine and the result is persisted. in is a single JSON-marshalable value (one struct of fields, or struct{} if the step has no payload); it is stored on the record for inspect. stepID must be unique within the run — a duplicate panics. Concurrent calls on the same StepRunner are safe: start several steps, then Get them in any order, or select on their Done channels for first-of-N.
Step IDs are the resume key: never rename a stepID once a run has started. A renamed step is treated as a new, unrelated step — the old result is orphaned and fn runs again under the new name. To re-run the same stepID after a code or param change, pass WithStepVersion with a new string.
func (*StepRun[O]) Done ¶ added in v0.1.3
func (r *StepRun[O]) Done() <-chan struct{}
Done reports readiness. It is closed once the step completes or fails — on a cache hit it is already closed when RunStep returns. Use it in a select across multiple StepRun handles to react to whichever finishes first, then call Get to retrieve that handle's result or error.
type StepRunner ¶
type StepRunner struct {
// contains filtered or unexported fields
}
StepRunner is scoped to a single run and provides RunStep, StepToken, and run-context accessors. Concurrent RunStep calls are safe: fan out work by calling RunStep multiple times before Get-ing any handle.
func (*StepRunner) Logger ¶ added in v0.1.2
func (s *StepRunner) Logger() *slog.Logger
Logger returns the engine logger pre-scoped with taskID and runID.
func (*StepRunner) RunID ¶ added in v0.1.2
func (s *StepRunner) RunID() string
RunID returns the runID of the current run.
func (*StepRunner) StepToken ¶ added in v0.1.2
func (s *StepRunner) StepToken(ctx context.Context) string
StepToken returns an opaque token encoding taskID, runID, and the current stepID. Must be called inside a RunStep function with that function's ctx (the stepID is carried on ctx, not on the StepRunner, so concurrent steps each get their own token). Pass the token to an external caller so they can later call CompleteStep. With WithStepTokenKey the token is HMAC-signed and may expire (default 24h, or WithDefaultStepTokenTTL / WithStepTokenTTL).
func (*StepRunner) TaskID ¶ added in v0.1.2
func (s *StepRunner) TaskID() string
TaskID returns the taskID of the current run.
type StepStatus ¶
type StepStatus string
StepStatus is the lifecycle state of a single step checkpoint.
const ( // StepStatusRunning means the step function is currently executing. // Watch-only: it is superseded by a terminal or waiting status once the // step finishes, and is never the record used to replay RunStep/GetStep — // an unfinished StepStatusRunning after a crash is treated as missing // (fn re-runs), not stuck. StepStatusRunning StepStatus = "running" // StepStatusWaiting means the step is suspended and awaiting CompleteStep. StepStatusWaiting StepStatus = "waiting" // StepStatusCompleted means the step succeeded and its result is cached. StepStatusCompleted StepStatus = "completed" // StepStatusFailed means the step returned an error or panicked. StepStatusFailed StepStatus = "failed" )
type Task ¶
type Task[I, O any] interface { // Exec performs the task logic. ctx is cancelled when the task timeout // elapses or the engine is closed. s is the StepRunner bound to this // run; wrap all memoised work in RunStep calls. Panics are recovered // by the engine, recorded in TaskInfo.PanicTrace, and surfaced to Get. Exec(ctx context.Context, s *StepRunner, input I) (O, error) }
Task is the execution contract for a durable task. Implementations should treat any work with external side effects as a RunStep; non-deterministic logic outside of RunStep calls may not be replayed correctly. I is the single input type (not variadic args); O is the output type. Bundle several fields in one struct. A task with no payload uses I = struct{} and RunTask(..., struct{}{}).
type TaskFunc ¶
type TaskFunc[I, O any] func(ctx context.Context, s *StepRunner, input I) (O, error)
TaskFunc adapts a plain function to Task[I, O].
type TaskInfo ¶
type TaskInfo struct {
TaskID string `json:"task_id"`
RunID string `json:"run_id"`
Name string `json:"name"`
Tags map[string]string `json:"tags"`
Status TaskStatus `json:"status"`
Error string `json:"error"`
PanicTrace string `json:"panic_trace"`
CreatedAt time.Time `json:"created_at"`
StartedAt time.Time `json:"started_at"`
CompletedAt time.Time `json:"completed_at"`
UpdatedAt time.Time `json:"updated_at"`
// OwnerPID is the process that last wrote this meta while the run was
// Running. After a crash it still names the dead process, so operators
// can tell an orphaned in-flight run from one this engine started.
OwnerPID int `json:"owner_pid,omitempty"`
// JournalGeneration increments each time compactJournal rewrites
// journal.log. WatchSteps Offset/ByteOffset are only valid for the
// generation they were observed in.
JournalGeneration int `json:"journal_generation,omitempty"`
}
TaskInfo is the persistent metadata for a single run. The task input is stored separately in input.json (see RunTask), not on this struct.
type TaskOption ¶
type TaskOption func(*taskConfig)
TaskOption configures RegisterTask.
func WithName ¶
func WithName(name string) TaskOption
WithName sets a human-readable label stored on TaskInfo. It does not affect uniqueness or execution.
func WithTag ¶
func WithTag(k, v string) TaskOption
WithTag attaches an arbitrary key-value annotation to TaskInfo. Call multiple times to set multiple tags.
func WithTaskMaxRetries ¶ added in v0.1.2
func WithTaskMaxRetries(n int) TaskOption
WithTaskMaxRetries overrides the engine default for task-level retries (re-invoking the whole closure). Nil-vs-set is tracked so an explicit 0 disables a non-zero engine default.
func WithTaskTimeout ¶ added in v0.1.2
func WithTaskTimeout(d time.Duration) TaskOption
WithTaskTimeout overrides the engine default task deadline. An explicit 0 disables a non-zero engine default.
type TaskRun ¶ added in v0.1.2
type TaskRun[O any] struct { // contains filtered or unexported fields }
TaskRun is the handle returned by RunTask. RunID is available immediately; Get blocks until the run reaches a terminal state.
func RunTask ¶ added in v0.1.2
func RunTask[I, O any](ctx context.Context, e *Engine, taskID string, runID string, input I, opts ...RunOption) *TaskRun[O]
RunTask starts or resumes a run in a background goroutine and returns immediately. input is one typed value I — not variadic args. Put several fields on one struct; a task with no payload uses struct{} and struct{}{}. On first start the value is written to input.json. The same runID reloads that file and ignores the input argument. Pass an empty runID to resume the oldest active run for taskID, or to generate a new ULID if none is active. A completed or failed run returns the stored result without spawning a goroutine.
func (*TaskRun[O]) Get ¶ added in v0.1.2
Get blocks until the run completes or ctx is cancelled. The typed result is cached after the first successful wait so later calls do not re-decode.
func (*TaskRun[O]) RunID ¶ added in v0.1.2
RunID returns the resolved runID. Empty if the taskID was not registered or the runID was invalid — check Get for the error.
func (*TaskRun[O]) Status ¶ added in v0.1.2
func (r *TaskRun[O]) Status() TaskStatus
Status returns the current run status without blocking.
type TaskStatus ¶
type TaskStatus string
TaskStatus is the lifecycle state of a run.
const ( // StatusRunning means the run is actively executing. StatusRunning TaskStatus = "running" // StatusWaiting means at least one step is awaiting CompleteStep. StatusWaiting TaskStatus = "waiting" // StatusCompleted means the run finished successfully. StatusCompleted TaskStatus = "completed" // StatusFailed means the run terminated with an error or recovered panic. StatusFailed TaskStatus = "failed" )
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
cmd
|
|
|
durable-inspect
command
Command durable-inspect is a read-only viewer for a durable-go journal.
|
Command durable-inspect is a read-only viewer for a durable-go journal. |
|
Package durablepb contains the protobuf-generated wire types for journal.log entries.
|
Package durablepb contains the protobuf-generated wire types for journal.log entries. |