Skip to content

API reference

The public API is small on purpose. Everything you need to write workflows is four callables and one exception; the rest is what you use to observe runs and steps. Rendered from the source.

Core

The four core callables live in everystep.api and are re-exported from the everystep package (lazily, so importing everystep does not import Django). The blocks below document them at their definition site.

Mark a function as a durable workflow.

The function takes the arguments passed to schedule() and combines step calls, sequentially or through parallel(). Workflows are scheduled by name; the return value is the persisted workflow result. Called outside a running workflow, it runs as a plain function.

Source code in everystep/api.py
def workflow(func):
    """Mark a function as a durable workflow.

    The function takes the arguments passed to ``schedule()`` and combines
    ``step`` calls, sequentially or through ``parallel()``. Workflows are
    scheduled by name; the return value is the persisted workflow result.
    Called outside a running workflow, it runs as a plain function.
    """
    if not callable(func):
        raise TypeError("@workflow must decorate a function")

    @functools.wraps(func)
    def wrapper(*args, **kwargs):
        return func(*args, **kwargs)

    wrapper.__everystep__ = _WORKFLOW
    return registry.register_workflow(wrapper)

Mark a function as a durable step.

Inside a running workflow, the call is recorded: the outcome (result or exception) is persisted in SQL and served from the store on replay, so the function body only runs for unrecorded steps. Called outside a running workflow, it runs as a plain function with no recording. Pass everystep_id=... at the call site for a stable step identity.

With unsafe_to_repeat=True, the step's side effect must not happen twice. The engine records the step as started before running it, and if a replay finds a started step without an outcome — the worker died in the effect window — it does not re-execute the step: the run ends in the blocked status until a human resolves it with the everystep_resolve_step management command.

Source code in everystep/api.py
def step(func=None, *, unsafe_to_repeat=False):
    """Mark a function as a durable step.

    Inside a running workflow, the call is recorded: the outcome (result or
    exception) is persisted in SQL and served from the store on replay, so
    the function body only runs for unrecorded steps. Called outside a
    running workflow, it runs as a plain function with no recording. Pass
    ``everystep_id=...`` at the call site for a stable step identity.

    With ``unsafe_to_repeat=True``, the step's side effect must not happen
    twice. The engine records the step as started before running it, and if
    a replay finds a started step without an outcome — the worker died in
    the effect window — it does not re-execute the step: the run ends in
    the ``blocked`` status until a human resolves it with the
    ``everystep_resolve_step`` management command.
    """
    if func is not None and not callable(func):
        raise TypeError("@step must decorate a function")

    def decorate(fn):
        @functools.wraps(fn)
        def wrapper(*args, **kwargs):
            everystep_id = kwargs.pop("everystep_id", None)
            ctx = context.current()
            if ctx is None:
                return fn(*args, **kwargs)
            return run_step(ctx, fn, args, kwargs, everystep_id, unsafe_to_repeat)

        wrapper.__everystep__ = _STEP
        return registry.register_step(wrapper)

    if func is not None:
        return decorate(func)
    return decorate

Run zero-arg callables concurrently; return their results in order.

Each branch is a single step (a @step function or a lambda calling one), or a list of zero-arg callables run in order within the branch; such a sequence branch returns a tuple of its steps' results. If any branch raises, the first error is re-raised (single branch) or an ExceptionGroup is raised (several branches).

Pass everystep_id="..." to give the fork a stable name, so the branch step ids ("name.0.1", "name.1.1") do not shift when steps before it change.

Source code in everystep/api.py
def parallel(*branches, everystep_id=None):
    """Run zero-arg callables concurrently; return their results in order.

    Each branch is a single step (a @step function or a lambda calling one),
    or a list of zero-arg callables run in order within the branch; such a
    sequence branch returns a tuple of its steps' results.
    If any branch raises, the first error is re-raised (single branch) or an
    ExceptionGroup is raised (several branches).

    Pass everystep_id="..." to give the fork a stable name, so the branch step
    ids ("name.0.1", "name.1.1") do not shift when steps before it change.
    """
    if not branches:
        raise ValueError("parallel() requires at least one branch")
    branches = tuple(_as_branch(branch) for branch in branches)

    ctx = context.current()
    owns_ctx = False
    if ctx is None:
        ctx = Context(workflow_id=None, outcomes={}, persistent=False)
        context.set_current(ctx)
        owns_ctx = True
    try:
        fork_id = ctx.next_id(everystep_id)
        with traces.span("parallel", {"everystep.parallel.id": fork_id}, active=ctx.persistent):
            fork_ctx = traces.capture_context() if ctx.persistent else None
            with ThreadPoolExecutor(
                max_workers=len(branches), thread_name_prefix="everystep-parallel"
            ) as pool:
                futures = [
                    pool.submit(_run_branch, ctx, fork_id, index, branch, fork_ctx)
                    for index, branch in enumerate(branches)
                ]
                outputs = [future.result() for future in futures]
    finally:
        if owns_ctx:
            context.clear_current()

    errors = [output for output in outputs if isinstance(output, _Failed)]
    if not errors:
        return tuple(output.value for output in outputs)
    drains = [e for e in errors if isinstance(e.exc, (DrainOrphan, EffectUncertain, Terminal))]
    if drains:
        raise drains[0].exc
    if len(errors) == 1:
        raise errors[0].exc
    raise ExceptionGroup("everystep parallel branches failed", [e.exc for e in errors])

Insert a scheduled workflow into the current transaction.

The workflow becomes claimable by workers once the transaction commits. With an idempotency_key, a concurrent or repeated schedule with the same key returns the existing workflow instead of creating a new one.

Source code in everystep/api.py
def schedule(workflow_func, *args, idempotency_key=None):
    """Insert a scheduled workflow into the current transaction.

    The workflow becomes claimable by workers once the transaction commits.
    With an idempotency_key, a concurrent or repeated schedule with the same
    key returns the existing workflow instead of creating a new one.
    """
    if getattr(workflow_func, "__everystep__", None) != _WORKFLOW:
        raise TypeError(f"{workflow_func!r} is not a @workflow-decorated function")
    try:
        serde.dumps(list(args))
    except (TypeError, ValueError) as exc:
        raise EverystepError(f"workflow arguments are not serializable: {exc}") from exc

    defaults = {"args": list(args), "status": Workflow.Status.SCHEDULED}
    if idempotency_key is None:
        return Workflow.objects.create(name=name_of(workflow_func), **defaults)
    run, _created = Workflow.objects.get_or_create(
        name=name_of(workflow_func), idempotency_key=idempotency_key, defaults=defaults
    )
    return run

Bases: EverystepError

Raised from a step to stop the workflow for a known reason.

The step in flight finishes and is recorded; no new step starts and the engine will not retry the workflow. It ends in the stopped status with the reason and payload recorded on it, so the caller can take over (for example, re-scheduling with a different resource). Not a failure: it is never reported to Sentry.

Source code in everystep/errors.py
class Terminal(EverystepError):
    """Raised from a step to stop the workflow for a known reason.

    The step in flight finishes and is recorded; no new step starts and the
    engine will not retry the workflow. It ends in the `stopped` status with
    the reason and payload recorded on it, so the caller can take over (for
    example, re-scheduling with a different resource). Not a failure: it is
    never reported to Sentry.
    """

    def __init__(self, reason, payload=None):
        self.reason = reason
        self.payload = payload
        super().__init__(reason, payload)

    def __str__(self):
        return str(self.reason)

Execution context

The execution context of the current thread, or None outside a run.

Source code in everystep/context.py
def current():
    """The execution context of the current thread, or None outside a run."""
    return getattr(_local, "ctx", None)

Execution state for one workflow run, or one branch of one fork.

Steps are identified by a dotpath like "3" or "3.1.2". Unnamed steps take the next position in their scope; a step called with everystep_id takes that name as its segment instead, which makes its id stable against edits elsewhere in the body. Names must be unique within a scope (one body, or one branch). The shared outcomes dict maps step ids to {name, status, result, error} outcomes. It is seeded with all outcomes recorded before the current claim and updated as each step completes, so inside a step you can read any previously-completed step's outcome from the current context. Treat it as read-only.

Source code in everystep/context.py
class Context:
    """Execution state for one workflow run, or one branch of one fork.

    Steps are identified by a dotpath like "3" or "3.1.2". Unnamed steps
    take the next position in their scope; a step called with `everystep_id`
    takes that name as its segment instead, which makes its id stable
    against edits elsewhere in the body. Names must be unique within a
    scope (one body, or one branch). The shared `outcomes` dict maps step
    ids to `{name, status, result, error}` outcomes. It is seeded with all
    outcomes recorded before the current claim and updated as each step
    completes, so inside a step you can read any previously-completed
    step's outcome from the current context. Treat it as read-only.
    """

    def __init__(self, workflow_id, outcomes, persistent, prefix="", draining=None, workflow_name=None):
        self.workflow_id = workflow_id
        self.workflow_name = workflow_name
        self.outcomes = outcomes
        self.persistent = persistent
        self.prefix = prefix
        # Event set by the runner on stop; checked before each new step.
        self.draining = draining
        self.counter = itertools.count(1)
        self._used_ids = set()

    def next_id(self, name=None):
        # Always consume a position, so naming (or un-naming) a step never
        # shifts the ids of the other unnamed steps in the scope.
        tick = next(self.counter)
        if name is None:
            segment = str(tick)
        else:
            _validate_name(name)
            segment = name
        step_id = f"{self.prefix}{segment}"
        if len(step_id) > _MAX_ID_LENGTH:
            raise EverystepError(f"everystep_id {name!r} is too long: step id would exceed {_MAX_ID_LENGTH} characters")
        if step_id in self._used_ids:
            raise EverystepError(
                f"everystep_id {segment!r} is already used in this scope (step id "
                f"{step_id!r}); everystep_ids must be unique per body or branch"
            )
        self._used_ids.add(step_id)
        return step_id

    def branch(self, fork_id, index):
        return Context(
            workflow_id=self.workflow_id,
            outcomes=self.outcomes,
            persistent=self.persistent,
            prefix=f"{fork_id}.{index}.",
            draining=self.draining,
            workflow_name=self.workflow_name,
        )

Exceptions

Bases: Exception

Base class for everystep errors.

Source code in everystep/errors.py
class EverystepError(Exception):
    """Base class for everystep errors."""

Bases: EverystepError

A workflow body diverged from its previously recorded step identities.

Source code in everystep/errors.py
class WorkflowCodeError(EverystepError):
    """A workflow body diverged from its previously recorded step identities."""

Bases: EverystepError

A step marked unsafe to repeat may have performed its effect.

The step's outcome was never recorded (the worker died in the effect window, or another claimant holds it), so the engine cannot tell whether the effect happened. It refuses to execute the step again and ends the run in the blocked status — a holding state, not a failure: it is never reported to Sentry. The run is resolved by a human with the everystep_resolve_step management command.

Source code in everystep/errors.py
class EffectUncertain(EverystepError):
    """A step marked unsafe to repeat may have performed its effect.

    The step's outcome was never recorded (the worker died in the effect
    window, or another claimant holds it), so the engine cannot tell whether
    the effect happened. It refuses to execute the step again and ends the
    run in the `blocked` status — a holding state, not a failure: it is
    never reported to Sentry. The run is resolved by a human with the
    `everystep_resolve_step` management command.
    """

Bases: EverystepError

A recorded step failure whose original exception type is unavailable.

Source code in everystep/errors.py
class StepFailure(EverystepError):
    """A recorded step failure whose original exception type is unavailable."""

Bases: EverystepError

Raised from a test fault handler to simulate a process death mid-step.

The runner re-raises it without writing any further state, leaving the workflow exactly as it would be if the worker had died.

Source code in everystep/errors.py
class SimulatedCrash(EverystepError):
    """Raised from a test fault handler to simulate a process death mid-step.

    The runner re-raises it without writing any further state, leaving the
    workflow exactly as it would be if the worker had died.
    """

Bases: EverystepError

Raised at a step boundary when the runner is draining after a stop signal.

The step in flight at the signal finishes and is recorded; no new step starts. The worker requeues the workflow so any runner can claim it and resume it from the recorded steps.

Source code in everystep/errors.py
class DrainOrphan(EverystepError):
    """Raised at a step boundary when the runner is draining after a stop signal.

    The step in flight at the signal finishes and is recorded; no new step
    starts. The worker requeues the workflow so any runner can claim it and
    resume it from the recorded steps.
    """

Stored exceptions

Encode an exception into a JSON-serializable dict for storage.

Source code in everystep/serde.py
def encode_exception(exc):
    """Encode an exception into a JSON-serializable dict for storage."""
    payload = {
        "type": f"{type(exc).__module__}.{type(exc).__qualname__}",
        "message": str(exc),
    }
    try:
        payload["args"] = json.loads(dumps(list(exc.args)))
    except TypeError:
        payload["args"] = []
    return payload

Decode a stored exception dict into an exception instance.

Returns the original exception type when it is importable, otherwise a StepFailure carrying the stored message.

Source code in everystep/serde.py
def decode_exception(payload):
    """Decode a stored exception dict into an exception instance.

    Returns the original exception type when it is importable, otherwise a
    StepFailure carrying the stored message.
    """
    if not isinstance(payload, dict):
        return StepFailure(f"step failed (unrecognized failure record: {payload!r})")
    message = payload.get("message", "")
    type_path = payload.get("type", "")
    args = payload.get("args") or []
    exc_type = import_dotted(type_path) if type_path else None
    if exc_type is not None and isinstance(exc_type, type) and issubclass(exc_type, BaseException):
        attempts = [args] if args else ([message] if message else [])
        for attempt in attempts:
            try:
                return exc_type(*attempt)
            except Exception:
                continue
    if message:
        return StepFailure(f"{message} (original exception type {type_path!r} unavailable)")
    return StepFailure("step failed (original exception unavailable)")

Worker

Claim due workflows and execute them on a thread pool.

See the running section of the documentation for the claim loop, the name contract, and the SIGTERM drain behavior.

Source code in everystep/worker.py
class Worker:
    """Claim due workflows and execute them on a thread pool.

    See the running section of the documentation for the claim loop, the
    name contract, and the SIGTERM drain behavior.
    """

    def __init__(
        self, pool_size=4, poll=0.2, name=None, drain=30, metrics_port=0, metrics_bind="0.0.0.0"
    ):
        self.pool_size = pool_size
        self.poll = poll
        self.name = name or socket.gethostname()
        self.drain = drain
        self.metrics_port = metrics_port
        self.metrics_bind = metrics_bind
        self._stop = threading.Event()
        self._draining = threading.Event()
        self._active = {}
        self._executor = ThreadPoolExecutor(max_workers=pool_size, thread_name_prefix="everystep-w")
        self._metrics_server = None

    def stop(self):
        self._stop.set()
        self._draining.set()

    def run(self):
        _install_signal_handlers(self)
        try:
            metrics.worker_started(self.name, self.pool_size)
            self._start_metrics_server()
            self._catchup()
            while not self._stop.is_set():
                self._reap()
                capacity = self.pool_size - len(self._active)
                if capacity > 0:
                    claimed = claim_new(capacity, self.name)
                    for workflow in claimed:
                        self._active[self._executor.submit(self._execute, workflow)] = workflow
                    metrics.record_claims(self.name, len(claimed))
                metrics.set_inflight(self.name, len(self._active))
                self._stop.wait(self.poll)
        finally:
            leftovers = self._drain()
            if leftovers:
                self._abandon(leftovers)
                if threading.current_thread() is threading.main_thread():
                    # Standalone worker: exit now rather than let the
                    # interpreter join the abandoned pool threads at
                    # shutdown. An embedded worker just returns; the host
                    # process owns its own lifecycle.
                    os._exit(0)
            metrics.set_inflight(self.name, 0)
            if self._metrics_server is not None:
                self._metrics_server.server_close()
            connections.close_all()

    def _start_metrics_server(self):
        if not self.metrics_port:
            return
        self._metrics_server = metrics.start_http_server(self.metrics_port, self.metrics_bind)

    def _drain(self):
        """Stop queued work and wait up to `self.drain` seconds for the
        in-flight workflows to finish. Returns the futures still running
        when the deadline is hit. With drain of 0, waits indefinitely and
        returns nothing."""
        if not self.drain or not self._active:
            self._executor.shutdown(wait=True)
            return []
        self._executor.shutdown(wait=False, cancel_futures=True)
        _, not_done = wait(list(self._active), timeout=self.drain)
        return list(not_done)

    def _abandon(self, leftovers):
        """Requeue the runs still in flight when the drain deadline expired,
        so nothing is left claimed by this worker."""
        for future in leftovers:
            self._requeue(self._active[future].id)
        metrics.record_requeues(self.name, len(leftovers))
        logger.warning(
            "everystep worker: drain deadline of %ss expired with %d workflow(s) still "
            "in flight; they were requeued and will be picked up by the next available "
            "runner",
            self.drain, len(leftovers),
        )

    def _requeue(self, workflow_id):
        """Put a run back in the queue: scheduled and unclaimed, so any
        runner can claim it and resume it from the recorded steps."""
        updated = Workflow.objects.filter(
            id=workflow_id,
            status=Workflow.Status.RUNNING,
            claimed_by=self.name,
        ).update(status=Workflow.Status.SCHEDULED, claimed_by=None)
        if not updated:
            logger.warning(
                "everystep worker: could not requeue workflow %s: no longer claimed by %r",
                workflow_id, self.name,
            )

    def _catchup(self):
        """Re-claim the runs a crashed process with this name left behind.

        Runs once at startup, before the poll loop, so it cannot re-select
        workflows this process is already executing in its pool. After a
        clean shutdown it finds nothing: in-flight runs are requeued at
        shutdown.
        """
        workflows = resume_own(self.pool_size, self.name)
        for workflow in workflows:
            self._active[self._executor.submit(self._execute, workflow)] = workflow
        metrics.record_claims(self.name, len(workflows))

    def _reap(self):
        pending = {}
        for future, workflow in self._active.items():
            if future.done():
                exc = future.exception()
                if exc is not None:
                    logger.exception("everystep worker: unexpected worker failure: %s", exc)
            else:
                pending[future] = workflow
        self._active = pending

    def _execute(self, workflow):
        started = time.monotonic()
        try:
            execute(workflow.id, draining=self._draining)
        except Workflow.DoesNotExist:
            logger.warning("everystep worker: workflow %s no longer exists", workflow.id)
        except DrainOrphan:
            # Drained at a step boundary with everything so far recorded:
            # give the run back to the queue for any runner to resume.
            self._requeue(workflow.id)
        except Exception as exc:
            logger.exception("everystep worker: workflow %s crashed outside the runner", workflow.id)
            report_workflow_failure(exc, workflow_id=workflow.id)
            try:
                # Guarded on claimed_by: if this run was requeued while a
                # step was still in flight (drain deadline) and since
                # claimed by another runner, its late failure must not fail
                # that runner's live run.
                updated = Workflow.objects.filter(
                    id=workflow.id,
                    status=Workflow.Status.RUNNING,
                    claimed_by=self.name,
                ).update(
                    status=Workflow.Status.FAILED,
                    error=serde.encode_exception(exc),
                    completed_at=timezone.now(),
                )
                if updated and not isinstance(exc, (SimulatedCrash, Terminal)):
                    metrics.record_workflow_terminal(
                        workflow.name, "failed", time.monotonic() - started
                    )
            except Exception:
                logger.exception("everystep worker: could not mark workflow %s failed", workflow.id)
        finally:
            # Pool threads are long-lived and Django connections are
            # thread-local, so release this thread's connection to avoid
            # leaking one per executed workflow.
            connections.close_all()