Source code for flowli.api.catalog

"""The workflow catalog. See specs/09-http-api.md section 7.

This is what is left of Flowli v1's `taskflow` package: a JSON Schema
derived from the signature of a workflow function, plus the defaults a start
uses. It holds no state and reads the `Registry` of this process.
"""

from __future__ import annotations

import inspect
from dataclasses import dataclass
from typing import Any, get_type_hints

from pydantic import ValidationError, create_model

from flowli.runtime.registry import Registry


@dataclass(frozen=True, slots=True)
class CatalogEntry:
    name: str
    version: str
    summary: str | None
    description: str | None
    queue: str
    schema: dict[str, Any] | None

    @property
    def has_schema(self) -> bool:
        return self.schema is not None

    def summary_dict(self) -> dict[str, Any]:
        return {
            "name": self.name,
            "version": self.version,
            "summary": self.summary,
            "queue": self.queue,
            "has_schema": self.has_schema,
        }

    def detail_dict(self) -> dict[str, Any]:
        return {**self.summary_dict(), "description": self.description, "schema": self.schema}


def _docstring(fn: Any) -> tuple[str | None, str | None]:
    doc = inspect.getdoc(fn)
    if not doc:
        return None, None
    summary, _, rest = doc.partition("\n")
    return summary.strip() or None, doc if rest.strip() else None


def _model(name: str, fn: Any) -> Any | None:
    """A pydantic model of the arguments, or None when they cannot be typed.

    The first parameter is the `Context` and is skipped. A parameter with no
    usable annotation makes the whole schema unavailable: a half-checked
    argument list is worse than an unchecked one, because it looks checked.
    """
    try:
        hints = get_type_hints(fn)
    except Exception:
        return None
    parameters = list(inspect.signature(fn).parameters.values())[1:]
    fields: dict[str, Any] = {}
    for p in parameters:
        if p.kind in (p.VAR_POSITIONAL, p.VAR_KEYWORD):
            return None
        annotation = hints.get(p.name)
        if annotation is None:
            return None
        default = ... if p.default is inspect.Parameter.empty else p.default
        fields[p.name] = (annotation, default)
    try:
        return create_model(f"{name}_Args", **fields)
    except Exception:
        return None


[docs] class Catalog: """Every registered workflow, with its argument schema. Built once, at startup: a schema never changes while the process runs. """ def __init__(self, registry: Registry, *, default_queue: str = "default") -> None: self.registry = registry self.default_queue = default_queue self._entries: dict[tuple[str, str], CatalogEntry] = {} self._models: dict[tuple[str, str], Any] = {} self._build() def _build(self) -> None: for ref in self.registry.entries(): name, version = ref.name, ref.version summary, description = _docstring(ref.fn) model = _model(f"{name}_{version}", ref.fn) if model is not None: self._models[(name, version)] = model self._entries[(name, version)] = CatalogEntry( name=name, version=version, summary=summary, description=description, queue=self.default_queue, schema=None if model is None else model.model_json_schema(), )
[docs] def entries(self) -> list[CatalogEntry]: return list(self._entries.values())
[docs] def entry(self, name: str, version: str) -> CatalogEntry: ref = self.registry.get(name, version) # raises WorkflowNotRegistered return self._entries[(ref.name, ref.version)]
[docs] def validate(self, name: str, version: str, args: dict[str, Any]) -> dict[str, Any]: """The arguments, checked against the schema. Unchecked when there is none.""" model = self._models.get((name, version)) if model is None: return args return dict(model(**args).model_dump())
__all__ = ["Catalog", "CatalogEntry", "ValidationError"]