Skip to content

server

server

The clearance hub — varlink server + reader ingester + verdict exec.

Fans reader-emitted events (blocks, container lifecycle, shield state) out to every connected clearance client, and applies verdicts the clients send back by shelling out to terok-shield allow|deny. The only D-Bus in sight is what individual clients choose to use on their own (the desktop notifier reaches for org.freedesktop.Notifications out-of-band); the hub itself speaks plain unix-socket varlink.

Authorisation is structural: the socket is mode 0600 (same-UID only), and every Verdict call must cite a (container, request_id, dest) triple the hub actually emitted via connection_blocked. The triple is recorded at emit time and dropped on verdict or lifecycle change; anything that doesn't match is a UnknownRequest or VerdictTupleMismatch refusal.

VERDICT_TTL_CONFIG_KEY = 'verdict_ttl' module-attribute

VerdictTtlParam = float | None | _Unset module-attribute

ClearanceHub(*, clearance_socket=None, reader_socket=None, verdict_client=None, socket_context=None, verdict_ttl=_UNSET)

Server for the org.terok.Clearance1 interface.

Owns three pieces of state:

  • _subscribers — a set of bounded per-connection queues; the hub puts a ClearanceEvent on each one every time the reader ingester delivers an event. Slow clients see their oldest events dropped; fast clients aren't affected.
  • _live_verdicts — the request_id → (container, dest) map the Verdict method checks for the authz binding.
  • An EventIngester bound to the canonical reader socket.

Lifecycle: start brings everything up; stop tears it down under individual timeouts so a flaky bus or a stuck subscriber can't stall the per-container supervisor's teardown.

Configure the two sockets and the verdict-helper client.

verdict_client is injected so tests can stub out shield exec without spawning the helper process. Production callers leave it defaulted — a fresh VerdictClient pointing at the canonical helper socket.

socket_context — optional zero-arg callable returning a context manager around the varlink bind() (forwarded to bind_hardened). terok-sandbox passes its SELinux setsockcreatecon helper here so the hub socket is labelled terok_socket_t and confined containers can connectto it.

verdict_ttl — seconds after which a pending connection_blocked entry is expired from the authz-binding map. None means no expiry (backward-compatible). When not passed explicitly, the value is read from the clearance.verdict_ttl key in the layered terok config (~/.config/terok/config.yml); absent falls back to a 24-hour default so long-running containers don't accumulate unbounded state. Pass verdict_ttl=None explicitly to disable expiry regardless of the config file.

Config file example::

clearance:
  verdict_ttl: 24h   # or "never" to disable
Source code in src/terok_clearance/hub/server.py
def __init__(
    self,
    *,
    clearance_socket: Path | None = None,
    reader_socket: Path | None = None,
    verdict_client: VerdictClient | None = None,
    socket_context: Callable[[], AbstractContextManager[None]] | None = None,
    verdict_ttl: VerdictTtlParam = _UNSET,
) -> None:
    """Configure the two sockets and the verdict-helper client.

    ``verdict_client`` is injected so tests can stub out shield exec
    without spawning the helper process.  Production callers leave
    it defaulted — a fresh [`VerdictClient`][terok_clearance.hub.server.VerdictClient] pointing at the
    canonical helper socket.

    *socket_context* — optional zero-arg callable returning a context
    manager around the varlink ``bind()`` (forwarded to
    [`bind_hardened`][terok_clearance.wire.socket.bind_hardened]).
    terok-sandbox passes its SELinux ``setsockcreatecon`` helper here
    so the hub socket is labelled ``terok_socket_t`` and confined
    containers can ``connectto`` it.

    *verdict_ttl* — seconds after which a pending ``connection_blocked``
    entry is expired from the authz-binding map.  ``None`` means no
    expiry (backward-compatible).  When not passed explicitly, the
    value is read from the ``clearance.verdict_ttl`` key in the
    layered terok config (``~/.config/terok/config.yml``); absent
    falls back to a 24-hour default so long-running containers
    don't accumulate unbounded state.  Pass ``verdict_ttl=None``
    explicitly to disable expiry regardless of the config file.

    Config file example::

        clearance:
          verdict_ttl: 24h   # or "never" to disable
    """
    self._clearance_socket = clearance_socket or default_clearance_socket_path()
    self._reader_socket = reader_socket  # None → EventIngester picks its default.
    self._verdict_client = verdict_client or VerdictClient()
    self._socket_context = socket_context
    ttl: float | None
    if verdict_ttl is _UNSET:
        ttl = _read_verdict_ttl_from_config()
    elif verdict_ttl is None:
        ttl = None
    elif isinstance(verdict_ttl, (int, float)):
        ttl = float(verdict_ttl)
    else:  # pragma: no cover — defensive; type narrowing should prevent this
        ttl = _DEFAULT_VERDICT_TTL
    self._verdict_ttl: float | None = ttl

    self._subscribers: set[asyncio.Queue[ClearanceEvent]] = set()
    # request_id → (container, dest) the hub emitted in the matching
    # ConnectionBlocked; Verdict calls must cite a triple that matches.
    self._live_verdicts: dict[str, tuple[str, str]] = {}
    # request_id → monotonic timestamp when the entry was recorded;
    # used by the background cleanup to expire stale entries.
    self._live_verdicts_ts: dict[str, float] = {}
    self._cleanup_task: asyncio.Task[None] | None = None

    self._ingester: EventIngester | None = None
    self._varlink_server: VarlinkUnixServer | None = None

verdict_ttl property

Pending-verdict TTL in seconds, or None when expiry is disabled.

start() async

Bring the ingester + varlink server online and accept clients.

Transactional: if the varlink bind fails after the ingester is already listening, the ingester is stopped before the exception propagates so a half-started hub doesn't leak a live reader-side socket on systemd restart paths.

Source code in src/terok_clearance/hub/server.py
async def start(self) -> None:
    """Bring the ingester + varlink server online and accept clients.

    Transactional: if the varlink bind fails after the ingester is
    already listening, the ingester is stopped before the exception
    propagates so a half-started hub doesn't leak a live
    reader-side socket on systemd restart paths.
    """
    self._ingester = EventIngester(
        socket_path=self._reader_socket or _default_reader_socket(),
        on_event=self._relay_reader_event,
    )
    await self._ingester.start()
    try:
        registry = VarlinkInterfaceRegistry()
        registry.register_interface(
            Clearance1Interface(
                event_stream_factory=self._subscribe,
                apply_verdict=self._apply_verdict,
            )
        )
        registry.register_interface(
            VarlinkServiceInterface(
                vendor="terok",
                product="terok-clearance",
                version=_own_version(),
                url="https://github.com/terok-ai/terok-clearance",
                registry=registry,
            )
        )

        from terok_clearance.wire.socket import bind_hardened

        async def _factory(path: str) -> object:
            return await create_unix_server(registry.protocol_factory, path=path)

        self._varlink_server = await bind_hardened(
            _factory,
            self._clearance_socket,
            "clearance",
            socket_context=self._socket_context,
        )
    except BaseException:
        with contextlib.suppress(Exception):
            await self._ingester.stop()
        self._ingester = None
        raise
    # Start the background cleanup task for stale verdict entries.
    # Only runs when a TTL is configured; otherwise entries live until
    # a verdict or lifecycle event removes them (backward-compatible).
    if self._verdict_ttl is not None:
        self._cleanup_task = asyncio.create_task(self._cleanup_loop())
    _log.info(
        "clearance hub online at %s (verdict_ttl=%s)",
        self._clearance_socket,
        f"{self._verdict_ttl}s" if self._verdict_ttl is not None else "never",
    )

stop() async

Close the varlink server + ingester; drain subscriber queues.

Source code in src/terok_clearance/hub/server.py
async def stop(self) -> None:
    """Close the varlink server + ingester; drain subscriber queues."""
    if self._varlink_server is not None:
        # ``close()`` on its own only stops accepting new connections;
        # existing subscribers would sit forever in ``queue.get()`` and
        # ``wait_closed`` would hang until the timeout fires.
        # ``close_clients()`` walks the live transports and closes them,
        # which makes the server-side ``_call_async_method_more``'s
        # next ``send_reply`` fail with OSError — that in turn calls
        # ``generator.aclose()`` on the subscriber, propagating cleanly
        # through to our ``finally`` block.  This avoids the
        # assertion asyncvarlink fires when a streaming generator
        # ends "normally" with ``continues=True`` on the last reply.
        self._varlink_server.close()
        with contextlib.suppress(AttributeError):
            self._varlink_server.close_clients()
        with contextlib.suppress(TimeoutError, Exception):
            await asyncio.wait_for(self._varlink_server.wait_closed(), timeout=1.0)
        self._varlink_server = None
    if self._ingester is not None:
        with contextlib.suppress(Exception):
            await self._ingester.stop()
        self._ingester = None
    with contextlib.suppress(Exception):
        await self._verdict_client.stop()
    if self._cleanup_task is not None:
        self._cleanup_task.cancel()
        with contextlib.suppress(asyncio.CancelledError):
            await self._cleanup_task
        self._cleanup_task = None
    self._subscribers.clear()
    self._live_verdicts.clear()
    self._live_verdicts_ts.clear()

parse_verdict_ttl(value)

Parse a TTL string into seconds, or None for "no expiry".

Accepts plain seconds (3600), unit-suffixed (24h, 30m, 7d), and the special values never / 0 / none (all meaning "no expiry"). None input returns the default (24h). Raises ValueError on malformed input.

Source code in src/terok_clearance/hub/server.py
def parse_verdict_ttl(value: str | None) -> float | None:
    """Parse a TTL string into seconds, or ``None`` for "no expiry".

    Accepts plain seconds (``3600``), unit-suffixed (``24h``, ``30m``,
    ``7d``), and the special values ``never`` / ``0`` / ``none`` (all
    meaning "no expiry").  ``None`` input returns the default (24h).
    Raises ``ValueError`` on malformed input.
    """
    if value is None:
        return _DEFAULT_VERDICT_TTL
    v = value.strip().lower()
    if v in ("never", "none", "0", "off", "disabled"):
        return None
    m = _TTL_RE.fullmatch(v)
    if not m:
        raise ValueError(f"invalid TTL {value!r} (expected e.g. '3600', '24h', '30m', 'never')")
    num = float(m.group(1))
    unit = m.group(2)
    return num * _TTL_UNITS[unit]

serve() async

Run the hub service until SIGINT/SIGTERM.

The entry point terok-clearance-hub serve hands off here. Blocks forever on a signal-set asyncio.Event; the first SIGINT/SIGTERM flips it, then stop tears down the server under a timeout.

Source code in src/terok_clearance/hub/server.py
async def serve() -> None:  # pragma: no cover — integration path
    """Run the hub service until SIGINT/SIGTERM.

    The entry point ``terok-clearance-hub serve`` hands off here.  Blocks forever
    on a signal-set [`asyncio.Event`][asyncio.Event]; the first SIGINT/SIGTERM
    flips it, then [`stop`][terok_clearance.hub.server.ClearanceHub.stop]
    tears down the server under a timeout.
    """
    from terok_clearance.runtime.service import configure_logging, wait_for_shutdown_signal

    configure_logging()
    hub = ClearanceHub()
    await hub.start()
    try:
        await wait_for_shutdown_signal()
    finally:
        await hub.stop()