> ## Content Index
> Fetch the complete content index at: https://www.sfrt.io/llms.txt
> Use this file to discover other available public pages before exploring further.

# Towards self-healing data pipelines: dltHub inspector in ❄️ Cortex
- URL: https://www.sfrt.io/towards-self-healing-data-pipelines-dlthub-inspector-in-cortex/
- Published: 2026-09-28T06:41:10.000Z
- Updated: 2026-09-28T06:41:10.000Z
- Description: A Snowflake Cortex Agent for dltHub jobs
- Author: Martin Seifert
- Tags: Infrastructure, Lakehouse

TL/DR: Let's build `META.CORTEX.JOB_INSPECTOR`, a Snowflake Cortex Agent that pulls failed dltHub job run records, logs and pipeline traces, investigates what broke, and notifies me! Everything is in here: the egress, the functions, the agent spec, the dltHub job. This step-by-step setup guide covers rebuilding the whole thing from scratch.

---

On the title: I am aware a notification is not a self-healing pipeline 😜 However, I am currently working on the next step (basically allowing the same agent to fix the pipeline if possible), so I'm on the road towards that goal 🤓 Come along, the ride is pretty fun!

None of this would exist without dltHub, and not just as a footnote. More on that in a moment.

## Why this exists, and whose idea it actually was

A few days ago, I had the enormous pleasure of participating in an [internal dltHub hackathon](https://lnkd.in/p/ekHs3-6M?ref=sfrt.io) (unrelated side note: This is a bunch of really smart and talented people! No idea how I convinced them to allow me in...). Shortly thereafter, they switched on a private preview of background agents for my workspace. The first "verified" agent (and rather obvious candidate for this) is a `job-inspector`: an `AGENT.md` file (system prompt + tool and output-schema) that runs on dltHub's internal engine (I believe it's Modal but could be Tower or something similar just as well) as a [pydantic-ai](https://ai.pydantic.dev/?ref=sfrt.io) agent, talking to whichever LLM API the user configures. On a `job.fail` event it wakes up, reads the failed run's logs (plus the actual source code... benefit of running inside the workspace), sorts the failure into one of a handful of buckets, and proposes a fix. Genuinely good, and I wired the output to Teams:

[dltHub notifies Slack and email. But my org runs on Teams 😩dltHub ships Slack and email alerts out of the box. My org runs on Teams, so I built the third one myself and wired it into all \~70 pipelines with one decorator shadow.![](https://storage.ghost.io/c/7d/94/7d942fe1-7868-4a1a-b2c9-4eb415b1a546/content/images/icon/headshot_ring-55374fbb-38c3-41bc-af89-31157a0d7eb4.png)sfrt.ioMartin Seifert![](https://storage.ghost.io/c/7d/94/7d942fe1-7868-4a1a-b2c9-4eb415b1a546/content/images/thumbnail/warning-ed6aae93-e932-4e65-a3d8-07a89143bd54.png)](https://www.sfrt.io/dlthub-notifies-slack-and-email-but-my-org-runs-on-teams/)

The classification rubric in that `AGENT.md` (categories, confidence levels, earliest error first, separate the job's code from the platform's) is pretty much what this whole post is built around. I didn't invent any of it... I just ported it (almost verbatim) into a second agent running somewhere else 😜

That "somewhere else" is ❄️ Snowflake Cortex, for one (and a half) reason: The goal (outside the scope of this post!) is to arrive at an agent that will be capable of investigating the full pipeline chain (also: automatically heal the full pipeline, ofc. 😅). A dltHub agent can view the ingestion, but not what depends on the ingested data or which tasks broke because of a failure. Also, the Snowflake perimeter (that's the half reason): Even if procurement allowed Snowflake agents, they might not be too fond of allowing others, too. Hence, being able to run this functionality behind the perimeter might be valuable to some folks.

So the question became whether the same job-inspector could run as a Cortex Agent instead. Short answer: yes, and this post tells the story of how to get there.

Once again: everything downstream of this paragraph exists only because dltHub documents a REST API for runs, jobs and logs, ships `dlthub_sdk` in every workspace, and allowed me to copy the reasoning behind its own agent as a plain, readable `AGENT.md` (and by "allowed" I mean: it's a .md file in a publicly accessible [GitHub repo](https://github.com/dlt-hub/example-background-agents/tree/main?ref=sfrt.io) without any explicit license, yet... please don't sue me 😅). What follows is my Cortex-flavored port of dltHub's own work.

## The shape of it

![](https://storage.ghost.io/c/7d/94/7d942fe1-7868-4a1a-b2c9-4eb415b1a546/content/images/2026/09/2026-09-27-22_42_11-SharePoint-_-Eric-Trittelvitz-_-Pro-Juventute-_-martin.seifert@projuventute.ch-_.png)

The agent combines 4 areas, each covered below with enough detail to rebuild the whole thing from scratch: Snowflake's egress to dltHub, 5 tool functions, the agent itself, and the dltHub job that ties them together. A few decisions apply across all of those areas:

- **Read-only.** The agent's tools are functions that read: no shell, no file access, not even any warehouse writes. dltHub's own agent can run shell commands (read-only by instruction rather than by physical limit), the Cortex one can't do anything else.
- **Identify the failed run in code.** The dltHub inspector job resolves its own run's `prev_run_id` (the run whose failure woke it) before the agent is ever called. Having an LLM search for the right run burns tokens for nothing.
- **Cap every tool result.** I had a case of dltHub's agent spending 332k tokens on a single diagnosis, reading its way through the platform's installed source code. Each of my UDFs caps its output at 8 to 30k characters, and dependency-install noise is dropped from logs before the model sees any of it in an attempt to reduce the noise.
- **One JSON output object as the contract.** I use the same keys as dltHub's `job-inspector` output, so the same Teams-card function can serve both agents' results.
- **Never silent.** Should the agent call fail for any reason (auth, timeout, garbage JSON) a card still goes out saying the diagnosis is unavailable.
- **The two inspectors ignore each other.** I (for now) run both the dltHub inspector and my Cortex version of it. Each watches every job, except the inspector jobs... so a pair of broken inspectors doesn't happily trigger each other forever 😅

## Layer 1: getting Snowflake to talk to dltHub at all

dltHub via `dlthub_sdk` splits access to the workspace into two planes, each needs to be included in the Snowflake egress rule:

| plane         | host                                                  | what lives there                            | auth                                                                      |
| ------------- | ----------------------------------------------------- | ------------------------------------------- | ------------------------------------------------------------------------- |
| control plane | api.dlthub.com                                        | runs, jobs ("scripts"), workspace metadata  | Authorization: Bearer <api key>                                           |
| data plane    | the workspace's dataplane\_url (mine: eu1.dlthub.com) | run logs, telemetry (pipeline runs, traces) | short-lived token from GET /api/v1/workspaces/{ws}/dataplane-access-token |

I created a dedicated **viewer** API key (read-only, workspace-scoped, created separately from the key I use to trigger pipeline runs from Snowflake tasks) covers everything, including the data-plane token. 

I kind of expected a 403 on that last part, but no, it just worked 😅

Network rule, secret, external access integration in ❄️ ▾ use role sysadmin; create or replace network rule meta.integration.nr\_dlt mode = 'EGRESS' type = 'HOST\_PORT' value\_list = ('api.dlthub.com:443', 'eu1.dlthub.com:443'); -- control plane + data plane create or replace secret meta.integration.se\_dlthub\_viewer\_api\_key type = generic\_string secret\_string = ''; use role accountadmin; create or replace external access integration i\_dlthub allowed\_network\_rules = (meta.integration.nr\_dlt) allowed\_authentication\_secrets = (meta.integration.se\_dlthub\_viewer\_api\_key) enabled = true; grant usage on integration i\_dlthub to role sysadmin; 

If `nr_dlt` or `i_dlthub` already exist for another purpose (say, triggering dltHub runs from Snowflake tasks), use `alter ... set value_list / alter ... set allowed_authentication_secrets` instead, repeating the existing entries: both statements replace the whole list rather than append.

## Layer 2: five read-only tool functions

The agent's toolbox is 5 Python functions (they only read, so no Snowpark session needed here):

| tool name (agent)    | function                                       | returns                                                                                  |
| -------------------- | ---------------------------------------------- | ---------------------------------------------------------------------------------------- |
| Get\_Run             | F\_DLTHUB\_GET\_RUN(run\_id)                   | status, trigger, profile, times, prev\_run\_id, job ref, timeout, per-pipeline summaries |
| Get\_Run\_Logs       | F\_DLTHUB\_GET\_RUN\_LOGS(run\_id, max\_lines) | last N log lines, line\_num \[phase\] content                                            |
| List\_Runs           | F\_DLTHUB\_LIST\_RUNS(job\_ref, max\_runs)     | the job's recent runs (the neighbour check)                                              |
| Get\_Job             | F\_DLTHUB\_GET\_JOB(job\_ref)                  | job definition: triggers, execute/require spec, paused, last run                         |
| Get\_Pipeline\_Trace | F\_DLTHUB\_GET\_PIPELINE\_TRACE(run\_id)       | per dlt pipeline run: failed step, error message, each step's exception                  |

If the role creating those functions has the database role `SNOWFLAKE.PYPI_REPOSITORY_USER`, each function can utilize the dltHub SDK. Otherwise, a fallback to raw REST using `requests` instead would also work, but would require rebuilding the stuff the SDK takes care of:

- **Data-plane token.** Logs and telemetry live on a separate data plane, and calling it needs a short-lived token minted from the control plane. The SDK does that and finds the data-plane host from the workspace payload.
- **Job refs vs. uuids.** Listing a job's runs needs the job's uuid. The SDK takes the job ref.
- **Naive datetimes.** The `pipeline-runs` endpoint rejects timestamps with an offset. The SDK doesn't care.
- **Log parsing.** The logs endpoint returns NDJSON, one object per line. The SDK gives typed lines (`line_num`, `phase`, `content`).

Output caps are baked into every function. Tool results pile up in the conversation, so one careless full-log dump gets paid for again on every following turn. An 8 to 30k character cap, plus dropping setup-phase lines (dozens of `+ package==x.y` install entries that bury the actual traceback), keeps a typical diagnosis small: the logs of real failed runs came back at an average 5.7k characters.

DDL for all 5 UDFs in ❄️ ▾ 

```sql
create or replace function meta.cortex.f_dlthub_get_run(run_id varchar)
copy grants
returns varchar
language python
runtime_version = '3.13'
artifact_repository = snowflake.snowpark.pypi_shared_repository
packages = ('dlthub-client')
handler = 'main'
comment = 'JOB_INSPECTOR tool: one dltHub job run record'
external_access_integrations = (i_dlthub)
secrets = ('dlthub' = meta.integration.se_dlthub_viewer_api_key)
as
$$
import json
import _snowflake
import dlthub_sdk

WS = "[WORKSPACE_ID]"

def workspace():
    rt = dlthub_sdk.connect(token=_snowflake.get_generic_secret_string("dlthub"), base_url="https://api.dlthub.com")
    return rt.workspaces.get(id=WS)

def dump(obj, cap):
    return json.dumps(obj, default=lambda o: getattr(o, "value", str(o)))[:cap]

def main(run_id):
    try:
        d = workspace().job_runs.get(id=run_id).to_dict()
    except Exception as e:
        return f"ERROR: {type(e).__name__}: {str(e)[:300]}"
    d["pipelines"] = (d.get("pipelines") or [])[:10]
    return dump(d, 8000)
$$;
```

```sql
create or replace function meta.cortex.f_dlthub_list_runs(job_ref varchar, max_runs number)
copy grants
returns varchar
language python
runtime_version = '3.13'
artifact_repository = snowflake.snowpark.pypi_shared_repository
packages = ('dlthub-client')
handler = 'main'
comment = 'JOB_INSPECTOR tool: the most recent runs of one dltHub job, newest first, max 20'
external_access_integrations = (i_dlthub)
secrets = ('dlthub' = meta.integration.se_dlthub_viewer_api_key)
as
$$
import json
import _snowflake
import dlthub_sdk

WS = "[WORKSPACE_ID]"

def workspace():
    rt = dlthub_sdk.connect(token=_snowflake.get_generic_secret_string("dlthub"), base_url="https://api.dlthub.com")
    return rt.workspaces.get(id=WS)

def dump(obj, cap):
    return json.dumps(obj, default=lambda o: getattr(o, "value", str(o)))[:cap]

def main(job_ref, max_runs):
    try:
        job = workspace().jobs.get(ref=job_ref)
        limit = max(1, min(int(max_runs or 10), 20))
        keep = ("number", "id", "status", "trigger", "started_at", "ended_at", "duration_seconds")
        runs = [{k: r.to_dict().get(k) for k in keep} for r in job.runs.list(limit=limit)]
    except Exception as e:
        return f"ERROR: {type(e).__name__}: {str(e)[:300]}"
    return dump({"job_ref": job_ref, "runs": runs}, 8000)
$$;
```

```sql
create or replace function meta.cortex.f_dlthub_get_job(job_ref varchar)
copy grants
returns varchar
language python
runtime_version = '3.13'
artifact_repository = snowflake.snowpark.pypi_shared_repository
packages = ('dlthub-client')
handler = 'main'
comment = 'JOB_INSPECTOR tool: one dltHub job definition'
external_access_integrations = (i_dlthub)
secrets = ('dlthub' = meta.integration.se_dlthub_viewer_api_key)
as
$$
import json
import _snowflake
import dlthub_sdk

WS = "[WORKSPACE_ID]"

def workspace():
    rt = dlthub_sdk.connect(token=_snowflake.get_generic_secret_string("dlthub"), base_url="https://api.dlthub.com")
    return rt.workspaces.get(id=WS)

def dump(obj, cap):
    return json.dumps(obj, default=lambda o: getattr(o, "value", str(o)))[:cap]

def main(job_ref):
    try:
        job = workspace().jobs.get(ref=job_ref)
        d = job.to_dict()
        keep = ("job_ref", "name", "job_type", "paused", "archived", "default_trigger", "triggers", "next_run_at", "definition")
        out = {k: d.get(k) for k in keep}
        latest = [r.to_dict() for r in job.runs.list(limit=1)]
        out["latest_run"] = {k: latest[0].get(k) for k in ("number", "id", "status", "started_at")} if latest else None
    except Exception as e:
        return f"ERROR: {type(e).__name__}: {str(e)[:300]}"
    return dump(out, 12000)
$$;
```

```sql
create or replace function meta.cortex.f_dlthub_get_run_logs(run_id varchar, max_lines number)
copy grants
returns varchar
language python
runtime_version = '3.13'
artifact_repository = snowflake.snowpark.pypi_shared_repository
packages = ('dlthub-client')
handler = 'main'
comment = 'JOB_INSPECTOR tool: the last N log lines of a dltHub job run (default 200, max 1000)'
external_access_integrations = (i_dlthub)
secrets = ('dlthub' = meta.integration.se_dlthub_viewer_api_key)
as
$$
import _snowflake
import dlthub_sdk

WS = "[WORKSPACE_ID]"

def workspace():
    rt = dlthub_sdk.connect(token=_snowflake.get_generic_secret_string("dlthub"), base_url="https://api.dlthub.com")
    return rt.workspaces.get(id=WS)

def main(run_id, max_lines):
    try:
        lines = list(workspace().job_runs.get(id=run_id).logs())
    except Exception as e:
        # logs are consolidated only after the run ended: an unfinished run raises here
        return f"ERROR: {type(e).__name__}: {str(e)[:300]}"
    # dependency-install noise (hundreds of "+ package" lines) buries the traceback
    program = [l for l in lines if str(getattr(l.phase, "value", l.phase)) != "setup"]
    kept = program if program else lines
    n = max(1, min(int(max_lines or 200), 1000))
    tail = kept[-n:]
    header = f"total_lines={len(lines)} shown={len(tail)} setup_dropped={len(lines) - len(kept)}\n"
    body = "\n".join(f"{l.line_num} [{getattr(l.phase, 'value', l.phase)}] {l.content}" for l in tail)
    return header + body[-30000:]
$$;
```

```sql
create or replace function meta.cortex.f_dlthub_get_pipeline_trace(run_id varchar)
copy grants
returns varchar
language python
runtime_version = '3.13'
artifact_repository = snowflake.snowpark.pypi_shared_repository
packages = ('dlthub-client')
handler = 'main'
comment = 'JOB_INSPECTOR tool: per dlt pipeline run of a job run: status, failed step, error, step exceptions'
external_access_integrations = (i_dlthub)
secrets = ('dlthub' = meta.integration.se_dlthub_viewer_api_key)
as
$$
import json
import _snowflake
import dlthub_sdk

WS = "[WORKSPACE_ID]"

def workspace():
    rt = dlthub_sdk.connect(token=_snowflake.get_generic_secret_string("dlthub"), base_url="https://api.dlthub.com")
    return rt.workspaces.get(id=WS)

def dump(obj, cap):
    return json.dumps(obj, default=lambda o: getattr(o, "value", str(o)))[:cap]

def main(run_id):
    try:
        pipelines = workspace().telemetry.pipeline_runs
        out = []
        for p in list(pipelines.list(job_run_id=run_id, limit=3)):
            entry = {
                "pipeline_name": p.pipeline_name,
                "status": p.status,
                "error_step": p.error_step,
                "rows_loaded": p.rows_loaded,
                "error_message": (p.error_message or "")[:2000] or None,
            }
            # resolved_config_values deliberately left out: it holds configuration
            steps = (pipelines.trace(id=p.id).trace or {}).get("steps", [])
            entry["steps"] = [
                {"step": s.get("step"), "exception": (str(s.get("step_exception"))[:1500] if s.get("step_exception") else None)}
                for s in steps
            ]
            out.append(entry)
    except Exception as e:
        return f"ERROR: {type(e).__name__}: {str(e)[:300]}"
    if not out:
        return "no dlt pipeline run recorded for this job run (non-pipeline job, or it failed before any pipeline started)"
    return dump(out, 15000)
$$;
```

I built this for myself, so my workspace is hard coded in each function. If someone wants to make this more flexible, maybe parametrize the workspace id 😜

Errors return as strings starting with `ERROR:` rather than raised exceptions. A raised exception ends the agent's tool call with a generic failure, while an error string is something the model can actually read and reason about ("no stored log for this run yet"). A few months ago, I would have called that sloppy typing, but these days I call it a tool contract 😅

Before an agent ever touches those functions, plain SQL against a known failed run can test them:

```sql
select meta.cortex.f_dlthub_get_run('<failed run id>');
select meta.cortex.f_dlthub_get_run_logs('<failed run id>', 200);
select meta.cortex.f_dlthub_list_runs('jobs.__deployment__.<job>', 5);
select meta.cortex.f_dlthub_get_job('jobs.__deployment__.<job>');
select meta.cortex.f_dlthub_get_pipeline_trace('<run id of a job that ran a pipeline>');
```

With those 5 functions, the Cortex agent would have one disadvantage compared to the dltHub-native agent: It doesn't have access to the pipeline's source code. To my knowledge, the code (currently) isn't exposed via the dltHub SDK (or API or MCP or CLI), so an agent running outside the dltHub workspace would need to look at the code somewhere else. I have the source code in GitHub, so I built a sixth function (this time a procedure, actually) and agent-tool to allow the inspector to get it from there. This is not as portable as the other 5 functions/tools, therefore I separate it from them hereby. 

The agent setup and instructions below will not mention this sixth tool for that reason... If you want the source code in your agent, add it yourself 🤓 Also, this procedure doesn't allow a diff-check against what ran in dltHub (based on the traces), assuming the pipeline code in dltHub is actually what the agent looks at in GitHub. And why a procedure? This way I could re-use my existing GitHub integration in Snowflake without separately authenticating.

DDL for the 6th agent tool in ❄️ ▾ 

```sql
create or replace api integration i_github
  api_provider = git_https_api
  api_allowed_prefixes = ('https://github.com/[org]')
  allowed_authentication_secrets = all
  enabled = true;

create or replace secret meta.integration.se_github
  type = password
  username = '[github user]'
  password = '[fine-grained token, read-only on the one repo]';

create or replace git repository meta.integration.git_repo
  api_integration = i_github
  git_credentials = meta.integration.se_github
  origin = 'https://github.com/[org]/[repo]';

create or replace procedure meta.cortex.p_dlthub_read_source(path varchar)
copy grants
returns varchar
language python
runtime_version = '3.13'
packages = ('snowflake-snowpark-python')
handler = 'main'
comment = 'JOB_INSPECTOR tool: reads one pipeline source file (pipeline folder only) from the Git repo, line-numbered; fetches the repo first'
execute as owner
as
$$
import re

REPO = "META.INTEGRATION.GIT_REPO"
STAGE = f"@{REPO}/branches/master/"
ROOT = "dlt/dltHub/"          # the workspace root inside the repo
DENIED = {"tests", "_archive", "_local", "debug", "secrets"}
CAP = 30000

def repo_path(path):
    """Map a traceback frame or a relative path to a repo path, or raise ValueError."""
    p = (path or "").strip()
    if not p or not re.fullmatch(r"[A-Za-z0-9_./-]+", p):
        raise ValueError("path not allowed (only letters, digits and _ . / - are accepted)")
    if "/run/" in p:  # traceback frame, e.g. /tmp/dlt_run_x/run/facebook_ads/helpers.py
        p = p.split("/run/", 1)[1]
    elif p.startswith("/"):
        raise ValueError("path not allowed (absolute path)")
    if p.startswith("dlt/") and not p.startswith(ROOT):
        raise ValueError("path not allowed (only dlt/dltHub is readable)")
    if not p.startswith(ROOT):
        p = ROOT + p
    segments = p.split("/")
    if any(s in ("", ".", "..") or s.startswith(".") or s in DENIED for s in segments):
        raise ValueError("path not allowed")
    if not (p.endswith((".py", ".md")) or p == ROOT + "pyproject.toml"):
        raise ValueError("path not allowed (only .py, .md and pyproject.toml are readable)")
    return p

def main(session, path):
    try:
        p = repo_path(path)
    except ValueError as e:
        return f"ERROR: {e}"
    try:
        try:
            session.sql(f"alter git repository {REPO} fetch").collect()
            fetched = "true"
        except Exception:
            fetched = "false"  # e.g. GitHub credential expired: fall back to the checked-out commit
        try:
            f = session.file.get_stream(STAGE + p)
            try:
                content = f.read()
            finally:
                f.close()
        except Exception:
            return f"ERROR: {p} not found on master (fetched={fetched})"
        commit = next((r["commit_hash"][:8] for r in session.sql(f"show git branches in git repository {REPO}").collect() if r["name"] == "master"), "?")
    except Exception as e:
        return f"ERROR: {type(e).__name__}: {str(e)[:300]}"
    lines = content.decode("utf-8", errors="replace").splitlines()
    body = "\n".join(f"{i}: {line}" for i, line in enumerate(lines, 1))
    note = f"\n[truncated at {CAP} characters, file has {len(lines)} lines]" if len(body) > CAP else ""
    return f"path={p} repo_commit={commit} fetched={fetched} lines={len(lines)}\n{body[:CAP]}{note}"
$$;
```

## Layer 3: the agent

`META.CORTEX.JOB_INSPECTOR` runs on `claude-sonnet-5` with the five functions above wired in as `generic` tools. The instructions are, once more, a direct port of dltHub's `job-inspector` `AGENT.md`, minus everything about reading the workspace from a shell, because this agent has no shell:

- **Classification**: `config`, `credentials`, `upstream_data`, `code`, `resources`, `transient`, `unknown`, each with a one-line definition.
- **Confidence**: `high` when the earliest error names the cause and the evidence quotes it, `medium` when the cause is inferred with a plausible alternative left open, `low` for a guess.
- **Status**: `succeeded` means the cause was found (a fix isn't required), `failed` means it wasn't, and admitting that beats inventing a plausible story, `aborted` means the input didn't identify a run.
- **Method**: earliest error first, the job's own code separated from the platform's, neighbouring runs checked before anything gets called transient.

Added on top: a scope rule (investigate the workspace, don't attempt to reverse-engineer the platform), a tool budget (at most 8 tool calls as a target), and two sharper rules than the dltHub version: classify by where an error is *raised*, never by what its message claims, and never label a neighbouring run "the same failure" without comparing its actual error text.

A Teams card wants fields, not prose, so the response instruction demands exactly one JSON object with the same keys as dltHub's own inspector output, so a single card function can serve both agents:

```json
{
  "status": "succeeded | failed | aborted",
  "summary": "markdown for an on-call engineer",
  "digest": "max 5 plain-text lines for the Teams card",
  "failed_run_id": "...",
  "failed_job_ref": "...",
  "classification": "config | credentials | upstream_data | code | resources | transient | unknown",
  "confidence": "high | medium | low",
  "evidence": [{"source": "Get_Run_Logs <run id> line 38", "excerpt": "..."}],
  "proposed_fix": "...",
  "requires_human": true
}
```

Before this one, I didn't ask any Cortex agent for machine-readable output (`response_format` is not a feature of Cortex agents, yet). It works well, though (thanks, Sonnet 5), but I still made the caller (layer 4 below) still tolerate code fences or preambles by slicing from the first `{` to the last `}`.

The full agent spec ▾ 

```json
{
  "name": "JOB_INSPECTOR",
  "database_name": "META",
  "schema_name": "CORTEX",
  "comment": "Diagnoses a failed dltHub job run through read-only dltHub REST tools, classifies the failure and proposes a fix. Answers with one JSON object.",
  "profile": { "display_name": "dltHub Job-Inspector", "avatar": "SparklesAgentIcon" },
  "agent_spec": {
    "models": { "orchestration": "claude-sonnet-5" },
    "orchestration": { "budget": { "seconds": 600, "tokens": 400000 } },
    "instructions": {
      "orchestration": "IDENTITY: You are a job inspector for a dltHub Platform workspace. You run unattended, seconds after a job failed. An engineer reads your output only when the failure matters, so it must stand on its own.\n\nINPUT: the message names a failed run id and/or a job ref. Resolve in this order and stop at the first that works: (1) a run id: inspect that run, even if it turns out completed or running; (2) a job ref: take its latest failed run (List_Runs); (3) nothing usable: return status aborted, naming which input was missing. Never substitute an unrelated job.\n\nWORKFLOW (economical: aim for at most 8 tool calls): Get_Run first (status, trigger, profile, job_ref, timeout, pipeline summaries). Then Get_Run_Logs (start with max_lines 200; only ask for more if the earliest error is cut off). Work through the log earliest-error-first: logs cascade, the final traceback is usually a consequence of something further up. A traceback in workspace code (paths under the job's own run/ directory, e.g. run/<pipeline>.py) is the job's; a traceback in the runner or control plane after the job's work printed its completion is the platform's. CLASSIFY BY WHERE THE ERROR IS RAISED, NOT BY WHAT ITS MESSAGE CLAIMS: an exception raised in workspace code (the innermost frames under the job's run/ directory, e.g. a raise inside a dlt resource) is code, even when its message blames upstream data or an API. Classify upstream_data only when the log shows the upstream response itself (an HTTP status or body, an empty or malformed payload the code received and reported) and the code merely surfaced it. A message is a claim, not evidence. Check neighbours with List_Runs before calling anything transient: recurrence is the test. When you cite a neighbour as failing the same way, compare its actual error text (Get_Run_Logs on it): a neighbour that failed with a different error is not the same failure. Read Get_Job only when config looks suspect (profile, trigger, timeout, dependency groups). For a pipeline job, Get_Pipeline_Trace names the failed step (extract/normalize/load) and its exception.\n\nSCOPE: investigate the workspace, not the platform. If the run record, logs, job definition and pipeline trace don't explain the failure, say so and stop; that is a platform question for a human, not something to reverse-engineer.\n\nCONSTRAINTS: read-only by construction (your tools only read). Evidence or admit it: every classification must cite something you actually read; if you cannot find supporting output, confidence low and say in summary what you could not establish. Never invent a plausible cause. One run at a time: compare against neighbours when it helps, do not sweep the whole job history.\n\nCLASSIFICATION: config = missing or wrong setting, wrong profile, bad trigger or manifest. credentials = auth failure reaching a source or destination. upstream_data = the job ran correctly; the data it received was wrong, late or absent. code = an exception in workspace code; the traceback points into the pipeline or transformation. resources = out-of-memory kill, timeout, or disk exhaustion. transient = network blip, rate limit, or a platform-side failure the previous run did not have and the next likely will not; only after checking the neighbouring runs. unknown = no cause established; confidence must be low.\n\nCONFIDENCE: high = the earliest error names the cause directly and evidence quotes it. medium = the cause is inferred from surrounding evidence (neighbouring runs, job definition) and a plausible alternative remains. low = a guess or unknown.\n\nSTATUS: succeeded = you established what went wrong (a cause without a remedy still counts). failed = you read record, logs and job definition and still cannot say what went wrong: classification unknown, confidence low, summary says what you ruled out and where a human should start. aborted = the input did not identify a run.",
      "response": "Answer with EXACTLY ONE JSON object and nothing else: no prose before or after, no markdown code fence. Keys:\n- status: \"succeeded\" | \"failed\" | \"aborted\"\n- summary: markdown string for an on-call engineer: what failed, why, what to do; concise, bullets where they fit; say why the confidence follows from the evidence and what stayed unverified.\n- digest: plain-text string for a Teams card, at most 5 lines each starting with \"- \": what failed, resolution (already fixed, or what to do), watch out (only if there is a real caveat). Terse fragments, no headers, no code blocks. A compression of summary, never a different claim.\n- failed_run_id: string, the run id you actually inspected (empty if aborted)\n- failed_job_ref: string, the job ref of that run\n- classification: \"config\" | \"credentials\" | \"upstream_data\" | \"code\" | \"resources\" | \"transient\" | \"unknown\"\n- confidence: \"high\" | \"medium\" | \"low\"\n- evidence: array of {\"source\": string (e.g. \"Get_Run_Logs <run id> line 38\"), \"excerpt\": string (quoted log text)}, earliest error first; empty only when aborted\n- proposed_fix: string, what a human should do next (you never apply it)\n- requires_human: boolean, true when a person has to act before the job can succeed again\nDo not use em dashes.",
      "sample_questions": [
        { "question": "Inspect failed run <run id> of job jobs.__deployment__.<job>." },
        { "question": "Inspect the latest failed run of job jobs.__deployment__.<job>." }
      ]
    },
    "tools": [
      { "tool_spec": { "type": "generic", "name": "Get_Run",
          "description": "Reads one dltHub job run record: number, status, trigger, profile, start/end, duration, prev_run_id, job_ref, the job's timeout, and per-pipeline summaries (rows, status). Call first.",
          "input_schema": { "type": "object",
            "properties": { "run_id": { "description": "dltHub job run id (uuid)", "type": "string" } },
            "required": ["run_id"] } } },
      { "tool_spec": { "type": "generic", "name": "Get_Run_Logs",
          "description": "Reads the last N log lines of a dltHub job run as 'line_num [phase] content'. Setup-phase dependency-install lines are dropped when the program itself logged anything. Start with 200; raise only if the earliest error is cut off.",
          "input_schema": { "type": "object",
            "properties": {
              "run_id": { "description": "dltHub job run id (uuid)", "type": "string" },
              "max_lines": { "description": "how many trailing lines to return, 1-1000 (default 200)", "type": "integer" } },
            "required": ["run_id", "max_lines"] } } },
      { "tool_spec": { "type": "generic", "name": "List_Runs",
          "description": "Lists the most recent runs of one dltHub job, newest first (number, id, status, trigger, start/end). Use for the neighbour check and to find a job's latest failed run.",
          "input_schema": { "type": "object",
            "properties": {
              "job_ref": { "description": "dltHub job ref, e.g. jobs.__deployment__.my_job", "type": "string" },
              "max_runs": { "description": "how many runs, 1-20", "type": "integer" } },
            "required": ["job_ref", "max_runs"] } } },
      { "tool_spec": { "type": "generic", "name": "Get_Job",
          "description": "Reads one dltHub job definition: triggers, execute/require spec (timeout, instance, dependency groups), paused/archived, last run. Only when config looks suspect.",
          "input_schema": { "type": "object",
            "properties": { "job_ref": { "description": "dltHub job ref, e.g. jobs.__deployment__.my_job", "type": "string" } },
            "required": ["job_ref"] } } },
      { "tool_spec": { "type": "generic", "name": "Get_Pipeline_Trace",
          "description": "For a job run that ran dlt pipelines: per pipeline its status, failed step (extract/normalize/load), error message and each step's exception. Says so when the run had no pipeline.",
          "input_schema": { "type": "object",
            "properties": { "run_id": { "description": "dltHub job run id (uuid)", "type": "string" } },
            "required": ["run_id"] } } }
    ],
    "tool_resources": {
      "Get_Run":            { "type": "function", "identifier": "META.CORTEX.F_DLTHUB_GET_RUN",            "name": "F_DLTHUB_GET_RUN(VARCHAR)",                  "execution_environment": { "type": "warehouse", "warehouse": "<TOOL_WAREHOUSE>" } },
      "Get_Run_Logs":       { "type": "function", "identifier": "META.CORTEX.F_DLTHUB_GET_RUN_LOGS",       "name": "F_DLTHUB_GET_RUN_LOGS(VARCHAR, NUMBER)",     "execution_environment": { "type": "warehouse", "warehouse": "<TOOL_WAREHOUSE>" } },
      "List_Runs":          { "type": "function", "identifier": "META.CORTEX.F_DLTHUB_LIST_RUNS",          "name": "F_DLTHUB_LIST_RUNS(VARCHAR, NUMBER)",        "execution_environment": { "type": "warehouse", "warehouse": "<TOOL_WAREHOUSE>" } },
      "Get_Job":            { "type": "function", "identifier": "META.CORTEX.F_DLTHUB_GET_JOB",            "name": "F_DLTHUB_GET_JOB(VARCHAR)",                  "execution_environment": { "type": "warehouse", "warehouse": "<TOOL_WAREHOUSE>" } },
      "Get_Pipeline_Trace": { "type": "function", "identifier": "META.CORTEX.F_DLTHUB_GET_PIPELINE_TRACE", "name": "F_DLTHUB_GET_PIPELINE_TRACE(VARCHAR)",       "execution_environment": { "type": "warehouse", "warehouse": "<TOOL_WAREHOUSE>" } }
    }
  }
}
```

The spec lives as JSON in the repo, and the DDL is generated from it (`json.dumps(agent_spec)` between `$$` quotes):

```sql
create or replace agent META.CORTEX.JOB_INSPECTOR
  copy grants
  with profile = '{"display_name": "dltHub Job-Inspector", "avatar": "SparklesAgentIcon"}'
  from specification $$<json.dumps(agent_spec)>$$;
```

On grants: a key-pair JWT call always executes under the calling service user's **default role**. In my setup that role already owns the agent and the five functions, so nothing extra was required. A dedicated service role would need `usage` on the agent, the functions, the tool warehouse and the `SNOWFLAKE.CORTEX_USER` database role, all granted to its default role.

## Layer 4: the dltHub job that calls the agent

`job_inspector_cortex` is a plain dltHub batch job, roughly 120 lines in `__deployment__.py`, and most of it is error handling - no agent loop runs in it. dltHub's own `job-inspector` handles that on its side, via `@run.agent` and pydantic-ai talking to a configured model endpoint. This job only has to fire one HTTP request at Snowflake and handle the answer.

The `__deployment__.py` looks like this:

```python
from dlt._workspace.deployment.decorators import job
from dlt._workspace.deployment.trigger import job_fail

__all__ = [
    "some_pipeline",
    "another_pipeline",
    # ...
    "job_inspector",          # dltHub's own background agent
    "job_inspector_cortex",
]

_INSPECTOR_EXCLUDED = {"job_inspector", "job_inspector_cortex"}
_INSPECTOR_TRIGGERS = [job_fail(f"__deployment__.{n}") for n in __all__ if n not in _INSPECTOR_EXCLUDED]
```

A `job.fail` trigger names the failed **job** rather than the run, but the platform sets the inspector run's own `prev_run_id` to the run that fired it. Resolving that in code (before the agent is called) keeps the model from spending a tool call hunting for the right run:

```python
def _resolve_failed_run(run_context, failed_run_id: str, failed_job_ref: str) -> tuple:
    """(run id, job ref): explicit inputs, else this run's prev_run_id + the job ref from the trigger."""
    trigger = str((run_context or {}).get("trigger") or "")
    job_ref = failed_job_ref or (trigger.split(":", 1)[1] if trigger.startswith("job.fail:") else "")
    if failed_run_id:
        return failed_run_id, job_ref
    try:
        import os
        import dlthub_sdk
        from dlt._workspace._workspace_context import active

        cfg = active().runtime_config
        runtime = dlthub_sdk.connect(token=os.environ["RUNTIME__API_KEY"], base_url=cfg.api_base_url)
        own = runtime.workspaces.get(id=cfg.workspace_id).job_runs.get(id=run_context["run_id"]).to_dict()
        return str(own.get("prev_run_id") or ""), job_ref
    except Exception as err:
        print(f"Could not resolve prev_run_id, the agent will look up the job's latest failed run: {err}")
        return "", job_ref
```

`dlthub_sdk` ships with `dlt[hub]`, so it's already in every dltHub workspace. One minor surprise while writing this: a plain `@job` receives no platform API credential from the runner. Agent jobs get scoped creds through their `access=` declaration, but a plain job has no such parameter... So I had to add the viewer `RUNTIME__API_KEY` as a workspace ENV var in dltHub, which feels a little dog-hunting-its-own-tail-ish 😅

Snowflake authentication reuses the service user my dlt pipelines already sign in as, with the same private key already sitting in dlt's own secrets:

The dltHub job incl. JWT auth and agent call ▾ def \_snowflake\_jwt() -> str: import base64, hashlib, json, time import dlt from cryptography.hazmat.primitives import hashes, serialization from cryptography.hazmat.primitives.asymmetric import padding base = "destination.snowflake.credentials" key = serialization.load\_pem\_private\_key( dlt.secrets\[f"{base}.private\_key"\].encode(), password=dlt.secrets\[f"{base}.private\_key\_passphrase"\].encode(), ) spki = key.public\_key().public\_bytes(serialization.Encoding.DER, serialization.PublicFormat.SubjectPublicKeyInfo) qualified = f".{str(dlt.secrets\[f'{base}.username'\]).upper()}" now = int(time.time()) claims = { "iss": f"{qualified}.SHA256:{base64.b64encode(hashlib.sha256(spki).digest()).decode()}", "sub": qualified, "iat": now, "exp": now + 3600, } def b64(raw: bytes) -> bytes: return base64.urlsafe\_b64encode(raw).rstrip(b"=") signing\_input = b64(json.dumps({"alg": "RS256", "typ": "JWT"}).encode()) + b"." + b64( json.dumps(claims, separators=(",", ":")).encode() ) return (signing\_input + b"." + b64(key.sign(signing\_input, padding.PKCS1v15(), hashes.SHA256()))).decode() import requests \_CORTEX\_AGENT\_URL = ( "https://..snowflakecomputing.com" "/api/v2/databases/META/schemas/CORTEX/agents/JOB\_INSPECTOR:run" ) @job(trigger=\_INSPECTOR\_TRIGGERS, execute={"timeout": "10m", "concurrency": 4}) def job\_inspector\_cortex(run\_context=None, failed\_run\_id: str = "", failed\_job\_ref: str = ""): import json run\_id, job\_ref = \_resolve\_failed\_run(run\_context, failed\_run\_id, failed\_job\_ref) try: if not (run\_id or job\_ref): raise ValueError("no failed run id or job ref to inspect (not triggered by job.fail, no inputs)") ask = f"Inspect failed run {run\_id} of job {job\_ref}." if run\_id else f"Inspect the latest failed run of job {job\_ref}." with \_snowflake\_proxy(): # my Snowflake user sits behind a network policy; drop this without one resp = requests.post( \_CORTEX\_AGENT\_URL, headers={ "Authorization": f"Bearer {\_snowflake\_jwt()}", "X-Snowflake-Authorization-Type": "KEYPAIR\_JWT", "Content-Type": "application/json", "Accept": "application/json", }, json={"messages": \[{"role": "user", "content": \[{"type": "text", "text": ask}\]}\], "stream": False}, timeout=540, ) if resp.status\_code != 200: raise RuntimeError(f"Cortex agent HTTP {resp.status\_code}: {resp.text\[:300\]}") texts = \[c.get("text", "") for c in resp.json().get("content", \[\]) if c.get("type") == "text"\] raw = texts\[-1\] if texts else "" # the response instruction asks for one bare JSON object; tolerate a stray fence or preamble output = json.loads(raw\[raw.find("{"): raw.rfind("}") + 1\]) except Exception as err: output = \_fallback\_output(err, run\_id, job\_ref) \_post\_diagnosis\_card(output, job\_ref, "cortex diagnosis") return output 

A few choices: `stream: False`, because the agent endpoint streams server-sent events by default and a batch job wants only the final answer (the last text item in the response's `content` list, since earlier ones can be interim text wrapped around tool calls). `concurrency: 4`, because dltHub's default of 1 means an inspector triggered while a previous inspector is still running would simply skip the second inspector. And a 540-second HTTP timeout inside a 10-minute job, so the request gives up before the platform kills the job outright, leaving room for the fallback card to still go out.

Any failure along the way, nothing to inspect, auth, HTTP, timeout, an answer that won't parse, becomes a schema-valid stand-in so the card still goes out:

```python
def _fallback_output(err: Exception, failed_run_id: str, failed_job_ref: str) -> dict:
    reason = f"{type(err).__name__}: {err}"
    return {
        "status": "aborted",
        "summary": f"AI diagnosis unavailable: {reason}",
        "digest": f"- Diagnosis unavailable ({reason}): investigate manually.",
        "failed_run_id": failed_run_id,
        "failed_job_ref": failed_job_ref,
        "classification": "unknown",
        "confidence": "low",
        "evidence": [],
        "requires_human": True,
    }
```

The Teams card is a plain Adaptive Card posted through a workflow webhook. That part that was already covered:

[dltHub notifies Slack and email. But my org runs on Teams 😩dltHub ships Slack and email alerts out of the box. My org runs on Teams, so I built the third one myself and wired it into all \~70 pipelines with one decorator shadow.![](https://storage.ghost.io/c/7d/94/7d942fe1-7868-4a1a-b2c9-4eb415b1a546/content/images/icon/headshot_ring-a67a82d7-5c79-4deb-84cb-ad5ca0a3fe22.png)sfrt.ioMartin Seifert![](https://storage.ghost.io/c/7d/94/7d942fe1-7868-4a1a-b2c9-4eb415b1a546/content/images/thumbnail/warning-ca7738c9-f9aa-4cee-969b-e2a2dd55ef93.png)](https://www.sfrt.io/dlthub-notifies-slack-and-email-but-my-org-runs-on-teams/)

## Does it actually work?

None of this replaces dltHub's own `job-inspector`. It runs next to it for now, on infrastructure that was already there for other reasons, and it exists only because dltHub wrote the rubric, shipped the SDK and opened the API in the first place.

With both inspectors running, both cards usually land within about a minute:

![](https://storage.ghost.io/c/7d/94/7d942fe1-7868-4a1a-b2c9-4eb415b1a546/content/images/2026/09/image.png)

Most of the times, they agree, but sometimes they arrive at different verdicts (like in the example above). Hence, for now, a human (me) stays in the loop. Whichever of the two stays long-term, I don't know... I just had a few days to test and not enough real pipeline failures (actually, only one expired credential case so far, both versions identified it correctly). But I can already say, both beat what was here before: nothing. 😎

Next up: auto-healing... Stay tuned 😉

---

#### License

This is free and unencumbered software released into the public domain. Anyone is free to copy, modify, publish, use, compile, sell, or distribute this software, either in source code form or as a compiled binary, for any purpose, commercial or non-commercial, and by any means.  
In jurisdictions that recognize copyright laws, the author or authors of this software dedicate any and all copyright interest in the software to the public domain. We make this dedication for the benefit of the public at large and to the detriment of our heirs and successors. We intend this dedication to be an overt act of relinquishment in perpetuity of all present and future rights to this software under copyright law.

THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.

For more information, please refer to [https://unlicense.org/](https://unlicense.org/?ref=sfrt.io)