Zum Inhalt springen
aviral gupta

// Projekt Vertiefung · etwa 14 Stunden Arbeit

Asynchroner Job-Runner

Sie bauen jobrunner, ein Kommandozeilenwerkzeug, das die Jobs einer TOML-Datei nebenläufig ausführt, wie ein kleiner Build-Server oder Batch-Planer. Jeder Job ist ein Befehl, der ohne Shell läuft. Eine Semaphore begrenzt, wie viele gleichzeitig laufen, jeder Job hat sein eigenes Zeitlimit, Fehlschläge werden mit exponentiellem Backoff wiederholt, und --fail-fast bricht nach dem ersten Fehlschlag den Rest ab. Im Kern steht ein generischer Runner, typisiert mit der Typparameter-Syntax von Python 3.12, der jeden asynchronen Job ausführt; die meisten Tests nutzen daher Schein-Jobs und dauern Millisekunden. Das Projekt verbindet die ganze Stufe Vertiefung: asyncio mit TaskGroup, Zeitlimits und Abbruch, Generics mit mypy, tomllib, sichere Unterprozesse und cProfile.

Was das fertige Programm kann

  • Job, Result und RetryPolicy sind eingefrorene Dataclasses. Job und Result sind generisch, geschrieben class Job[T], und Runner.run[T] gibt list[Result[T]] zurück, sodass mypy weiß, was jedes Ergebnis enthält.
  • RetryPolicy hat attempts, backoff, factor, max_backoff und retry_on. Die Pause nach Versuch n ist backoff * factor ** (n - 1), höchstens max_backoff. Falsche Einstellungen lösen ValueError aus.
  • Runner(concurrency, fail_fast=False).run(jobs) führt jeden Job in einer asyncio.TaskGroup aus, nie mehr als concurrency gleichzeitig (eine asyncio.Semaphore), und gibt die Ergebnisse in der Reihenfolge der Jobs zurück. Doppelte Namen lösen ValueError aus.
  • Jeder Versuch läuft in asyncio.timeout(job.timeout). Ein Versuch, dem die Zeit ausgeht, ist TIMED_OUT, einer mit Ausnahme FAILED; beide werden wiederholt, solange Versuche übrig sind und retry_on zum Fehler passt.
  • Ein Job, der auf die Wiederholung wartet, schläft außerhalb der Semaphore und belegt so keinen Platz. Die Schlaffunktion wird an Runner übergeben, damit Tests das Backoff aufzeichnen, statt zu warten.
  • Mit fail_fast bricht der erste endgültig gescheiterte Job die anderen ab, die als CANCELLED gemeldet werden; wartende Jobs starten nie. Ein Abbruch von run() von außen bricht jeden Job ab und wird nie verschluckt.
  • Die Jobs stehen in einer TOML-Datei mit [[job]]-Tabellen (name, command, timeout, attempts, backoff) und einer optionalen Tabelle [defaults]. Jeder Fehler löst ConfigError aus, der Datei und Job nennt.
  • Befehle laufen mit asyncio.create_subprocess_exec, nie über eine Shell; ein erstes Wort "python" bedeutet sys.executable. Ein Exit-Status ungleich 0 löst CommandFailed aus, und ein abgebrochener Befehl wird beendet.
  • jobrunner JOBDATEI [-c N] [--fail-fast] [--profile DATEI] [-v] gibt eine Tabelle und eine Zusammenfassung aus und protokolliert Wiederholungen und Fehlschläge auf stderr. main gibt 0 zurück, 1 bei einem gescheiterten Job oder 2 bei fehlerhafter Jobdatei.
  • --profile DATEI führt den ganzen Lauf unter cProfile.Profile aus und speichert die Statistik für pstats. mypy --check-untyped-defs . meldet keine Fehler.

Aufbau der Startdateien

pyproject.toml
Paket-Metadaten, das Build-Backend und der Konsolenbefehl jobrunner. Vollständig.
jobrunner/__init__.py
Kennzeichnet das Paket und enthält __version__. Vollständig.
jobrunner/__main__.py
Damit startet python -m jobrunner die Kommandozeile. Vollständig.
jobrunner/model.py
Status, RetryPolicy, Job und Result. Status ist vollständig; die Methoden von RetryPolicy und Result.ok sind Stubs, und Job und Result sind noch nicht generisch.
jobrunner/runner.py
Die Klasse Runner: run und _run_job sind Stubs. describe, das einen Fehler in einen kurzen Text verwandelt, ist vollständig.
jobrunner/commands.py
Befehls-Jobs und die Jobdatei. CommandOutput, CommandFailed und command_action sind vollständig; run_command und load_jobs sind Stubs.
jobrunner/report.py
Summary, summarise und format_table. Stubs.
jobrunner/cli.py
Der argparse-Parser, die Logging-Einrichtung und main(argv). Die Logging-Einrichtung ist geschrieben; der Rest besteht aus Stubs und TODOs.
samples/jobs.toml
Eine Beispiel-Jobdatei: zwei Jobs, die gelingen, einer, der scheitert und wiederholt wird, und einer, dem die Zeit ausgeht.

Etappen

  1. Etappe 1

    Das typisierte Modell

    Schreiben Sie die Prüfungen von RetryPolicy, delay und should_retry sowie Result.ok. Machen Sie Job und Result mit class Job[T] und class Result[T] generisch, sodass ein Job[int] ein Result[int] ergibt.

    Prüfungen, die nach dieser Etappe bestehen:

    • RetryPolicy verdoppelt die Wartezeit nach jedem Versuch bis max_backoff und lehnt falsche Einstellungen ab
    • Job, Result und Runner.run sind generisch, geschrieben mit der Typparameter-Syntax von Python 3.12
  2. Etappe 2

    Jobs gemeinsam ausführen

    Schreiben Sie Runner.run: Namen prüfen, eine Semaphore anlegen (Lektion A3.5) und je Job einen Task in einer asyncio.TaskGroup starten. In _run_job läuft jeder Versuch in der Semaphore und in asyncio.timeout; speichern Sie sein Result.

    Prüfungen, die nach dieser Etappe bestehen:

    • run() führt jeden Job aus und gibt die Ergebnisse in der Reihenfolge der Jobs zurück
    • Nie laufen mehr als concurrency Jobs gleichzeitig, und sie überlappen sich
    • Doppelte Jobnamen und eine concurrency unter 1 lösen ValueError aus
    • Ein Job, der länger als sein Zeitlimit braucht, wird gestoppt und als abgelaufen gemeldet
  3. Etappe 3

    Wiederholungen mit Backoff

    Wiederholen Sie, solange Versuche übrig sind und retry_on passt, und warten Sie nach dem Verlassen der Semaphore mit self._sleep(job.retry.delay(attempt)). Behalten Sie den letzten Fehler im Result.

    Prüfungen, die nach dieser Etappe bestehen:

    • Ein fehlschlagender Job wird mit wachsenden Pausen wiederholt; andere Fehler oder keine Versuche mehr lassen ihn scheitern
    • Während ein Job auf die Wiederholung wartet, belegt er keinen Platz, sodass ein anderer Job laufen kann
  4. Etappe 4

    Abbruch

    Fangen Sie asyncio.CancelledError in _run_job ab, speichern Sie CANCELLED und lösen Sie den Fehler erneut aus. Für fail_fast setzen Sie das Event stopping (A3.5) und lösen _FailFast aus, solange Sie den Platz halten; fangen Sie es mit except* um die TaskGroup.

    Prüfungen, die nach dieser Etappe bestehen:

    • Mit fail_fast bricht der erste Fehlschlag die anderen Jobs ab, die als abgebrochen gemeldet werden
    • Ein Abbruch von run() von außen bricht jeden Job ab und wird nicht verschluckt
  5. Etappe 5

    Befehle und Jobdateien

    Schreiben Sie run_command mit asyncio.create_subprocess_exec und communicate (A3.4), und beenden Sie den Prozess mit kill(), wenn der Task abgebrochen wird. Dann load_jobs mit tomllib: [defaults] in jeden Job übernehmen und für jeden Fehler ConfigError auslösen.

    Prüfungen, die nach dieser Etappe bestehen:

    • run_command führt einen Befehl ohne Shell aus und löst CommandFailed aus, wenn er mit einem Fehler endet
    • Ein Befehl, der sein Zeitlimit überschreitet, wird beendet und als abgelaufen gemeldet
    • load_jobs liest [[job]]-Tabellen und ergänzt die Werte aus [defaults]
    • Eine fehlende Datei, falsches TOML oder ein falscher Job lösen ConfigError mit einer klaren Meldung aus
  6. Etappe 6

    Bericht, Kommandozeile und Profil

    Schreiben Sie summarise, Summary.__str__ und format_table, dann die Optionen und main. Umschließen Sie asyncio.run für --profile mit cProfile.Profile, und führen Sie mypy --check-untyped-defs . aus, bis es sauber ist.

    Prüfungen, die nach dieser Etappe bestehen:

    • summarise zählt die Ergebnisse je Status, und format_table gibt eine Zeile je Job aus
    • jobrunner JOBDATEI gibt eine Tabelle und eine Zusammenfassung aus, protokolliert Wiederholungen und gibt 1 zurück, wenn ein Job scheiterte
    • Wenn alle Jobs gelingen, ist das Ergebnis 0; --fail-fast bricht nach einem Fehlschlag den Rest ab
    • Eine fehlerhafte Jobdatei ergibt 2 mit einer Meldung auf stderr; falsche Argumente enden mit Status 2
    • --profile DATEI speichert cProfile-Statistiken, die pstats lesen kann
    • samples/jobs.toml endet mit 2 erfolgreichen, 1 gescheiterten und 1 abgelaufenen Job

Dieses Projekt nutzt Teile von Python, die im Browser nicht laufen. Sie bauen es deshalb auf Ihrem Rechner.

Auf dem eigenen Rechner bauen

Legen Sie einen Ordner mit diesen Startdateien an, installieren Sie Python 3.14 und arbeiten Sie die Etappen ab. Die Abnahmetests starten Sie jederzeit mit:

Starter als eine .zip-Datei herunterladen (Starter-Dateien, test_main.py und learnrun.py)
python learnrun.py test

Unter macOS und Linux tippen Sie python3, wo in diesen Befehlen python steht, wie in der ersten Lektion.

learnrun.py herunterladen

pyproject.toml

[build-system]
requires = ["setuptools >= 77.0.3"]
build-backend = "setuptools.build_meta"

[project]
name = "jobrunner"
version = "1.0.0"
description = "Run jobs concurrently with asyncio: timeouts, retries with backoff, and cancellation."
requires-python = ">= 3.14"
dependencies = []

[project.scripts]
jobrunner = "jobrunner.cli:main"

[tool.setuptools]
packages = ["jobrunner"]

jobrunner/__init__.py

"""Run jobs concurrently with asyncio: timeouts, retries with backoff, and cancellation."""

__version__ = "1.0.0"

jobrunner/__main__.py

"""python -m jobrunner runs the command line."""

import sys

from .cli import main

sys.exit(main())

jobrunner/model.py

"""The domain: jobs, retry policies and results, generic in what a job returns."""

from collections.abc import Awaitable, Callable
from dataclasses import dataclass, field
from enum import StrEnum
from typing import Any


class Status(StrEnum):
    """How a job ended."""

    SUCCEEDED = "succeeded"
    FAILED = "failed"
    TIMED_OUT = "timed out"
    CANCELLED = "cancelled"


@dataclass(frozen=True, slots=True)
class RetryPolicy:
    """How often to try a job, and how long to wait between tries.

    The wait before try n + 1 is backoff * factor ** (n - 1), capped at
    max_backoff: 0.1, 0.2, 0.4, ... with the defaults.
    """

    attempts: int = 1
    backoff: float = 0.1
    factor: float = 2.0
    max_backoff: float = 30.0
    retry_on: tuple[type[Exception], ...] = (Exception,)

    def __post_init__(self) -> None:
        # TODO: raise ValueError if attempts < 1, backoff or max_backoff < 0,
        # or factor < 1.
        pass

    def delay(self, attempt: int) -> float:
        """Seconds to wait after failed try number attempt (counting from 1)."""
        raise NotImplementedError

    def should_retry(self, error: BaseException) -> bool:
        """Whether this kind of error is worth another try (isinstance with retry_on)."""
        raise NotImplementedError


# TODO: make Job generic in what its action returns, with the type parameter
# syntax: class Job[T], and action: Callable[[], Awaitable[T]].
@dataclass(frozen=True, slots=True)
class Job:
    """A named unit of work. action is called afresh for every try."""

    name: str
    action: Callable[[], Awaitable[Any]]
    timeout: float | None = None
    retry: RetryPolicy = field(default_factory=RetryPolicy)


# TODO: make Result generic too: class Result[T], with value: T | None.
@dataclass(frozen=True, slots=True)
class Result:
    """What happened to one job."""

    name: str
    status: Status
    value: Any = None
    error: BaseException | None = None
    attempts: int = 0
    elapsed: float = 0.0

    @property
    def ok(self) -> bool:
        """True if the job succeeded."""
        raise NotImplementedError

jobrunner/runner.py

"""Running jobs concurrently: a TaskGroup, a semaphore, timeouts, retries and cancellation."""

import asyncio
import logging
from collections.abc import Awaitable, Callable, Iterable

from .model import Job, Result, Status

logger = logging.getLogger(__name__)

type Sleep = Callable[[float], Awaitable[None]]


class _FailFast(Exception):
    """Raised inside the TaskGroup to cancel the other jobs after a failure."""


class Runner:
    """Runs jobs at most `concurrency` at a time and reports a Result for each."""

    def __init__(self, concurrency: int = 4, *, fail_fast: bool = False, sleep: Sleep = asyncio.sleep) -> None:
        # TODO: raise ValueError if concurrency is below 1.
        self.concurrency = concurrency
        self.fail_fast = fail_fast
        # Injected so tests can record the backoff instead of waiting for it.
        self._sleep = sleep

    # TODO: make run generic: async def run[T](self, jobs: Iterable[Job[T]]) -> list[Result[T]]
    async def run(self, jobs: Iterable[Job]) -> list[Result]:
        """Run every job; return their results in the order the jobs were given.

        With fail_fast, the first job that fails or times out cancels the
        others, which are reported as cancelled. If the caller cancels run()
        itself, every job is cancelled and the cancellation propagates.
        """
        # TODO:
        # - raise ValueError if two jobs have the same name;
        # - one asyncio.Semaphore(self.concurrency) for all jobs;
        # - create a task per job (self._run_job) in an asyncio.TaskGroup;
        # - catch _FailFast with except* around the TaskGroup;
        # - return the results in the order of the jobs.
        raise NotImplementedError

    async def _run_job(self, job: Job, slots: asyncio.Semaphore, stopping: asyncio.Event, results: dict[str, Result]) -> None:
        """Try job until it succeeds or runs out of attempts; store its Result in results."""
        # TODO:
        # - each try inside `async with slots:` and `async with asyncio.timeout(job.timeout):`;
        # - TimeoutError -> Status.TIMED_OUT, another Exception -> Status.FAILED;
        # - retry if attempts remain and job.retry.should_retry(error), after
        #   `await self._sleep(job.retry.delay(attempt))` OUTSIDE the slot;
        # - with fail_fast, a final failure sets stopping and raises _FailFast
        #   while still holding the slot; a job that finds stopping set when it
        #   gets its slot is cancelled without starting;
        # - on asyncio.CancelledError store Status.CANCELLED and re-raise it.
        raise NotImplementedError


def describe(error: BaseException | None) -> str:
    """A short description of an error, for logs and reports."""
    if error is None:
        return ""
    if isinstance(error, TimeoutError) and not str(error):
        return "timed out"
    return f"{type(error).__name__}: {error}" if str(error) else type(error).__name__

jobrunner/commands.py

"""Jobs that run commands, and loading them from a TOML job file."""

import asyncio
import sys
import tomllib
from collections.abc import Awaitable, Callable, Sequence
from dataclasses import dataclass
from pathlib import Path
from typing import Any

from .model import Job, RetryPolicy


class ConfigError(Exception):
    """The job file is missing, is not valid TOML, or describes a job wrongly."""


@dataclass(frozen=True, slots=True)
class CommandOutput:
    """What a finished command printed, and its exit status."""

    argv: tuple[str, ...]
    returncode: int
    stdout: str
    stderr: str


class CommandFailed(Exception):
    """A command exited with a status other than 0."""

    def __init__(self, output: CommandOutput) -> None:
        last = output.stderr.strip().splitlines()[-1:] or [""]
        super().__init__(f"exit status {output.returncode}" + (f": {last[0]}" if last[0] else ""))
        self.output = output


async def run_command(argv: Sequence[str]) -> CommandOutput:
    """Run a command without a shell and wait for it; raise CommandFailed if it fails.

    A command whose first word is "python" runs with the interpreter running
    this program. If the waiting task is cancelled (by a timeout, say), the
    process is killed before the cancellation goes on.
    """
    # TODO: replace a first word "python" with sys.executable;
    # asyncio.create_subprocess_exec with stdout and stderr piped;
    # await process.communicate(), and on asyncio.CancelledError kill the
    # process, await process.wait() and re-raise; decode the output.
    raise NotImplementedError


def command_action(argv: Sequence[str]) -> Callable[[], Awaitable[CommandOutput]]:
    """A job action that runs argv afresh each time it is called."""
    frozen = tuple(argv)

    async def action() -> CommandOutput:
        return await run_command(frozen)

    return action


JOB_KEYS = {"name", "command", "timeout", "attempts", "backoff"}


def _number(table: dict[str, Any], key: str, where: str, *, minimum: float) -> float | None:
    value = table.get(key)
    if value is None:
        return None
    if isinstance(value, bool) or not isinstance(value, int | float) or value < minimum:
        raise ConfigError(f"{where}: {key} must be a number of at least {minimum}, not {value!r}")
    return float(value)


def load_jobs(path: Path) -> list[Job]:
    """Read the jobs of a TOML job file: an optional [defaults] table and [[job]] tables.

    Raise ConfigError for a file that cannot be read, invalid TOML, an unknown
    key, a job without a name or a list of strings as its command, a bad
    number, or a name used twice.
    """
    # TODO: tomllib.load needs the file opened in binary mode ("rb").
    # Each job: Job(name, command_action(command), timeout, RetryPolicy(...)),
    # with values from [defaults] where the job has none (backoff defaults to 0.5).
    raise NotImplementedError

jobrunner/report.py

"""Summing up a run: counts per status and a table of results."""

from collections import Counter
from collections.abc import Sequence
from dataclasses import dataclass

from .model import Result, Status
from .runner import describe


@dataclass(frozen=True, slots=True)
class Summary:
    """The shape of a finished run."""

    total: int
    counts: dict[Status, int]
    attempts: int
    slowest: str | None

    @property
    def ok(self) -> bool:
        """True if every job succeeded."""
        raise NotImplementedError

    def __str__(self) -> str:
        """For example "4 jobs: 2 succeeded, 1 failed, 1 timed out" (statuses in Status order, none that are 0)."""
        raise NotImplementedError


# TODO: generic in T, like Runner.run.
def summarise(results: Sequence[Result]) -> Summary:
    """Count the results per status, add up the attempts and find the slowest job."""
    raise NotImplementedError


def format_table(results: Sequence[Result]) -> str:
    """A header JOB STATUS TRIES SECONDS DETAIL, then one aligned row per job."""
    raise NotImplementedError

jobrunner/cli.py

"""The command line: jobrunner [options] JOBFILE."""

import argparse
import asyncio
import cProfile
import logging
import sys
from pathlib import Path

from . import __version__
from .commands import ConfigError, load_jobs
from .report import format_table, summarise
from .runner import Runner

logger = logging.getLogger("jobrunner")

EXIT_OK = 0
EXIT_FAILED = 1  # at least one job did not succeed
EXIT_USAGE = 2  # a bad job file; argparse also exits with 2 on bad arguments


def positive_int(text: str) -> int:
    """An argparse type: a whole number of at least 1."""
    # TODO: raise argparse.ArgumentTypeError for anything else.
    raise NotImplementedError


def build_parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(
        prog="jobrunner",
        description="Run the jobs of a TOML job file concurrently, with timeouts and retries.",
    )
    parser.add_argument("jobfile", type=Path, metavar="JOBFILE", help="a TOML file of [[job]] tables")
    # TODO: -c/--concurrency (default 4), --fail-fast, --profile FILE,
    # -v/--verbose (a count) and --version.
    return parser


def configure_logging(verbosity: int) -> None:
    """Warnings (retries, failures) and, with -v, progress go to standard error."""
    handler = logging.StreamHandler(sys.stderr)
    handler.setFormatter(logging.Formatter("%(levelname)s: %(message)s"))
    logger.handlers[:] = [handler]
    logger.setLevel(logging.WARNING if verbosity == 0 else logging.INFO if verbosity == 1 else logging.DEBUG)
    logger.propagate = False


def main(argv: list[str] | None = None) -> int:
    """Run the tool with argv (default: sys.argv[1:]); return the exit status."""
    # TODO: parse, configure logging, load the jobs (a ConfigError is logged
    # and returns EXIT_USAGE), asyncio.run(Runner(...).run(jobs)), inside
    # `with cProfile.Profile() as profiler:` when --profile is given, then
    # print the table and the summary and return EXIT_OK or EXIT_FAILED.
    raise NotImplementedError


if __name__ == "__main__":
    sys.exit(main())

samples/jobs.toml

# A job file for jobrunner. Try: python -m jobrunner samples/jobs.toml -v
# A command whose first word is "python" runs with the Python running jobrunner.

[defaults]
timeout = 5.0
backoff = 0.2

[[job]]
name = "compile"
command = ["python", "-c", "print('compiled 12 modules')"]

[[job]]
name = "unit-tests"
command = ["python", "-c", "import time; time.sleep(0.05); print('84 passed')"]

[[job]]
name = "lint"
command = ["python", "-c", "import sys; sys.exit('style: line 12 is too long')"]
attempts = 2

[[job]]
name = "slow-download"
command = ["python", "-c", "import time; time.sleep(30)"]
timeout = 0.3

Abnahmetests

Das Projekt ist fertig, wenn jede Prüfung in test_main.py besteht. Lesen Sie sie vor dem Start: Sie sind die Spezifikation, als Code geschrieben.

test_main.py

import asyncio
import contextlib
import io
import json
import pstats
import sys
import tempfile
import time
from pathlib import Path

from jobrunner.cli import main
from jobrunner.commands import CommandFailed, ConfigError, command_action, load_jobs, run_command
from jobrunner.model import Job, Result, RetryPolicy, Status
from jobrunner.report import format_table, summarise
from jobrunner.runner import Runner


def job(name, value=None, *, delay=0.0, error=None, **options):
    """Ein Job, der delay Sekunden wartet und dann error auslöst oder value zurückgibt (sonst seinen Namen)."""
    async def action():
        await asyncio.sleep(delay)
        if error is not None:
            raise error
        return name if value is None else value
    return Job(name, action, **options)


def blocked(name, log, **options):
    """Ein Job, der ewig wartet und in log vermerkt, wenn er abgebrochen wird."""
    async def action():
        try:
            await asyncio.Event().wait()
        finally:
            log.append(name)
    return Job(name, action, **options)


def run(jobs, **options):
    return asyncio.run(Runner(**options).run(jobs))


def cli(argv):
    """Ruft main(argv) auf und gibt (Status, stdout, stderr) zurück."""
    out, err = io.StringIO(), io.StringIO()
    with contextlib.redirect_stdout(out), contextlib.redirect_stderr(err):
        status = main(argv)
    return status, out.getvalue(), err.getvalue()


def py(code):
    """Eine Befehlszeile, die diesen Python-Code ausführt, als TOML-Array."""
    return json.dumps(["python", "-c", code])


def job_file(folder, text):
    path = Path(folder) / "jobs.toml"
    path.write_text(text, encoding="utf-8")
    return path


def test_retry_policy():
    """RetryPolicy verdoppelt die Wartezeit nach jedem Versuch bis max_backoff und lehnt falsche Einstellungen ab"""
    policy = RetryPolicy(attempts=5, backoff=0.1, factor=2, max_backoff=0.5)
    got = [round(policy.delay(n), 6) for n in range(1, 5)]
    assert got == [0.1, 0.2, 0.4, 0.5], f"die Wartezeiten nach den Versuchen 1 bis 4 sind {got!r}, erwartet [0.1, 0.2, 0.4, 0.5]"
    assert RetryPolicy(retry_on=(ConnectionError,)).should_retry(ConnectionResetError()), "ConnectionResetError ist ein ConnectionError und sollte wiederholt werden"
    assert not RetryPolicy(retry_on=(ConnectionError,)).should_retry(ValueError()), "ValueError ist kein ConnectionError und sollte nicht wiederholt werden"
    for bad in [{"attempts": 0}, {"backoff": -1}, {"factor": 0.5}]:
        try:
            RetryPolicy(**bad)
        except ValueError:
            continue
        raise AssertionError(f"RetryPolicy(**{bad!r}) sollte ValueError auslösen")


def test_generic_types():
    """Job, Result und Runner.run sind generisch, geschrieben mit der Typparameter-Syntax von Python 3.12"""
    for thing in [Job, Result, Runner.run]:
        params = getattr(thing, "__type_params__", ())
        assert len(params) == 1, f"{thing.__qualname__} hat die Typparameter {params!r}, erwartet einer, etwa [T]"
    result = Result("x", Status.SUCCEEDED, 42, attempts=1)
    assert result.ok and not Result("y", Status.FAILED).ok, "Result.ok sollte nur für einen erfolgreichen Job True sein"


def test_runs_all_in_order():
    """run() führt jeden Job aus und gibt die Ergebnisse in der Reihenfolge der Jobs zurück"""
    results = run([job("slow", delay=0.03), job("fast", delay=0.0), job("middle", delay=0.01)])
    got = [(r.name, r.status, r.value, r.attempts) for r in results]
    want = [("slow", Status.SUCCEEDED, "slow", 1), ("fast", Status.SUCCEEDED, "fast", 1), ("middle", Status.SUCCEEDED, "middle", 1)]
    assert got == want, f"run ergab {got!r}, erwartet {want!r}"


def test_concurrency_limit():
    """Nie laufen mehr als concurrency Jobs gleichzeitig, und sie überlappen sich"""
    running = 0
    most = 0

    async def action():
        nonlocal running, most
        running += 1
        most = max(most, running)
        await asyncio.sleep(0.01)
        running -= 1

    started = time.perf_counter()
    run([Job(f"job{n}", action) for n in range(6)], concurrency=2)
    took = time.perf_counter() - started
    assert most == 2, f"höchstens {most} Jobs liefen gleichzeitig, erwartet genau 2"
    assert took < 0.5, f"sechs Jobs zu 0,01 s, je zwei gleichzeitig, dauerten {took:.2f} s"


def test_rejects_bad_input():
    """Doppelte Jobnamen und eine concurrency unter 1 lösen ValueError aus"""
    try:
        run([job("same"), job("same")])
    except ValueError:
        pass
    else:
        raise AssertionError("zwei Jobs namens 'same' sollten in run() ValueError auslösen")
    try:
        Runner(0)
    except ValueError:
        pass
    else:
        raise AssertionError("Runner(0) sollte ValueError auslösen")


def test_timeout():
    """Ein Job, der länger als sein Zeitlimit braucht, wird gestoppt und als abgelaufen gemeldet"""
    log = []
    results = run([blocked("stuck", log, timeout=0.02), job("quick")])
    stuck, quick = results
    assert stuck.status is Status.TIMED_OUT, f"der hängende Job ist {stuck.status!r}, erwartet Status.TIMED_OUT"
    assert isinstance(stuck.error, TimeoutError), f"sein Fehler ist {stuck.error!r}, erwartet ein TimeoutError"
    assert log == ["stuck"], "der hängende Job hätte nach Ablauf seiner Zeit abgebrochen werden sollen"
    assert quick.status is Status.SUCCEEDED, f"der andere Job ist {quick.status!r}, erwartet Status.SUCCEEDED"


def test_retries_with_backoff():
    """Ein fehlschlagender Job wird mit wachsenden Pausen wiederholt; andere Fehler oder keine Versuche mehr lassen ihn scheitern"""
    waits = []

    async def fake_sleep(seconds):
        waits.append(seconds)

    calls = 0

    async def flaky():
        nonlocal calls
        calls += 1
        if calls < 3:
            raise ConnectionError(f"try {calls} failed")
        return "done"

    policy = RetryPolicy(attempts=4, backoff=0.1, retry_on=(ConnectionError,))
    [result] = asyncio.run(Runner(sleep=fake_sleep).run([Job("flaky", flaky, retry=policy)]))
    assert (result.status, result.value, result.attempts) == (Status.SUCCEEDED, "done", 3), f"flaky endete {result.status!r} mit {result.value!r} nach {result.attempts} Versuchen, erwartet succeeded, 'done', 3"
    assert [round(w, 6) for w in waits] == [0.1, 0.2], f"der Runner wartete {waits!r} zwischen den Versuchen, erwartet [0.1, 0.2]"
    waits.clear()
    results = asyncio.run(Runner(sleep=fake_sleep).run([
        job("wrong", error=ValueError("bad data"), retry=policy),
        job("down", error=ConnectionError("no route"), retry=RetryPolicy(attempts=3, backoff=0)),
    ]))
    got = [(r.status, r.attempts) for r in results]
    assert got == [(Status.FAILED, 1), (Status.FAILED, 3)], f"erhalten {got!r}; ein ValueError sollte sofort scheitern, ein ConnectionError nach allen 3 Versuchen"
    assert isinstance(results[1].error, ConnectionError), f"gespeichert ist der Fehler {results[1].error!r}, erwartet der letzte ConnectionError"


def test_backoff_frees_the_slot():
    """Während ein Job auf die Wiederholung wartet, belegt er keinen Platz, sodass ein anderer Job laufen kann"""
    other_ran = None

    async def sleep_until_other_ran(seconds):
        await other_ran.wait()

    calls = 0

    async def fails_once():
        nonlocal calls
        calls += 1
        if calls == 1:
            raise ConnectionError("first try fails")
        return "second try"

    async def other():
        other_ran.set()
        return "other"

    async def scenario():
        nonlocal other_ran
        other_ran = asyncio.Event()
        runner = Runner(1, sleep=sleep_until_other_ran)
        jobs = [Job("retrying", fails_once, retry=RetryPolicy(attempts=2)), Job("other", other)]
        async with asyncio.timeout(1):
            return await runner.run(jobs)

    try:
        results = asyncio.run(scenario())
    except TimeoutError:
        raise AssertionError("mit concurrency 1 blieb der Lauf hängen: der wartende Job behielt seinen Platz") from None
    got = [(r.name, r.status) for r in results]
    assert got == [("retrying", Status.SUCCEEDED), ("other", Status.SUCCEEDED)], f"erhalten {got!r}"


def test_fail_fast():
    """Mit fail_fast bricht der erste Fehlschlag die anderen Jobs ab, die als abgebrochen gemeldet werden"""
    log = []
    results = run([blocked("waiting", log), job("broken", error=RuntimeError("boom")), blocked("queued", log)], concurrency=2, fail_fast=True)
    got = [(r.name, r.status) for r in results]
    want = [("waiting", Status.CANCELLED), ("broken", Status.FAILED), ("queued", Status.CANCELLED)]
    assert got == want, f"fail_fast ergab {got!r}, erwartet {want!r}"
    assert log == ["waiting"], f"während des Laufs abgebrochen wurden {log!r}, erwartet ['waiting'] (queued startete nie)"
    assert results[2].attempts == 0, f"der wartende Job startete nie, meldet aber {results[2].attempts} Versuche"


def test_external_cancellation():
    """Ein Abbruch von run() von außen bricht jeden Job ab und wird nicht verschluckt"""
    log = []

    async def scenario():
        async with asyncio.timeout(0.02):
            await Runner(3).run([blocked(name, log) for name in ["a", "b", "c"]])

    try:
        asyncio.run(scenario())
    except TimeoutError:
        pass
    else:
        raise AssertionError("das äußere Zeitlimit sollte mit TimeoutError enden; hat run() den Abbruch verschluckt?")
    assert sorted(log) == ["a", "b", "c"], f"abgebrochen wurden {sorted(log)!r}, erwartet ['a', 'b', 'c']"


def test_run_command():
    """run_command führt einen Befehl ohne Shell aus und löst CommandFailed aus, wenn er mit einem Fehler endet"""
    output = asyncio.run(run_command(["python", "-c", "print('hello from', 'python')"]))
    assert (output.returncode, output.stdout.strip()) == (0, "hello from python"), f"erhalten Exit-Status {output.returncode} und Ausgabe {output.stdout!r}"
    assert output.argv[0] == sys.executable, f"'python' sollte durch sys.executable ersetzt werden, aber argv ist {output.argv!r}"
    try:
        asyncio.run(run_command(["python", "-c", "import sys; sys.exit('disk full')"]))
    except CommandFailed as err:
        assert err.output.returncode == 1 and "disk full" in str(err), f"CommandFailed meldet {str(err)!r} mit Status {err.output.returncode}"
    else:
        raise AssertionError("ein Befehl, der mit Status 1 endet, sollte CommandFailed auslösen")


def test_command_timeout():
    """Ein Befehl, der sein Zeitlimit überschreitet, wird beendet und als abgelaufen gemeldet"""
    sleeper = Job("sleeper", command_action(["python", "-c", "import time; time.sleep(30)"]), timeout=0.2)
    started = time.perf_counter()
    [result] = run([sleeper])
    took = time.perf_counter() - started
    assert result.status is Status.TIMED_OUT, f"der schlafende Befehl ist {result.status!r}, erwartet Status.TIMED_OUT"
    assert took < 5, f"der Lauf dauerte {took:.1f} s; der Befehl hätte nach 0,2 s beendet werden sollen"


def test_load_jobs():
    """load_jobs liest [[job]]-Tabellen und ergänzt die Werte aus [defaults]"""
    text = f"""
[defaults]
timeout = 10
attempts = 2

[[job]]
name = "build"
command = {py("print(1)")}

[[job]]
name = "deploy"
command = ["python", "deploy.py"]
timeout = 2.5
attempts = 1
backoff = 0
"""
    with tempfile.TemporaryDirectory() as folder:
        jobs = load_jobs(job_file(folder, text))
    got = [(j.name, j.timeout, j.retry.attempts) for j in jobs]
    assert got == [("build", 10.0, 2), ("deploy", 2.5, 1)], f"load_jobs ergab (Name, Zeitlimit, Versuche) {got!r}"
    assert jobs[1].retry.backoff == 0, f"deploy hat backoff {jobs[1].retry.backoff}, erwartet 0"


def test_load_jobs_rejects():
    """Eine fehlende Datei, falsches TOML oder ein falscher Job lösen ConfigError mit einer klaren Meldung aus"""
    bad = {
        "not toml": "[[job]\nname = 1",
        "command": '[[job]]\nname = "a"\ncommand = "echo hi"\n',
        "twice": '[[job]]\nname = "a"\ncommand = ["x"]\n[[job]]\nname = "a"\ncommand = ["y"]\n',
        "retries": '[[job]]\nname = "a"\ncommand = ["x"]\nretries = 3\n',
        "timeout": '[[job]]\nname = "a"\ncommand = ["x"]\ntimeout = -1\n',
    }
    with tempfile.TemporaryDirectory() as folder:
        try:
            load_jobs(Path(folder) / "missing.toml")
        except ConfigError as err:
            assert "missing.toml" in str(err), f"die Meldung {str(err)!r} sollte die Datei nennen"
        else:
            raise AssertionError("eine fehlende Jobdatei sollte ConfigError auslösen")
        for word, text in bad.items():
            try:
                load_jobs(job_file(folder, text))
            except ConfigError as err:
                message = str(err)
                assert ("TOML" if word == "not toml" else word) in message, f"für {text!r} sollte die Meldung {message!r} {word!r} erwähnen"
            else:
                raise AssertionError(f"load_jobs akzeptierte {text!r}; erwartet ConfigError")


def test_summary_and_table():
    """summarise zählt die Ergebnisse je Status, und format_table gibt eine Zeile je Job aus"""
    results = [
        Result("a", Status.SUCCEEDED, 1, attempts=1, elapsed=0.5),
        Result("b", Status.FAILED, error=RuntimeError("exit status 1"), attempts=3, elapsed=2.0),
        Result("c", Status.TIMED_OUT, error=TimeoutError(), attempts=1, elapsed=1.0),
        Result("d", Status.SUCCEEDED, 4, attempts=2, elapsed=0.25),
    ]
    summary = summarise(results)
    assert str(summary) == "4 jobs: 2 succeeded, 1 failed, 1 timed out", f"die Zusammenfassung lautet {str(summary)!r}"
    assert (summary.attempts, summary.slowest, summary.ok) == (7, "b", False), f"attempts, slowest und ok sind {(summary.attempts, summary.slowest, summary.ok)!r}, erwartet (7, 'b', False)"
    lines = format_table(results).splitlines()
    assert len(lines) == 5 and lines[0].split()[:2] == ["JOB", "STATUS"], f"die Tabelle ist {lines!r}; erwartet eine Kopfzeile und 4 Zeilen"
    got = [line.split()[:3] for line in lines[1:]]
    want = [["a", "succeeded", "1"], ["b", "failed", "3"], ["c", "timed", "out"], ["d", "succeeded", "2"]]
    assert got == want, f"die Zeilen beginnen mit {got!r}, erwartet {want!r}"


def test_cli_runs_file():
    """jobrunner JOBDATEI gibt eine Tabelle und eine Zusammenfassung aus, protokolliert Wiederholungen und gibt 1 zurück, wenn ein Job scheiterte"""
    text = f"""
[defaults]
backoff = 0

[[job]]
name = "hello"
command = {py("print('hello')")}

[[job]]
name = "broken"
command = {py("import sys; sys.exit(3)")}
attempts = 2

[[job]]
name = "count"
command = {py("print(sum(range(10)))")}
"""
    with tempfile.TemporaryDirectory() as folder:
        status, out, err = cli([str(job_file(folder, text)), "-c", "2"])
    lines = out.splitlines()
    assert status == 1, f"main gab {status} zurück, erwartet 1, weil ein Job scheiterte"
    assert lines[-1] == "3 jobs: 2 succeeded, 1 failed", f"die letzte Zeile ist {lines[-1]!r}"
    rows = [line.split()[:3] for line in lines[1:-1]]
    assert rows == [["hello", "succeeded", "1"], ["broken", "failed", "2"], ["count", "succeeded", "1"]], f"die Zeilen beginnen mit {rows!r}"
    assert "retrying" in err and "broken" in err, f"die Standardfehlerausgabe war {err!r}; erwartet eine Warnung, dass broken wiederholt wird"


def test_cli_fail_fast():
    """Wenn alle Jobs gelingen, ist das Ergebnis 0; --fail-fast bricht nach einem Fehlschlag den Rest ab"""
    good = f'[[job]]\nname = "a"\ncommand = {py("pass")}\n[[job]]\nname = "b"\ncommand = {py("pass")}\n'
    mixed = f'[[job]]\nname = "sleeper"\ncommand = {py("import time; time.sleep(30)")}\n[[job]]\nname = "broken"\ncommand = {py("raise SystemExit(1)")}\n'
    with tempfile.TemporaryDirectory() as folder:
        status, out, _ = cli([str(job_file(folder, good))])
        assert (status, out.splitlines()[-1]) == (0, "2 jobs: 2 succeeded"), f"main gab {status} zurück und endete mit {out.splitlines()[-1:]!r}"
        started = time.perf_counter()
        status, out, _ = cli([str(job_file(folder, mixed)), "--fail-fast"])
        took = time.perf_counter() - started
    assert (status, out.splitlines()[-1]) == (1, "2 jobs: 1 failed, 1 cancelled"), f"mit --fail-fast gab main {status} zurück und endete mit {out.splitlines()[-1:]!r}"
    assert took < 5, f"--fail-fast dauerte {took:.1f} s; der schlafende Befehl hätte abgebrochen werden sollen"


def test_cli_bad_input():
    """Eine fehlerhafte Jobdatei ergibt 2 mit einer Meldung auf stderr; falsche Argumente enden mit Status 2"""
    with tempfile.TemporaryDirectory() as folder:
        status, out, err = cli([str(Path(folder) / "nope.toml")])
    assert (status, out) == (2, ""), f"für eine fehlende Jobdatei gab main {status} zurück und gab {out!r} aus, erwartet 2 und nichts"
    assert "nope.toml" in err, f"die Meldung {err!r} sollte nope.toml nennen"
    for argv in [[], ["jobs.toml", "-c", "0"], ["jobs.toml", "--concurrency", "many"]]:
        try:
            with contextlib.redirect_stderr(io.StringIO()):
                main(argv)
        except SystemExit as exc:
            assert exc.code == 2, f"für {argv!r} war der Exit-Status {exc.code!r}, erwartet 2"
        else:
            raise AssertionError(f"main({argv!r}) kehrte zurück, statt mit Status 2 zu enden")


def test_cli_profile():
    """--profile DATEI speichert cProfile-Statistiken, die pstats lesen kann"""
    text = f'[[job]]\nname = "a"\ncommand = {py("pass")}\n'
    with tempfile.TemporaryDirectory() as folder:
        target = Path(folder) / "run.prof"
        status, _, _ = cli([str(job_file(folder, text)), "--profile", str(target)])
        assert status == 0 and target.exists(), f"main gab {status} zurück; die Profildatei existiert: {target.exists()}"
        stats = pstats.Stats(str(target), stream=io.StringIO())
    functions = {name for _, _, name in stats.stats}
    assert "run" in functions and "_run_job" in functions, "das Profil sollte Runner.run und Runner._run_job enthalten"


def test_sample_jobs():
    """samples/jobs.toml endet mit 2 erfolgreichen, 1 gescheiterten und 1 abgelaufenen Job"""
    status, out, _ = cli(["samples/jobs.toml"])
    last = out.splitlines()[-1] if out else ""
    assert (status, last) == (1, "4 jobs: 2 succeeded, 1 failed, 1 timed out"), f"für samples/jobs.toml gab main {status} zurück und endete mit {last!r}"

Das fertige Programm starten

python -m venv .venv, activate it, then python -m pip install -e . and jobrunner samples/jobs.toml -v (or python -m jobrunner samples/jobs.toml); add --profile run.prof and read the file with pstats
Referenzlösung

Versuchen Sie zuerst die Etappen. Diese Lösung besteht alle Abnahmetests und die Typprüfung.

pyproject.toml

[build-system]
requires = ["setuptools >= 77.0.3"]
build-backend = "setuptools.build_meta"

[project]
name = "jobrunner"
version = "1.0.0"
description = "Run jobs concurrently with asyncio: timeouts, retries with backoff, and cancellation."
requires-python = ">= 3.14"
dependencies = []

[project.scripts]
jobrunner = "jobrunner.cli:main"

[tool.setuptools]
packages = ["jobrunner"]

jobrunner/__init__.py

"""Run jobs concurrently with asyncio: timeouts, retries with backoff, and cancellation."""

__version__ = "1.0.0"

jobrunner/__main__.py

"""python -m jobrunner runs the command line."""

import sys

from .cli import main

sys.exit(main())

jobrunner/model.py

"""The domain: jobs, retry policies and results, generic in what a job returns."""

from collections.abc import Awaitable, Callable
from dataclasses import dataclass, field
from enum import StrEnum


class Status(StrEnum):
    """How a job ended."""

    SUCCEEDED = "succeeded"
    FAILED = "failed"
    TIMED_OUT = "timed out"
    CANCELLED = "cancelled"


@dataclass(frozen=True, slots=True)
class RetryPolicy:
    """How often to try a job, and how long to wait between tries.

    The wait before try n + 1 is backoff * factor ** (n - 1), capped at
    max_backoff: 0.1, 0.2, 0.4, ... with the defaults.
    """

    attempts: int = 1
    backoff: float = 0.1
    factor: float = 2.0
    max_backoff: float = 30.0
    retry_on: tuple[type[Exception], ...] = (Exception,)

    def __post_init__(self) -> None:
        if self.attempts < 1:
            raise ValueError(f"attempts must be at least 1, not {self.attempts}")
        if self.backoff < 0 or self.max_backoff < 0:
            raise ValueError("backoff and max_backoff cannot be negative")
        if self.factor < 1:
            raise ValueError(f"factor must be at least 1, not {self.factor}")

    def delay(self, attempt: int) -> float:
        """Seconds to wait after failed try number attempt (counting from 1)."""
        return min(self.backoff * self.factor ** (attempt - 1), self.max_backoff)

    def should_retry(self, error: BaseException) -> bool:
        """Whether this kind of error is worth another try."""
        return isinstance(error, self.retry_on)


@dataclass(frozen=True, slots=True)
class Job[T]:
    """A named unit of work. action is called afresh for every try."""

    name: str
    action: Callable[[], Awaitable[T]]
    timeout: float | None = None
    retry: RetryPolicy = field(default_factory=RetryPolicy)


@dataclass(frozen=True, slots=True)
class Result[T]:
    """What happened to one job."""

    name: str
    status: Status
    value: T | None = None
    error: BaseException | None = None
    attempts: int = 0
    elapsed: float = 0.0

    @property
    def ok(self) -> bool:
        return self.status is Status.SUCCEEDED

jobrunner/runner.py

"""Running jobs concurrently: a TaskGroup, a semaphore, timeouts, retries and cancellation."""

import asyncio
import logging
from collections.abc import Awaitable, Callable, Iterable

from .model import Job, Result, Status

logger = logging.getLogger(__name__)

type Sleep = Callable[[float], Awaitable[None]]


class _FailFast(Exception):
    """Raised inside the TaskGroup to cancel the other jobs after a failure."""


class Runner:
    """Runs jobs at most `concurrency` at a time and reports a Result for each."""

    def __init__(self, concurrency: int = 4, *, fail_fast: bool = False, sleep: Sleep = asyncio.sleep) -> None:
        if concurrency < 1:
            raise ValueError(f"concurrency must be at least 1, not {concurrency}")
        self.concurrency = concurrency
        self.fail_fast = fail_fast
        # Injected so tests can record the backoff instead of waiting for it.
        self._sleep = sleep

    async def run[T](self, jobs: Iterable[Job[T]]) -> list[Result[T]]:
        """Run every job; return their results in the order the jobs were given.

        With fail_fast, the first job that fails or times out cancels the
        others, which are reported as cancelled. If the caller cancels run()
        itself, every job is cancelled and the cancellation propagates.
        """
        jobs = list(jobs)
        names = [job.name for job in jobs]
        if len(set(names)) != len(names):
            raise ValueError("job names must be unique")
        slots = asyncio.Semaphore(self.concurrency)
        stopping = asyncio.Event()
        results: dict[str, Result[T]] = {}
        try:
            async with asyncio.TaskGroup() as group:
                for job in jobs:
                    group.create_task(self._run_job(job, slots, stopping, results), name=job.name)
        except* _FailFast:
            logger.warning("stopped after the first failure (fail fast)")
        return [results[name] for name in names]

    async def _run_job[T](
        self, job: Job[T], slots: asyncio.Semaphore, stopping: asyncio.Event, results: dict[str, Result[T]]
    ) -> None:
        loop = asyncio.get_running_loop()
        started = loop.time()
        attempt = 0

        def finish(status: Status, value: T | None = None, error: BaseException | None = None) -> None:
            elapsed = loop.time() - started
            results[job.name] = Result(job.name, status, value, error, attempt, elapsed)
            log = logger.info if status is Status.SUCCEEDED else logger.warning
            log("job %s %s after %d attempt(s) in %.2f s", job.name, status, attempt, elapsed)

        try:
            while True:
                async with slots:
                    if stopping.is_set():
                        # A job that got its slot just as fail fast began never starts.
                        finish(Status.CANCELLED)
                        return
                    attempt += 1
                    logger.debug("job %s: attempt %d", job.name, attempt)
                    status: Status
                    error: Exception
                    try:
                        async with asyncio.timeout(job.timeout):
                            value = await job.action()
                    except TimeoutError as err:
                        status, error = Status.TIMED_OUT, err
                    except Exception as err:
                        status, error = Status.FAILED, err
                    else:
                        finish(Status.SUCCEEDED, value)
                        return
                    if attempt >= job.retry.attempts or not job.retry.should_retry(error):
                        finish(status, error=error)
                        if self.fail_fast:
                            # Still holding the slot, so no queued job starts after this.
                            stopping.set()
                            raise _FailFast(job.name)
                        return
                delay = job.retry.delay(attempt)
                logger.warning(
                    "job %s: attempt %d of %d %s (%s); retrying in %.2f s",
                    job.name, attempt, job.retry.attempts, status, describe(error), delay,
                )
                # Outside `async with slots`: a job waiting to retry frees its slot.
                await self._sleep(delay)
        except asyncio.CancelledError:
            finish(Status.CANCELLED)
            raise


def describe(error: BaseException | None) -> str:
    """A short description of an error, for logs and reports."""
    if error is None:
        return ""
    if isinstance(error, TimeoutError) and not str(error):
        return "timed out"
    return f"{type(error).__name__}: {error}" if str(error) else type(error).__name__

jobrunner/commands.py

"""Jobs that run commands, and loading them from a TOML job file."""

import asyncio
import sys
import tomllib
from collections.abc import Awaitable, Callable, Sequence
from dataclasses import dataclass
from pathlib import Path
from typing import Any

from .model import Job, RetryPolicy


class ConfigError(Exception):
    """The job file is missing, is not valid TOML, or describes a job wrongly."""


@dataclass(frozen=True, slots=True)
class CommandOutput:
    """What a finished command printed, and its exit status."""

    argv: tuple[str, ...]
    returncode: int
    stdout: str
    stderr: str


class CommandFailed(Exception):
    """A command exited with a status other than 0."""

    def __init__(self, output: CommandOutput) -> None:
        last = output.stderr.strip().splitlines()[-1:] or [""]
        super().__init__(f"exit status {output.returncode}" + (f": {last[0]}" if last[0] else ""))
        self.output = output


async def run_command(argv: Sequence[str]) -> CommandOutput:
    """Run a command without a shell and wait for it; raise CommandFailed if it fails.

    A command whose first word is "python" runs with the interpreter running
    this program. If the waiting task is cancelled (by a timeout, say), the
    process is killed before the cancellation goes on.
    """
    args = [sys.executable, *argv[1:]] if argv[0] == "python" else list(argv)
    process = await asyncio.create_subprocess_exec(
        *args, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE
    )
    try:
        stdout, stderr = await process.communicate()
    except asyncio.CancelledError:
        process.kill()
        await process.wait()
        raise
    returncode = await process.wait()
    output = CommandOutput(tuple(args), returncode, stdout.decode(errors="replace"), stderr.decode(errors="replace"))
    if returncode != 0:
        raise CommandFailed(output)
    return output


def command_action(argv: Sequence[str]) -> Callable[[], Awaitable[CommandOutput]]:
    """A job action that runs argv afresh each time it is called."""
    frozen = tuple(argv)

    async def action() -> CommandOutput:
        return await run_command(frozen)

    return action


JOB_KEYS = {"name", "command", "timeout", "attempts", "backoff"}


def _number(table: dict[str, Any], key: str, where: str, *, minimum: float) -> float | None:
    value = table.get(key)
    if value is None:
        return None
    if isinstance(value, bool) or not isinstance(value, int | float) or value < minimum:
        raise ConfigError(f"{where}: {key} must be a number of at least {minimum}, not {value!r}")
    return float(value)


def _job(table: object, defaults: dict[str, Any], where: str) -> Job[CommandOutput]:
    if not isinstance(table, dict):
        raise ConfigError(f"{where}: a job must be a table")
    unknown = set(table) - JOB_KEYS
    if unknown:
        raise ConfigError(f"{where}: unknown key(s) {', '.join(sorted(unknown))}")
    name = table.get("name")
    if not isinstance(name, str) or not name:
        raise ConfigError(f"{where}: every job needs a name")
    command = table.get("command")
    if not isinstance(command, list) or not command or not all(isinstance(word, str) for word in command):
        raise ConfigError(f"{where} ({name}): command must be a non-empty list of strings")
    merged = defaults | table
    timeout = _number(merged, "timeout", where, minimum=0.001)
    attempts = _number(merged, "attempts", where, minimum=1)
    backoff = _number(merged, "backoff", where, minimum=0)
    if attempts is not None and not attempts.is_integer():
        raise ConfigError(f"{where}: attempts must be a whole number")
    retry = RetryPolicy(
        attempts=int(attempts) if attempts is not None else 1,
        backoff=backoff if backoff is not None else 0.5,
    )
    return Job(name, command_action(command), timeout, retry)


def load_jobs(path: Path) -> list[Job[CommandOutput]]:
    """Read the jobs of a TOML job file: an optional [defaults] table and [[job]] tables."""
    try:
        with path.open("rb") as f:
            data = tomllib.load(f)
    except OSError as err:
        raise ConfigError(f"{path}: cannot read it ({err.strerror})") from err
    except tomllib.TOMLDecodeError as err:
        raise ConfigError(f"{path}: not valid TOML ({err})") from err
    defaults = data.get("defaults", {})
    if not isinstance(defaults, dict) or set(defaults) - (JOB_KEYS - {"name", "command"}):
        raise ConfigError(f"{path}: [defaults] may only set timeout, attempts and backoff")
    tables = data.get("job", [])
    if not isinstance(tables, list) or not tables:
        raise ConfigError(f"{path}: no [[job]] tables")
    jobs = [_job(table, defaults, f"{path}, job {n}") for n, table in enumerate(tables, start=1)]
    seen: set[str] = set()
    for job in jobs:
        if job.name in seen:
            raise ConfigError(f"{path}: the job name {job.name!r} is used twice")
        seen.add(job.name)
    return jobs

jobrunner/report.py

"""Summing up a run: counts per status and a table of results."""

from collections import Counter
from collections.abc import Sequence
from dataclasses import dataclass

from .model import Result, Status
from .runner import describe


@dataclass(frozen=True, slots=True)
class Summary:
    """The shape of a finished run."""

    total: int
    counts: dict[Status, int]
    attempts: int
    slowest: str | None

    @property
    def ok(self) -> bool:
        return self.counts.get(Status.SUCCEEDED, 0) == self.total

    def __str__(self) -> str:
        """For example "4 jobs: 2 succeeded, 1 failed, 1 timed out"."""
        parts = [f"{self.counts[status]} {status}" for status in Status if self.counts.get(status)]
        noun = "job" if self.total == 1 else "jobs"
        return f"{self.total} {noun}: " + (", ".join(parts) if parts else "nothing to do")


def summarise[T](results: Sequence[Result[T]]) -> Summary:
    """Count the results per status, add up the attempts and find the slowest job."""
    counts = Counter(result.status for result in results)
    slowest = max(results, key=lambda result: result.elapsed, default=None)
    return Summary(
        total=len(results),
        counts=dict(counts),
        attempts=sum(result.attempts for result in results),
        slowest=slowest.name if slowest else None,
    )


def format_table[T](results: Sequence[Result[T]]) -> str:
    """One aligned row per job: name, status, attempts, seconds and the error, if any."""
    header = ("JOB", "STATUS", "TRIES", "SECONDS", "DETAIL")
    rows = [
        (r.name, str(r.status), str(r.attempts), f"{r.elapsed:.2f}", describe(r.error) if not r.ok else "")
        for r in results
    ]
    width = [max(len(row[i]) for row in [header, *rows]) for i in range(4)]
    lines = []
    for name, status, tries, seconds, detail in [header, *rows]:
        line = f"{name:<{width[0]}}  {status:<{width[1]}}  {tries:>{width[2]}}  {seconds:>{width[3]}}  {detail}"
        lines.append(line.rstrip())
    return "\n".join(lines)

jobrunner/cli.py

"""The command line: jobrunner [options] JOBFILE."""

import argparse
import asyncio
import cProfile
import logging
import sys
from pathlib import Path

from . import __version__
from .commands import ConfigError, load_jobs
from .report import format_table, summarise
from .runner import Runner

logger = logging.getLogger("jobrunner")

EXIT_OK = 0
EXIT_FAILED = 1  # at least one job did not succeed
EXIT_USAGE = 2  # a bad job file; argparse also exits with 2 on bad arguments


def positive_int(text: str) -> int:
    """An argparse type: a whole number of at least 1."""
    try:
        value = int(text)
    except ValueError:
        raise argparse.ArgumentTypeError(f"not a whole number: {text!r}") from None
    if value < 1:
        raise argparse.ArgumentTypeError(f"must be at least 1, not {value}")
    return value


def build_parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(
        prog="jobrunner",
        description="Run the jobs of a TOML job file concurrently, with timeouts and retries.",
    )
    parser.add_argument("jobfile", type=Path, metavar="JOBFILE", help="a TOML file of [[job]] tables")
    parser.add_argument("-c", "--concurrency", type=positive_int, default=4, help="jobs to run at once (default: 4)")
    parser.add_argument("--fail-fast", action="store_true", help="cancel the other jobs when one fails")
    parser.add_argument("--profile", type=Path, metavar="FILE", help="profile the run with cProfile and save the stats to FILE")
    parser.add_argument("-v", "--verbose", action="count", default=0, help="log progress (-vv for every attempt)")
    parser.add_argument("--version", action="version", version=f"%(prog)s {__version__}")
    return parser


def configure_logging(verbosity: int) -> None:
    """Warnings (retries, failures) and, with -v, progress go to standard error."""
    handler = logging.StreamHandler(sys.stderr)
    handler.setFormatter(logging.Formatter("%(levelname)s: %(message)s"))
    logger.handlers[:] = [handler]
    logger.setLevel(logging.WARNING if verbosity == 0 else logging.INFO if verbosity == 1 else logging.DEBUG)
    logger.propagate = False


def main(argv: list[str] | None = None) -> int:
    """Run the tool with argv (default: sys.argv[1:]); return the exit status."""
    args = build_parser().parse_args(argv)
    configure_logging(args.verbose)
    try:
        jobs = load_jobs(args.jobfile)
    except ConfigError as err:
        logger.error("%s", err)
        return EXIT_USAGE
    runner = Runner(args.concurrency, fail_fast=args.fail_fast)
    logger.info("running %d job(s), %d at a time", len(jobs), args.concurrency)
    if args.profile:
        with cProfile.Profile() as profiler:
            results = asyncio.run(runner.run(jobs))
        profiler.dump_stats(args.profile)
        logger.info("profile saved to %s", args.profile)
    else:
        results = asyncio.run(runner.run(jobs))
    summary = summarise(results)
    print(format_table(results))
    print(summary)
    return EXIT_OK if summary.ok else EXIT_FAILED


if __name__ == "__main__":
    sys.exit(main())

samples/jobs.toml

# A job file for jobrunner. Try: python -m jobrunner samples/jobs.toml -v
# A command whose first word is "python" runs with the Python running jobrunner.

[defaults]
timeout = 5.0
backoff = 0.2

[[job]]
name = "compile"
command = ["python", "-c", "print('compiled 12 modules')"]

[[job]]
name = "unit-tests"
command = ["python", "-c", "import time; time.sleep(0.05); print('84 passed')"]

[[job]]
name = "lint"
command = ["python", "-c", "import sys; sys.exit('style: line 12 is too long')"]
attempts = 2

[[job]]
name = "slow-download"
command = ["python", "-c", "import time; time.sleep(30)"]
timeout = 0.3

Weiterentwickeln

  • Ergänzen Sie Jitter: Multiplizieren Sie jedes Backoff mit einem Zufallsfaktor aus einem random.Random, das RetryPolicy erhält, damit viele gescheiterte Jobs nicht alle im selben Moment wiederholen.
  • Ergänzen Sie depends_on für Jobs und starten Sie einen Job erst, wenn die Jobs, von denen er abhängt, gelungen sind. Ordnen Sie die Jobs und finden Sie Zyklen mit graphlib.TopologicalSorter.
  • Schreiben Sie die Ausgabe jedes Befehls Zeile für Zeile mit process.stdout.readline in eine Logdatei je Job, statt sie mit communicate ganz im Speicher zu halten.
  • Ergänzen Sie eine Option --json, die die Ergebnisse als JSON schreibt, typisieren Sie jede Zeile als TypedDict, und halten Sie mypy sauber.
  • Profilieren Sie einen Lauf mit 200 kurzen Jobs, finden Sie mit pstats, nach kumulierter Zeit sortiert, wohin die Zeit geht, und messen Sie nach jeder Änderung erneut.

Projekte sind Übung: Ihre Prüfungen laufen im Browser oder auf Ihrem Rechner und zählen nie für eine Bescheinigung.