Pipelines
Run a client's Airbyte syncs as one daily unit and get a signed webhook when all of that client's data has landed.
Overview
By default every Airbyte connection runs on its own schedule. That works until something has to run after the data lands: a transform, a report, an export. Twelve connections on twelve crons never produce a moment that means "today's data is here", so a consumer either guesses a time or polls.
A pipeline makes one client's syncs a single daily run:
- You choose which of the client's connections the run waits for, and the UTC hour it starts.
- Mythic sets those connections to
manualin Airbyte and triggers them itself at that hour. - Every five minutes Mythic reads each job's state from Airbyte. When every connection that could sync has synced, the run succeeds.
- Mythic POSTs one signed completion webhook to your URL.
A pipeline belongs to one client. Each client can have one pipeline, and each pipeline runs at most once per UTC day.
A pipeline is created for you: the client's first connection (OAuth or direct, any platform) creates it enabled with run_hour_utc 2, unless that connection sets an explicit sync_frequency. Every later connection created without an explicit sync_frequency joins it on manual. An explicit sync_frequency opts out — no pipeline is created for it and nothing is enrolled — and deleting a connection removes it from the scope. Everything else on this page is for reading runs and changing the setup afterwards.
All paths below are relative to https://mythic-analytics.gulp.workers.dev/client/v1/airbyte. Reads accept an agency key (ak_) or a location secret key (sk_). Every write needs an agency key. The MCP server exposes the same operations; the tool for each step is named below.
Quickstart
List the client's connections
GET /pipelines/{clientId}/candidates returns the client's connections with in_scope (already in the pipeline) and suggested (currently active). MCP: list_airbyte_pipeline_candidates.
curl https://mythic-analytics.gulp.workers.dev/client/v1/airbyte/pipelines/538d4ced-447c-40b9-9afa-d220f9afcc84/candidates \
-H "Authorization: Bearer ak_..."
{
"success": true,
"data": [
{
"id": "c44b56c8-35b8-44b6-b8b8-82167f54f697",
"connection_id": "4d874ed9-c54d-4b02-84c6-fcd34857a60e",
"platform": "meta_ads",
"display_name": "Hello Conversions - Meta Ads",
"status": "active",
"last_sync_at": "2026-09-30T07:41:12Z",
"sync_frequency": "24h",
"managed_by_pipeline": false,
"in_scope": false,
"suggested": true
}
]
}
Use suggested to pre-select, and have a person confirm the set. Which connections belong to a client can't be decided reliably from an Airbyte workspace alone.
Create the pipeline
PUT /pipelines/{clientId} with the connection ids (Mythic's id, not Airbyte's connection_id). MCP: set_airbyte_pipeline.
curl -X PUT https://mythic-analytics.gulp.workers.dev/client/v1/airbyte/pipelines/538d4ced-447c-40b9-9afa-d220f9afcc84 \
-H "Authorization: Bearer ak_..." -H "Content-Type: application/json" \
-d '{
"enabled": true,
"connection_ids": ["c44b56c8-35b8-44b6-b8b8-82167f54f697"],
"run_hour_utc": 8,
"webhook_url": "https://reports.example.com/hooks/mythic",
"webhook_secret": "generate-a-long-random-string"
}'
{
"success": true,
"data": {
"id": "5d0c2f1e-8a41-4a8e-9c1b-2f7e1d3a9b60",
"client_id": "538d4ced-447c-40b9-9afa-d220f9afcc84",
"enabled": true,
"connection_ids": ["c44b56c8-35b8-44b6-b8b8-82167f54f697"],
"run_hour_utc": 8,
"webhook_url": "https://reports.example.com/hooks/mythic",
"has_webhook_secret": true,
"schedule_changes": {
"taken_over": [
{ "id": "c44b56c8-35b8-44b6-b8b8-82167f54f697", "taken": true,
"stashed": { "scheduleType": "cron", "cronExpression": "0 0 2 * * ? UTC" } }
],
"handed_back": [],
"errors": []
}
}
}
stashed is the schedule the connection had. Mythic gives it back when the connection leaves the pipeline or the pipeline is disabled.
Start a run now (optional)
Runs start on their own at run_hour_utc. To start one immediately, POST /pipelines/{clientId}/runs. MCP: start_airbyte_pipeline_run.
curl -X POST https://mythic-analytics.gulp.workers.dev/client/v1/airbyte/pipelines/538d4ced-447c-40b9-9afa-d220f9afcc84/runs \
-H "Authorization: Bearer ak_..."
Check progress
GET /pipelines/{clientId}/runs lists runs newest first. GET /pipelines/{clientId}/runs/{runId} returns one. MCP: list_airbyte_pipeline_runs, get_airbyte_pipeline_run. The run's shape is the same as the webhook's data below.
Schedules
A pipeline owns its connections' schedules. While it is enabled, every in-scope connection is manual in Airbyte and only Mythic triggers it. If someone adds a cron back in the Airbyte UI, the connection syncs twice a day. Mythic resets drifted schedules to manual and logs the drift, but change schedules through this API rather than the Airbyte UI.
- Taking over is all or nothing. Mythic changes Airbyte first and saves the pipeline second. If Airbyte is unreachable (
502 airbyte_unavailable) or any connection can't be set tomanual(502 schedule_takeover_failed, reasons indetails.errors), nothing is saved and every connection changed during the request gets its schedule back. - Giving back is reported, not hidden. Removing a connection from
connection_ids, or sending"enabled": false, restores the stashed schedule. If Airbyte refuses, the response still succeeds but lists the connection inschedule_changes.errorsand adds a top-levelwarning. That connection is onmanualwith nothing triggering it until you retry the samePUT. - Omitted fields keep their stored value. A new pipeline starts disabled with
run_hour_utc2. - Deleting a connection removes it from the pipeline. The delete response includes
removed_from_pipeline. - New connections join automatically. A connection created without an explicit
sync_frequencyfor a client that has a pipeline is createdmanualand added to it.
Change the hour, drop a connection, or remove the webhook secret:
curl -X PUT .../pipelines/538d4ced-447c-40b9-9afa-d220f9afcc84 \
-H "Authorization: Bearer ak_..." -H "Content-Type: application/json" \
-d '{"run_hour_utc": 6, "connection_ids": ["c44b56c8-35b8-44b6-b8b8-82167f54f697"], "webhook_secret": null}'
Turn the pipeline off and give every connection its schedule back:
curl -X PUT .../pipelines/538d4ced-447c-40b9-9afa-d220f9afcc84 \
-H "Authorization: Bearer ak_..." -H "Content-Type: application/json" \
-d '{"enabled": false}'
Run lifecycle
| Status | Meaning |
|---|---|
running | Syncs are in flight or waiting to be triggered. |
held | At least one sync failed after its retries. The run keeps retrying every five minutes and completes on its own once the sync passes. |
succeeded | Every connection that could sync did. The completion webhook fires. |
failed | The run wedged: nothing in flight and no progress for 36 hours. |
skipped | Superseded by a manual run. |
A run decides it is done with these rules, in order:
- Skipped connections never block. A connection disabled in Airbyte, inactive in Mythic, or not due today (a weekly
sync_dow) is skipped with areason. - A first-time backfill doesn't block, because it can run for days. A backfill that fails does block.
- Anything still running → wait.
- Anything failed → hold. Each failed sync is retried twice before the run waits for a person.
- Otherwise → succeeded, and the webhook fires.
Handling a held run
A held run means data you expected is missing. Fix the cause (a revoked ad account, an expired credential), then choose one:
| Action | Request | MCP | Effect |
|---|---|---|---|
| Retry failed syncs | POST /pipelines/{clientId}/runs/{runId}/retry | airbyte_pipeline_run_action with retry | Resets the failed steps and reopens the run. Works on held and failed runs only, and not while another run of the pipeline is in flight. |
| Start over | POST /pipelines/{clientId}/runs | start_airbyte_pipeline_run | Today's held run restarts in place with every connection re-triggered (restarted: true). A held run from an earlier day is marked skipped and a new run starts. |
| Complete anyway | POST /pipelines/{clientId}/runs/{runId}/force | airbyte_pipeline_run_action with force | Closes a held run as succeeded and sends the webhook with forced: true. |
Forcing is a per-run decision, never a setting. A forced run carries forced: true in its webhook permanently, because anything built on it was built on incomplete data.
The completion webhook
When a run succeeds, Mythic POSTs to webhook_url:
POST /hooks/mythic HTTP/1.1
Content-Type: application/json
X-Mythic-Run-Id: 7b1f0c52-3d8e-4a19-b6f2-0e9d4c1a8f37
X-Mythic-Signature: sha256=5f2b...c91e
{
"type": "airbyte.pipeline.completed",
"delivered_at": "2026-10-01T08:47:03.112Z",
"data": {
"run_id": "7b1f0c52-3d8e-4a19-b6f2-0e9d4c1a8f37",
"client_id": "538d4ced-447c-40b9-9afa-d220f9afcc84",
"run_date": "2026-10-01",
"status": "succeeded",
"complete": true,
"forced": false,
"held_reason": null,
"triggered_by": "schedule",
"started_at": "2026-10-01T08:00:41.905Z",
"finished_at": "2026-10-01T08:45:00.000Z",
"totals": { "total": 2, "succeeded": 1, "failed": 0, "skipped": 1, "records_committed": 18432 },
"connections": [
{
"connection_id": "4d874ed9-c54d-4b02-84c6-fcd34857a60e",
"name": "Hello Conversions - Meta Ads",
"platform": "meta_ads",
"status": "succeeded",
"reason": null,
"job_id": 108421337,
"records": 18432,
"bytes": 40211876,
"duration_seconds": 2659,
"first_time": false,
"last_success_at": "2026-10-01T08:45:00.000Z"
},
{
"connection_id": "9a3e61b0-51f4-4c8e-a0c2-6d7b18e2f4d9",
"name": "Hello Conversions - Google Ads",
"platform": "google_ads",
"status": "skipped",
"reason": "disabled in Airbyte (inactive)",
"job_id": null,
"records": null,
"bytes": null,
"duration_seconds": null,
"first_time": false,
"last_success_at": "2026-09-22T02:14:09.000Z"
}
]
}
}
- Check
complete, not onlystatus.complete: falsewithstatus: "succeeded"means some connections were skipped. Read eachconnections[].reasonto see whether that matters. - Check
last_success_aton skipped connections. The example's Google Ads data is nine days old even though the run succeeded. - Delivery is at least once. Deduplicate on
X-Mythic-Run-Id, and return a 2xx for a duplicate. - Respond within 10 seconds. Queue the work and return. A slow response counts as a failure and is retried.
- Failed deliveries are retried on later ticks, up to five attempts. After that, re-send with
POST /pipelines/{clientId}/runs/{runId}/redeliver(MCP:airbyte_pipeline_run_actionwithredeliver). Only asucceededrun can be re-sent.
Verifying the signature
X-Mythic-Signature is sha256= followed by the hex HMAC-SHA256 of the raw request body using your webhook_secret. Verify it before parsing the JSON.
import crypto from 'node:crypto';
import express from 'express';
const app = express();
const seen = new Set(); // use your database in production
app.post('/hooks/mythic', express.raw({ type: 'application/json' }), (req, res) => {
const expected = 'sha256=' + crypto
.createHmac('sha256', process.env.MYTHIC_WEBHOOK_SECRET)
.update(req.body)
.digest('hex');
const got = req.get('X-Mythic-Signature') || '';
if (got.length !== expected.length
|| !crypto.timingSafeEqual(Buffer.from(got), Buffer.from(expected))) {
return res.status(401).end();
}
const runId = req.get('X-Mythic-Run-Id');
if (seen.has(runId)) return res.status(200).end(); // duplicate delivery
seen.add(runId);
const { data } = JSON.parse(req.body);
queueReportBuild(data); // do the work after responding
res.status(202).end();
});
import hashlib, hmac, json, os
from flask import Flask, request
app = Flask(__name__)
seen = set() # use your database in production
@app.post("/hooks/mythic")
def mythic_hook():
raw = request.get_data()
expected = "sha256=" + hmac.new(
os.environ["MYTHIC_WEBHOOK_SECRET"].encode(), raw, hashlib.sha256
).hexdigest()
if not hmac.compare_digest(request.headers.get("X-Mythic-Signature", ""), expected):
return "", 401
run_id = request.headers["X-Mythic-Run-Id"]
if run_id in seen:
return "", 200 # duplicate delivery
seen.add(run_id)
queue_report_build(json.loads(raw)["data"]) # do the work after responding
return "", 202
Timing
- Runs are advanced every five minutes, so completion is detected up to about five minutes after the last sync finishes.
- If Mythic misses
run_hour_utc, the run starts as soon as it can. A run missed late in the UTC day starts in the first hours of the next day, dated for the day it was meant for. - One run per client per UTC day. After a run finishes, the next one is tomorrow's scheduled run. To re-sync a single connection sooner, use
POST /connections/{id}/sync.
Errors
Errors use the standard envelope. code is stable; details carries the conflicting run_id or per-connection errors where relevant.
{
"success": false,
"error": "A run already exists for 2026-10-01 (succeeded). Runs are one per client per UTC day; the next one starts tomorrow.",
"code": "run_exists_today",
"details": { "run_id": "7b1f0c52-3d8e-4a19-b6f2-0e9d4c1a8f37", "status": "succeeded" }
}
| Status | Code | Meaning |
|---|---|---|
400 | invalid_body | A field has the wrong type or format. |
400 | unknown_connections | A connection_ids entry is not this client's connection. |
400 | empty_scope | The pipeline has no connections. |
400 | not_retryable | Retry needs a held or failed run. |
400 | not_held | Force needs a held run. |
400 | not_completed / no_webhook | Redeliver needs a succeeded run and a webhook_url. |
403 | Writes need an agency key (ak_). | |
404 | no_pipeline / run_not_found / client_not_found | Nothing by that id for this agency. |
409 | run_in_flight | A run is running. Starting or retrying another would sync the same connections twice. |
409 | run_exists_today | Today already has a finished run. |
409 | run_changed | The five-minute tick updated the run during your request. Read it again and retry. |
500 | save_failed | The pipeline could not be saved. Schedules changed during the request were given back. |
502 | airbyte_unavailable / schedule_takeover_failed | Airbyte could not be changed. Nothing was saved. |