Skip to content

teff.observability.builder

teff.observability.builder

Build a :class:GraphObserver from a workflow's observability: block.

The block is declarative data inside a workflow.yaml::

observability:
  db: ./data/traces.db          # our SQLite dashboard (relative to the workflow)
  export:                       # optional fan-out to remote sinks
    - type: langfuse
      public_key_env: LANGFUSE_PUBLIC_KEY
      secret_key_env: LANGFUSE_SECRET_KEY
      host: https://cloud.langfuse.com
    - type: langsmith
      api_key_env: LANGCHAIN_API_KEY
      project: my-project
    - type: webhook
      url: https://hooks.example.com/traces

Relative paths resolve against the workflow file's directory and secrets come from environment variables, never from the file itself.

Functions:

Name Description
build_observability

Assemble a :class:GraphObserver from an observability: block.

build_observer_factory

Like :func:build_observability but for repeated runs.

build_remote_exporter

Construct a remote exporter from one observability.export entry.

build_observability

build_observability(config, *, base_dir='.', graph=None, name='workflow')

Assemble a :class:GraphObserver from an observability: block.

Returns None when the block is missing or declares no sinks. The observer is wired with the graph topology so remote exporters and the dashboard can render the flow.

Source code in teff/observability/builder.py
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
def build_observability(
    config: dict[str, Any] | None,
    *,
    base_dir: str = ".",
    graph=None,
    name: str = "workflow",
) -> GraphObserver | None:
    """Assemble a :class:`GraphObserver` from an ``observability:`` block.

    Returns ``None`` when the block is missing or declares no sinks.  The
    observer is wired with the graph topology so remote exporters and the
    dashboard can render the flow.
    """
    if not config:
        return None
    exporters = _build_exporters(config, base_dir=base_dir)
    if not exporters:
        return None
    return GraphObserver(
        name,
        exporter=CompositeExporter(exporters),
        topology=topology_from_graph(graph) if graph is not None else None,
    )

build_observer_factory

build_observer_factory(config, *, base_dir='.', graph=None, name='workflow')

Like :func:build_observability but for repeated runs.

Returns a zero-arg callable that yields a fresh observer per run while sharing one set of exporters — the right shape for a daemon that traces every tick. None when observability is not configured.

Source code in teff/observability/builder.py
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
def build_observer_factory(
    config: dict[str, Any] | None,
    *,
    base_dir: str = ".",
    graph=None,
    name: str = "workflow",
) -> "Callable[[], GraphObserver] | None":
    """Like :func:`build_observability` but for repeated runs.

    Returns a zero-arg callable that yields a *fresh* observer per run while
    sharing one set of exporters — the right shape for a daemon that traces
    every tick.  ``None`` when observability is not configured.
    """
    if not config:
        return None
    exporters = _build_exporters(config, base_dir=base_dir)
    if not exporters:
        return None
    composite = CompositeExporter(exporters)
    topology = topology_from_graph(graph) if graph is not None else None

    def factory() -> GraphObserver:
        return GraphObserver(name, exporter=composite, topology=topology)

    return factory

build_remote_exporter

build_remote_exporter(spec, *, base_dir='.')

Construct a remote exporter from one observability.export entry.

Source code in teff/observability/builder.py
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
def build_remote_exporter(
    spec: dict[str, Any], *, base_dir: str = "."
) -> TraceExporter:
    """Construct a remote exporter from one ``observability.export`` entry."""
    kind = spec.get("type")
    timeout = float(spec.get("timeout", 10.0))
    retries = int(spec.get("retries", 3))
    backoff = float(spec.get("backoff", 1.0))
    if kind == "webhook":
        url = spec.get("url")
        if not url and spec.get("url_env"):
            url = os.environ.get(spec["url_env"])
        if not url:
            raise ConfigError(
                "observability.export.webhook: 'url' or 'url_env' is required"
            )
        return HttpExporter(
            url,
            headers=dict(spec.get("headers") or {}),
            timeout=timeout,
            retries=retries,
            backoff=backoff,
        )
    if kind == "langfuse":
        host = spec.get("host")
        if not host:
            raise ConfigError("observability.export.langfuse: 'host' is required")
        public_key = _require_env(
            spec, "public_key_env", "LANGFUSE_PUBLIC_KEY", "langfuse"
        )
        secret_key = _require_env(
            spec, "secret_key_env", "LANGFUSE_SECRET_KEY", "langfuse"
        )
        return LangfuseExporter(
            host,
            public_key,
            secret_key,
            timeout=timeout,
            retries=retries,
            backoff=backoff,
        )
    if kind == "langsmith":
        api_url = (
            spec.get("api_url")
            or os.environ.get("LANGCHAIN_ENDPOINT")
            or "https://api.smith.langchain.com"
        )
        api_key = _require_env(spec, "api_key_env", "LANGCHAIN_API_KEY", "langsmith")
        project = spec.get("project") or os.environ.get("LANGCHAIN_PROJECT")
        return LangsmithExporter(
            api_url,
            api_key,
            project=project,
            timeout=timeout,
            retries=retries,
            backoff=backoff,
        )
    raise ConfigError(f"observability.export: unknown exporter type {kind!r}")