Skip to content

docker_engine

tit.jobs.docker_engine

Dependency-free Docker Engine API client — stdlib http.client over AF_UNIX.

The Node-side twin of this module is desktop/src/main/docker/engine.ts; both talk the same Docker Engine REST API subset directly over the Unix socket, never the docker CLI and never the third-party docker PyPI package (r5-docker-integration.md §5 — that package is not installed in idossha/ti-toolbox:v3.0.0 today, and adding it would only grow a 19 GB image further for something ~15 lines of stdlib already does).

Scope, and why it is narrower than the TypeScript client: this module is used from inside the tit container, where QSIPrep/QSIRecon are spawned as sibling containers via Docker-outside-of-Docker (DooD) — the root docker-compose.yml mounts the host's /var/run/docker.sock straight through, a Unix path, regardless of what OS the host actually is. So unlike the Electron client (which must speak to a real Windows named pipe on a Windows host), this module never needs npipe transport at all — Windows support here would mean supporting it as a pywin32 dependency for a case that cannot occur in this architecture, which is why it is explicitly out of scope, not merely unimplemented. If a future redesign ever has Electron spawn QSIPrep containers directly instead of via DooD (r5 §5's "bigger architectural change, out of scope"), that would need a real npipe-capable Python client; this one is not it.

Discovery: DOCKER_HOST (unix:// only — tcp:///npipe:// are parsed for completeness but this module never drives the tcp/npipe internals http.client would need), else the well-known mount path /var/run/docker.sock. No docker context inspect shell-out here (unlike discover.ts): the Python side always runs inside the tit container, where the socket location is a known, fixed bind mount, not an environment to be discovered the way a user's host machine is (r5 §5's "the Python client only ever needs the Unix-socket transport").

API version negotiation: identical strategy to engine.ts — :meth:DockerEngineClient.version hits the always-unversioned GET /version, reads ApiVersion from the response, and every later call is prefixed /v<ApiVersion>/....

Framing: :meth:DockerEngineClient.logs demultiplexes the same 8-byte [STREAM_TYPE, 0,0,0, SIZE(be32)][SIZE bytes] frames as frames.ts's LogFrameDecoder (mirrored here as :class:LogFrameDecoder, pure and independently unit-tested); pull progress and /events are newline-delimited JSON (:class:NdjsonDecoder, mirrors NdjsonDecoder).

Job-container convenience: :func:run_job_container always accepts a job_id and sets Labels["tit.job_id"] — the exact label tit/jobs/runner.py's stop_docker_siblings(job_id) filters on (docker ps -q --filter label=tit.job_id=<id>). tit/pre/qsi/docker_builder.py's _label_args() now sets this same label on both the qsiprep and qsirecon docker run invocations it builds (confirmed: rg -n "tit\.job_id" tit/pre/qsi/docker_builder.py — two hits), so stop_docker_siblings already finds its matching containers for QSI jobs today; the gap this paragraph used to describe (r5 §1.4 / skeptic-3 claim #3) is closed on the labeling side. What is not closed: docker_builder.py still shells out to the docker CLI directly rather than through this client (rg -n '"docker"' tit/pre/qsi/docker_builder.py) — migrating it onto :func:run_job_container is still open work, out of scope here. Mount/env shape (/data read-only BIDS input, /out output, /work scratch) mirrors tit/pre/qsi/docker_builder.py's DockerPaths so the two stay drop-in compatible; see run_qsiprep_example at the bottom of this file for the exact call shape a migrated run_qsiprep() would make (documentation only — this module is not wired into tit/pre/qsi/*.py yet).

DockerEngineError

DockerEngineError(kind: str, message: str, status_code: int | None = None)

Bases: Exception

Carries a stable kind (mirrors DockerErrorKind in engine.ts) and, for an HTTP- level failure, the status_code that produced it.

Source code in tit/jobs/docker_engine.py
def __init__(self, kind: str, message: str, status_code: int | None = None) -> None:
    super().__init__(message)
    self.kind = kind
    self.status_code = status_code

LogFrameDecoder

LogFrameDecoder()

Demultiplexes Docker's [STREAM_TYPE,0,0,0,SIZE(be32)][SIZE bytes] log/attach framing.

Source code in tit/jobs/docker_engine.py
def __init__(self) -> None:
    self._buf = b""

NdjsonDecoder

NdjsonDecoder()

Splits a byte stream into complete JSON objects, one per newline-delimited line.

Source code in tit/jobs/docker_engine.py
def __init__(self) -> None:
    self._carry = ""

DockerEngineClient

DockerEngineClient(conn: DockerConnection, timeout_s: float | None = DEFAULT_TIMEOUT_S)

Covers the endpoints a QSIPrep/QSIRecon DooD spawn needs: pull, create/start/wait/ stop/kill/remove/inspect, logs (demuxed), events. No exec/attach-hijack (unused, r5 §1.2).

Source code in tit/jobs/docker_engine.py
def __init__(
    self, conn: DockerConnection, timeout_s: float | None = DEFAULT_TIMEOUT_S
) -> None:
    self._conn = conn
    self._timeout_s = timeout_s
    self._api_version: str | None = None

ping

ping(timeout_s: float | None = SHORT_TIMEOUT_S) -> bool

Liveness probe. Bounded by timeout_s (not the client's default): a wedged daemon that accepts connections and never answers must report "not reachable" promptly rather than block the caller.

Source code in tit/jobs/docker_engine.py
def ping(self, timeout_s: float | None = SHORT_TIMEOUT_S) -> bool:
    """Liveness probe. Bounded by *timeout_s* (not the client's default): a wedged daemon
    that accepts connections and never answers must report "not reachable" promptly rather
    than block the caller."""
    try:
        status, resp, conn = self._request("GET", "/_ping", timeout_s=timeout_s)
        try:
            resp.read()
        finally:
            conn.close()
        return status == 200
    except (DockerEngineError, OSError):
        return False

system_df

system_df(timeout_s: float | None | object = 'default') -> dict[str, Any]

GET /system/df -- the Engine-API equivalent of docker system df.

Read-only, used by the System page's Docker-health panel. It is a genuinely expensive call on the daemon side (it walks the image graph), which is why the caller polls it on a multi-second TTL and passes an explicit timeout rather than the client default.

Source code in tit/jobs/docker_engine.py
def system_df(
    self, timeout_s: float | None | object = "default"
) -> dict[str, Any]:
    """``GET /system/df`` -- the Engine-API equivalent of ``docker system df``.

    Read-only, used by the System page's Docker-health panel. It is a genuinely expensive call
    on the daemon side (it walks the image graph), which is why the caller polls it on a
    multi-second TTL and passes an explicit timeout rather than the client default.
    """
    self._ensure_negotiated()
    return self._json_request(
        "GET", self._vpath("/system/df"), timeout_s=timeout_s
    )

list_images

list_images(timeout_s: float | None | object = 'default') -> list[dict[str, Any]]

GET /images/json -- docker images. Read-only; see :meth:system_df.

Source code in tit/jobs/docker_engine.py
def list_images(
    self, timeout_s: float | None | object = "default"
) -> list[dict[str, Any]]:
    """``GET /images/json`` -- ``docker images``. Read-only; see :meth:`system_df`."""
    self._ensure_negotiated()
    result = self._json_request(
        "GET", self._vpath("/images/json"), timeout_s=timeout_s
    )
    return list(result or [])

pull_image

pull_image(image: str, tag: str, on_progress: Callable[[dict[str, Any]], None] | None = None) -> None

POST /images/create — streams NDJSON progress; raises on a non-2xx status or a {"error": ...} object appearing mid-stream (Docker reports some pull failures inside an HTTP-200 NDJSON stream, not as an HTTP error status).

Source code in tit/jobs/docker_engine.py
def pull_image(
    self,
    image: str,
    tag: str,
    on_progress: Callable[[dict[str, Any]], None] | None = None,
) -> None:
    """``POST /images/create`` — streams NDJSON progress; raises on a non-2xx status *or* a
    ``{"error": ...}`` object appearing mid-stream (Docker reports some pull failures inside an
    HTTP-200 NDJSON stream, not as an HTTP error status)."""
    self._ensure_negotiated()
    status, resp, conn = self._request(
        "POST",
        self._vpath("/images/create"),
        query={"fromImage": image, "tag": tag},
        timeout_s=NO_TIMEOUT,
    )
    try:
        if status >= 400:
            raise _api_error_from_body(status, resp.read())
        decoder = NdjsonDecoder()
        while True:
            chunk = resp.read(65536)
            if not chunk:
                break
            for event in decoder.push(chunk):
                if event.get("error"):
                    raise DockerEngineError("unknown", str(event["error"]))
                if on_progress:
                    on_progress(event)
    finally:
        conn.close()

list_containers

list_containers(*, filters: dict[str, list[str]] | None = None, all_containers: bool = False, timeout_s: float | None | object = 'default') -> list[dict[str, Any]]

GET /containers/json — the Engine-API equivalent of docker ps [--filter …].

filters is the Engine's own filter map (e.g. {"label": ["tit.job_id=abc"]}); it is JSON-encoded into the query string by :meth:_request. Every caller in this repo passes an explicit timeout_s (startup reconciliation must never block on a wedged daemon).

Source code in tit/jobs/docker_engine.py
def list_containers(
    self,
    *,
    filters: dict[str, list[str]] | None = None,
    all_containers: bool = False,
    timeout_s: float | None | object = "default",
) -> list[dict[str, Any]]:
    """``GET /containers/json`` — the Engine-API equivalent of ``docker ps [--filter …]``.

    *filters* is the Engine's own filter map (e.g. ``{"label": ["tit.job_id=abc"]}``); it is
    JSON-encoded into the query string by :meth:`_request`. Every caller in this repo passes
    an explicit *timeout_s* (startup reconciliation must never block on a wedged daemon).
    """
    self._ensure_negotiated()
    result = self._json_request(
        "GET",
        self._vpath("/containers/json"),
        query={
            "all": "1" if all_containers else "0",
            "filters": filters or None,
        },
        timeout_s=timeout_s,
    )
    return list(result or [])

discover

discover(env: dict[str, str] | None = None) -> DockerConnection

DOCKER_HOST=unix://... if set, else the well-known DooD bind-mount path. See the module docstring for why this is intentionally narrower than discover.ts's multi-candidate, CLI-assisted search — this module always runs inside a container with a known socket mount.

Source code in tit/jobs/docker_engine.py
def discover(env: dict[str, str] | None = None) -> DockerConnection:
    """``DOCKER_HOST=unix://...`` if set, else the well-known DooD bind-mount path. See the module
    docstring for why this is intentionally narrower than ``discover.ts``'s multi-candidate,
    CLI-assisted search — this module always runs inside a container with a known socket mount.
    """
    resolved_env = os.environ if env is None else env
    host = resolved_env.get("DOCKER_HOST")
    if host and host.startswith("unix://"):
        path = host[len("unix://") :] or DEFAULT_SOCKET_PATH
        return DockerConnection(socket_path=path)
    return DockerConnection(socket_path=DEFAULT_SOCKET_PATH)

run_job_container

run_job_container(client: DockerEngineClient, *, image: str, cmd: list[str], job_id: str | None = None, name: str | None = None, env: dict[str, str] | None = None, binds: Iterable[str] = (), labels: dict[str, str] | None = None, platform: str | None = None, cpus: float | None = None, memory_bytes: int | None = None, auto_remove: bool = False) -> str

create_container + start_container, returning the container id immediately. Always sets Labels["tit.job_id"] when job_id is given — see the module docstring for how tit/pre/qsi/docker_builder.py's own _label_args() sets the same label independently today, and what migrating it onto this client would still change (the CLI shell-out).

Source code in tit/jobs/docker_engine.py
def run_job_container(
    client: DockerEngineClient,
    *,
    image: str,
    cmd: list[str],
    job_id: str | None = None,
    name: str | None = None,
    env: dict[str, str] | None = None,
    binds: Iterable[str] = (),
    labels: dict[str, str] | None = None,
    platform: str | None = None,
    cpus: float | None = None,
    memory_bytes: int | None = None,
    auto_remove: bool = False,
) -> str:
    """``create_container`` + ``start_container``, returning the container id immediately. Always
    sets ``Labels["tit.job_id"]`` when *job_id* is given — see the module docstring for how
    ``tit/pre/qsi/docker_builder.py``'s own ``_label_args()`` sets the same label independently
    today, and what migrating it onto this client would still change (the CLI shell-out).
    """
    resolved_labels = dict(labels or {})
    if job_id:
        resolved_labels["tit.job_id"] = job_id
    spec = ContainerCreateSpec(
        image=image,
        cmd=cmd,
        env=env or {},
        labels=resolved_labels,
        platform=platform,
        binds=list(binds),
        auto_remove=auto_remove,
        nano_cpus=int(cpus * 1e9) if cpus is not None else None,
        memory_bytes=memory_bytes,
    )
    container_id = client.create_container(spec, name=name)
    client.start_container(container_id)
    return container_id

run_qsiprep_example

run_qsiprep_example(client: DockerEngineClient, job_id: str, host_project_dir: str, subject_id: str) -> str

Documentation only (not called by production code): the call shape a tit/pre/qsi/docker_builder.py migration would make, using the exact same /data:ro / /out / /work mount points as today's DockerPaths dataclass, and the same tit.job_id label build_qsiprep_cmd already sets via _label_args — the migration this documents is about the CLI shell-out, not the label.

Source code in tit/jobs/docker_engine.py
def run_qsiprep_example(
    client: DockerEngineClient, job_id: str, host_project_dir: str, subject_id: str
) -> str:
    """Documentation only (not called by production code): the call shape a
    ``tit/pre/qsi/docker_builder.py`` migration would make, using the exact same
    ``/data``:ro / ``/out`` / ``/work`` mount points as today's ``DockerPaths`` dataclass, and the
    same ``tit.job_id`` label ``build_qsiprep_cmd`` already sets via ``_label_args`` — the
    migration this documents is about the CLI shell-out, not the label."""
    return run_job_container(
        client,
        image="pennlinc/qsiprep:26.0.0",
        cmd=["/data", "/out", "participant", "--participant-label", subject_id],
        job_id=job_id,
        env={"OMP_NUM_THREADS": "4"},
        binds=[
            f"{host_project_dir}:/data:ro",
            f"{host_project_dir}/derivatives/qsiprep:/out",
            f"{host_project_dir}/derivatives/.qsiprep_work:/work",
        ],
        platform="linux/amd64",
    )