airbyte_ops_mcp.mcp.prod_db_ops
MCP tools for querying the Airbyte Cloud Prod DB Replica.
This module provides MCP tools that wrap the query functions from airbyte_ops_mcp.prod_db_access.queries for use by AI agents.
MCP reference
MCP primitives registered by the prod_db_ops module of the airbyte-internal-ops server: 21 tool(s), 0 prompt(s), 0 resource(s).
Tools (21)
query_connector_pin_stats
Query connector versions that have at least one scoped configuration pin.
Returns versions from the prod DB that are referenced by at least one
scoped_configuration pin (key = 'connector_version'). Each version
appears exactly once with per-scope pin breakdown (actor, workspace, org).
If neither filter is provided, returns the global superset across all connectors.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
connector_definition_id |
string | null |
no | null |
Connector definition UUID to filter by (optional). Mutually exclusive with connector_canonical_name. |
connector_canonical_name |
string | null |
no | null |
Connector canonical name (e.g. source-postgres) to filter by. Resolved to a definition ID via the registry. Mutually exclusive with connector_definition_id. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"connector_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector definition UUID to filter by (optional). Mutually exclusive with `connector_canonical_name`."
},
"connector_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector canonical name (e.g. `source-postgres`) to filter by. Resolved to a definition ID via the registry. Mutually exclusive with `connector_definition_id`."
}
},
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"description": "A connector version that has at least one scoped configuration pin.",
"properties": {
"version_id": {
"description": "The actor_definition_version UUID",
"type": "string"
},
"connector_definition_id": {
"description": "The connector definition UUID",
"type": "string"
},
"connector_name": {
"description": "Human-readable connector name",
"type": "string"
},
"docker_repository": {
"description": "Docker repository path",
"type": "string"
},
"docker_image_tag": {
"description": "Docker image tag for this version",
"type": "string"
},
"last_published": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "ISO timestamp when this version was last published"
},
"pin_count": {
"description": "Total number of scoped_configuration rows pinning to this version",
"type": "integer"
},
"breaking_change_pins": {
"default": 0,
"description": "Number of actor-scoped pins created by breaking changes",
"type": "integer"
},
"rollout_pins": {
"default": 0,
"description": "Number of pins created by connector rollouts",
"type": "integer"
},
"actor_pins": {
"description": "Number of actor-scoped pins (excludes breaking change and rollout pins)",
"type": "integer"
},
"workspace_pins": {
"description": "Number of workspace-scoped pins",
"type": "integer"
},
"org_pins": {
"description": "Number of organization-scoped pins",
"type": "integer"
}
},
"required": [
"version_id",
"connector_definition_id",
"connector_name",
"docker_repository",
"docker_image_tag",
"pin_count",
"actor_pins",
"workspace_pins",
"org_pins"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_connector_population_summary
Summarize the applied vs potential pinning audience for a connector, by tier.
Answers "how many actors are pinned and how many are eligible for pinning, split by tier" — the population view analogous to a rollout's audience.
active: the potential audience — non-tombstoned actors of the definition that have at least one active connection (connection.status = 'active'). Inactive/disabled and deprecated connections are excluded, so this reflects the enabled, rollout-touchable population rather than every actor ever created.pinned_any: active actors that already have an effectiveconnector_versionpin at any scope (actor/workspace/org).eligible:activeminuspinned_any— actors available to pin.pinned_to_version: the applied audience for the requested version (only when a version identifier was provided).
Backed by scoped_configuration + actor/connection tables, so it is cheap
to compute. Accepts a version identifier (preferred, adds
pinned_to_version) or a definition-level identifier.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
connector_version_id |
string | null |
no | null |
Connector version UUID. When provided, the applied audience (pinned_to_version) is included. Provide this OR connector_name + connector_version OR a definition-level identifier. |
connector_name |
string | null |
no | null |
Canonical connector name (e.g. source-postgres). Used with connector_version to resolve the version UUID. |
connector_version |
string | null |
no | null |
Semver version tag (e.g. 0.3.59). Used with connector_name. |
connector_definition_id |
string | null |
no | null |
Connector definition UUID for a definition-level summary (no pinned_to_version breakdown). |
connector_canonical_name |
string | null |
no | null |
Canonical connector name resolved to a definition ID via the registry, for a definition-level summary. |
customer_tier_filter |
enum("TIER_0", "TIER_1", "TIER_2", "UNKNOWN", "ALL") |
no | "TIER_2" |
Which customer tiers to count. Defaults to TIER_2; pass ALL to include TIER_0/TIER_1 (revenue-critical) customers. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"connector_version_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector version UUID. When provided, the applied audience (`pinned_to_version`) is included. Provide this OR `connector_name` + `connector_version` OR a definition-level identifier."
},
"connector_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical connector name (e.g. `source-postgres`). Used with `connector_version` to resolve the version UUID."
},
"connector_version": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Semver version tag (e.g. `0.3.59`). Used with `connector_name`."
},
"connector_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector definition UUID for a definition-level summary (no `pinned_to_version` breakdown)."
},
"connector_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical connector name resolved to a definition ID via the registry, for a definition-level summary."
},
"customer_tier_filter": {
"default": "TIER_2",
"description": "Which customer tiers to count. Defaults to `TIER_2`; pass `ALL` to include TIER_0/TIER_1 (revenue-critical) customers.",
"enum": [
"TIER_0",
"TIER_1",
"TIER_2",
"UNKNOWN",
"ALL"
],
"type": "string"
}
},
"type": "object"
}
Show output JSON schema
{
"description": "Applied vs potential pinning audience for a connector, split by tier.",
"properties": {
"connector_definition_id": {
"description": "The connector definition UUID",
"type": "string"
},
"connector_version_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The version UUID, when a specific version was requested"
},
"docker_repository": {
"description": "Docker repository path",
"type": "string"
},
"docker_image_tag": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Docker image tag, when a version was requested"
},
"customer_tier_filter": {
"description": "Tier filter applied to the counts (`TIER_0`/`TIER_1`/`TIER_2`/`UNKNOWN`/`ALL`)",
"type": "string"
},
"active": {
"description": "Potential audience: enabled actors of the definition (those with at least one active connection, `status = 'active'`), by tier",
"properties": {
"tier_0_count": {
"default": 0,
"description": "Number of TIER_0 entries",
"type": "integer"
},
"tier_1_count": {
"default": 0,
"description": "Number of TIER_1 entries",
"type": "integer"
},
"tier_2_count": {
"default": 0,
"description": "Number of TIER_2 entries",
"type": "integer"
},
"unknown_count": {
"default": 0,
"description": "Number of UNKNOWN entries",
"type": "integer"
},
"total": {
"default": 0,
"description": "Total number of entries",
"type": "integer"
},
"warnings": {
"description": "Warnings raised while building this summary.",
"items": {
"type": "string"
},
"type": "array"
}
},
"type": "object"
},
"pinned_any": {
"description": "Active actors already pinned to any version, by tier",
"properties": {
"tier_0_count": {
"default": 0,
"description": "Number of TIER_0 entries",
"type": "integer"
},
"tier_1_count": {
"default": 0,
"description": "Number of TIER_1 entries",
"type": "integer"
},
"tier_2_count": {
"default": 0,
"description": "Number of TIER_2 entries",
"type": "integer"
},
"unknown_count": {
"default": 0,
"description": "Number of UNKNOWN entries",
"type": "integer"
},
"total": {
"default": 0,
"description": "Total number of entries",
"type": "integer"
},
"warnings": {
"description": "Warnings raised while building this summary.",
"items": {
"type": "string"
},
"type": "array"
}
},
"type": "object"
},
"eligible": {
"description": "Active actors not pinned to any version (available to pin), by tier",
"properties": {
"tier_0_count": {
"default": 0,
"description": "Number of TIER_0 entries",
"type": "integer"
},
"tier_1_count": {
"default": 0,
"description": "Number of TIER_1 entries",
"type": "integer"
},
"tier_2_count": {
"default": 0,
"description": "Number of TIER_2 entries",
"type": "integer"
},
"unknown_count": {
"default": 0,
"description": "Number of UNKNOWN entries",
"type": "integer"
},
"total": {
"default": 0,
"description": "Total number of entries",
"type": "integer"
},
"warnings": {
"description": "Warnings raised while building this summary.",
"items": {
"type": "string"
},
"type": "array"
}
},
"type": "object"
},
"pinned_to_version": {
"anyOf": [
{
"description": "Summary of tier distribution across a set of results.",
"properties": {
"tier_0_count": {
"default": 0,
"description": "Number of TIER_0 entries",
"type": "integer"
},
"tier_1_count": {
"default": 0,
"description": "Number of TIER_1 entries",
"type": "integer"
},
"tier_2_count": {
"default": 0,
"description": "Number of TIER_2 entries",
"type": "integer"
},
"unknown_count": {
"default": 0,
"description": "Number of UNKNOWN entries",
"type": "integer"
},
"total": {
"default": 0,
"description": "Total number of entries",
"type": "integer"
},
"warnings": {
"description": "Warnings raised while building this summary.",
"items": {
"type": "string"
},
"type": "array"
}
},
"type": "object"
},
{
"type": "null"
}
],
"default": null,
"description": "Applied audience: actors pinned to the requested version, by tier. `None` when no specific version was requested."
}
},
"required": [
"connector_definition_id",
"docker_repository",
"customer_tier_filter",
"active",
"pinned_any",
"eligible"
],
"type": "object"
}
query_connector_version_health_summary
Summarize actor health for a connector version into four buckets.
Answers "how many actors on a version are healthy / unhealthy / awaiting / disabled". Classification per actor over the lookback window:
healthy: at least one successful sync (the same success signal the autopilot health gate uses).unhealthy: failures and no successes.awaiting: ran but produced only non-terminal jobs (no result yet).disabled: (wheninclude_pinned_disabled) pinned to the version with no jobs at all in the window — the dormant/inactive audience.
Built on the attempt/version primitive — the version stamped into
jobs.config at job-creation time — not the current pin state, so it
reflects actors that actually ran this version. This scans jobs over the
window and is more expensive than the population summary; keep days
bounded and query per-version.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
connector_version_id |
string | null |
no | null |
Connector version UUID. Provide this OR connector_name + connector_version. |
connector_name |
string | null |
no | null |
Canonical connector name (e.g. source-postgres). Used with connector_version to resolve the version UUID. |
connector_version |
string | null |
no | null |
Semver version tag (e.g. 0.3.59). Used with connector_name. |
days |
integer |
no | 7 |
Number of days to look back (default: 7, max: 30) |
include_pinned_disabled |
boolean |
no | true |
If True (default), actors pinned to the version that ran no jobs in the window are counted as disabled. |
customer_tier_filter |
enum("TIER_0", "TIER_1", "TIER_2", "UNKNOWN", "ALL") |
no | "TIER_2" |
Which customer tiers to count. Defaults to TIER_2; pass ALL to include TIER_0/TIER_1 (revenue-critical) customers. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"connector_version_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector version UUID. Provide this OR `connector_name` + `connector_version`."
},
"connector_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical connector name (e.g. `source-postgres`). Used with `connector_version` to resolve the version UUID."
},
"connector_version": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Semver version tag (e.g. `0.3.59`). Used with `connector_name`."
},
"days": {
"default": 7,
"description": "Number of days to look back (default: 7, max: 30)",
"maximum": 30,
"minimum": 1,
"type": "integer"
},
"include_pinned_disabled": {
"default": true,
"description": "If `True` (default), actors pinned to the version that ran no jobs in the window are counted as `disabled`.",
"type": "boolean"
},
"customer_tier_filter": {
"default": "TIER_2",
"description": "Which customer tiers to count. Defaults to `TIER_2`; pass `ALL` to include TIER_0/TIER_1 (revenue-critical) customers.",
"enum": [
"TIER_0",
"TIER_1",
"TIER_2",
"UNKNOWN",
"ALL"
],
"type": "string"
}
},
"type": "object"
}
Show output JSON schema
{
"description": "Four-bucket health rollup for the actors running a connector version.",
"properties": {
"connector_version_id": {
"description": "The connector version UUID",
"type": "string"
},
"connector_definition_id": {
"description": "The connector definition UUID",
"type": "string"
},
"docker_repository": {
"description": "Docker repository path",
"type": "string"
},
"docker_image_tag": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Docker image tag for this version"
},
"days": {
"description": "Lookback window in days",
"type": "integer"
},
"customer_tier_filter": {
"description": "Tier filter applied to the counts (`TIER_0`/`TIER_1`/`TIER_2`/`UNKNOWN`/`ALL`)",
"type": "string"
},
"healthy": {
"description": "Actors with at least one successful sync",
"type": "integer"
},
"unhealthy": {
"description": "Actors with failures and no successes in the window",
"type": "integer"
},
"awaiting": {
"description": "Actors that ran but produced only non-terminal jobs (no result yet)",
"type": "integer"
},
"disabled": {
"description": "Actors pinned to the version that produced no jobs in the window \u2014 the dormant/inactive audience",
"type": "integer"
},
"total_actors": {
"description": "Total actors counted across all states",
"type": "integer"
},
"healthy_by_tier": {
"description": "Healthy actors by tier",
"properties": {
"tier_0_count": {
"default": 0,
"description": "Number of TIER_0 entries",
"type": "integer"
},
"tier_1_count": {
"default": 0,
"description": "Number of TIER_1 entries",
"type": "integer"
},
"tier_2_count": {
"default": 0,
"description": "Number of TIER_2 entries",
"type": "integer"
},
"unknown_count": {
"default": 0,
"description": "Number of UNKNOWN entries",
"type": "integer"
},
"total": {
"default": 0,
"description": "Total number of entries",
"type": "integer"
},
"warnings": {
"description": "Warnings raised while building this summary.",
"items": {
"type": "string"
},
"type": "array"
}
},
"type": "object"
},
"unhealthy_by_tier": {
"description": "Unhealthy actors by tier",
"properties": {
"tier_0_count": {
"default": 0,
"description": "Number of TIER_0 entries",
"type": "integer"
},
"tier_1_count": {
"default": 0,
"description": "Number of TIER_1 entries",
"type": "integer"
},
"tier_2_count": {
"default": 0,
"description": "Number of TIER_2 entries",
"type": "integer"
},
"unknown_count": {
"default": 0,
"description": "Number of UNKNOWN entries",
"type": "integer"
},
"total": {
"default": 0,
"description": "Total number of entries",
"type": "integer"
},
"warnings": {
"description": "Warnings raised while building this summary.",
"items": {
"type": "string"
},
"type": "array"
}
},
"type": "object"
},
"awaiting_by_tier": {
"description": "Awaiting-results actors by tier",
"properties": {
"tier_0_count": {
"default": 0,
"description": "Number of TIER_0 entries",
"type": "integer"
},
"tier_1_count": {
"default": 0,
"description": "Number of TIER_1 entries",
"type": "integer"
},
"tier_2_count": {
"default": 0,
"description": "Number of TIER_2 entries",
"type": "integer"
},
"unknown_count": {
"default": 0,
"description": "Number of UNKNOWN entries",
"type": "integer"
},
"total": {
"default": 0,
"description": "Total number of entries",
"type": "integer"
},
"warnings": {
"description": "Warnings raised while building this summary.",
"items": {
"type": "string"
},
"type": "array"
}
},
"type": "object"
},
"disabled_by_tier": {
"description": "Disabled actors by tier",
"properties": {
"tier_0_count": {
"default": 0,
"description": "Number of TIER_0 entries",
"type": "integer"
},
"tier_1_count": {
"default": 0,
"description": "Number of TIER_1 entries",
"type": "integer"
},
"tier_2_count": {
"default": 0,
"description": "Number of TIER_2 entries",
"type": "integer"
},
"unknown_count": {
"default": 0,
"description": "Number of UNKNOWN entries",
"type": "integer"
},
"total": {
"default": 0,
"description": "Total number of entries",
"type": "integer"
},
"warnings": {
"description": "Warnings raised while building this summary.",
"items": {
"type": "string"
},
"type": "array"
}
},
"type": "object"
}
},
"required": [
"connector_version_id",
"connector_definition_id",
"docker_repository",
"days",
"customer_tier_filter",
"healthy",
"unhealthy",
"awaiting",
"disabled",
"total_actors",
"healthy_by_tier",
"unhealthy_by_tier",
"awaiting_by_tier",
"disabled_by_tier"
],
"type": "object"
}
query_prod_actors_by_pinned_connector_version
List actors (sources/destinations) effectively pinned to a specific connector version.
Returns all actors that are effectively pinned to a specific connector version, considering all scope levels: actor-level pins, workspace-level pins, and organization-level pins (with actor > workspace > organization precedence). Useful for monitoring rollouts and understanding which customers are affected.
The actor_id field is the actor ID (superset of source_id/destination_id).
Returns list of dicts with keys: actor_id, connector_definition_id, origin_type, origin, description, created_at, expires_at, pin_scope_type, actor_name, workspace_id, workspace_name, organization_id, dataplane_group_id, dataplane_name
pin_scope_type is 'actor', 'workspace', or 'organization' indicating which scope level the effective pin came from.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
connector_version_id |
string |
yes | — | Connector version UUID to find pinned instances for |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"connector_version_id": {
"description": "Connector version UUID to find pinned instances for",
"type": "string"
}
},
"required": [
"connector_version_id"
],
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_connection_sync_activity
List recent sync jobs and attempts from the Prod DB Replica.
Returns one row per (job, attempt) pair for sync jobs whose updated_at
falls in [start_at, end_at), scoped to the provided organization,
workspace, or connection IDs. Designed for live operational lookups —
e.g. "what happened on this connection in the last hour" — not for
historical analysis.
Each row is enriched with customer_tier and is_eu for the owning
organization. Tier filtering is intentionally not applied — this is a
read-only observability query.
Input requirements:
- At least one of
organization_id,workspace_id, orconnection_idsmust be provided (any combination is accepted). start_atandend_atmust be timezone-aware andstart_at < end_at.
Key fields in each row:
job_id,attempt_id,attempt_numberjob_status,attempt_statusjob_started_at,job_updated_at,attempt_ended_atfailure_summary(JSON; populated when an attempt failed)connection_id,connection_name,connection_statussource_actor_id,source_actor_name,source_actor_definition_iddestination_actor_id,destination_actor_name,destination_actor_definition_idworkspace_id,workspace_name,organization_iddataplane_group_id,dataplane_namecustomer_tier,is_eu(added by tier enrichment)
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
start_at |
string |
yes | — | Inclusive start timestamp for the sync activity window. Must be timezone-aware (ISO 8601 with offset or Z). |
end_at |
string |
yes | — | Exclusive end timestamp for the sync activity window. Must be timezone-aware and strictly after start_at. |
organization_id |
string | enum("664c690e-5263-49ba-b01f-4a6759b3330a") | null |
no | null |
Optional organization UUID or alias. At least one of organization_id, workspace_id, or connection_ids is required. Accepts @airbyte-internal as an alias for the Airbyte internal org. |
workspace_id |
string | enum("266ebdfe-0d7b-4540-9817-de7e4505ba61") | null |
no | null |
Optional workspace UUID or alias. At least one of organization_id, workspace_id, or connection_ids is required. Accepts @devin-ai-sandbox as an alias for the Devin AI sandbox workspace. |
connection_ids |
array<string> | null |
no | null |
Optional list of connection UUIDs. At least one of organization_id, workspace_id, or connection_ids is required. |
status_filter |
enum("all", "succeeded", "failed") |
no | "all" |
Filter by job status: all (default), succeeded, or failed. Applied to jobs.status in the Prod DB Replica. |
limit |
integer |
no | 1000 |
Maximum number of attempt rows to return. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"start_at": {
"description": "Inclusive start timestamp for the sync activity window. Must be timezone-aware (ISO 8601 with offset or `Z`).",
"format": "date-time",
"type": "string"
},
"end_at": {
"description": "Exclusive end timestamp for the sync activity window. Must be timezone-aware and strictly after `start_at`.",
"format": "date-time",
"type": "string"
},
"organization_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Organization ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@airbyte-internal\") and its value\nis the actual organization UUID. Use `OrganizationAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"664c690e-5263-49ba-b01f-4a6759b3330a"
],
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Optional organization UUID or alias. At least one of `organization_id`, `workspace_id`, or `connection_ids` is required. Accepts `@airbyte-internal` as an alias for the Airbyte internal org."
},
"workspace_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Workspace ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@devin-ai-sandbox\") and its value\nis the actual workspace UUID. Use `WorkspaceAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"266ebdfe-0d7b-4540-9817-de7e4505ba61"
],
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Optional workspace UUID or alias. At least one of `organization_id`, `workspace_id`, or `connection_ids` is required. Accepts `@devin-ai-sandbox` as an alias for the Devin AI sandbox workspace."
},
"connection_ids": {
"anyOf": [
{
"items": {
"type": "string"
},
"type": "array"
},
{
"type": "null"
}
],
"default": null,
"description": "Optional list of connection UUIDs. At least one of `organization_id`, `workspace_id`, or `connection_ids` is required."
},
"status_filter": {
"description": "Filter by job status: `all` (default), `succeeded`, or `failed`. Applied to `jobs.status` in the Prod DB Replica.",
"enum": [
"all",
"succeeded",
"failed"
],
"type": "string",
"default": "all"
},
"limit": {
"default": 1000,
"description": "Maximum number of attempt rows to return.",
"type": "integer"
}
},
"required": [
"start_at",
"end_at"
],
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_connections_by_connector
Search for all connections using a specific source or destination connector type.
This tool queries the Airbyte Cloud Prod DB Replica directly for fast results. It finds all connections where the source or destination connector matches the specified type, regardless of how the connector is named by users.
Results are always enriched with customer_tier and is_eu fields. The customer_tier_filter parameter is required to ensure tier-aware querying.
Optionally filter by organization_id to limit results to a specific organization. Use '@airbyte-internal' as an alias for the Airbyte internal organization.
Set exclude_pinned=True to filter out connections that are already pinned to a
specific version. This is useful for 'prove fix' live connection testing workflows
where you want to find unpinned connections to test against.
Set enabled_schedules_only=True to restrict results to connections that are both
enabled (status='active') and on an automated schedule (not manual-trigger-only).
This is useful for canary prerelease workflows where you need connections that
will run organically during the monitoring window.
Returns a list of connection dicts with workspace context and clickable Cloud UI URLs. For source queries, returns: connection_id, connection_name, connection_url, source_id, source_name, source_definition_id, workspace_id, workspace_name, organization_id, dataplane_group_id, dataplane_name, pin_origin_type, pin_origin, pinned_version_id, pin_scope_type, customer_tier, is_eu. For destination queries, returns: connection_id, connection_name, connection_url, destination_id, destination_name, destination_definition_id, workspace_id, workspace_name, organization_id, dataplane_group_id, dataplane_name, pin_origin_type, pin_origin, pinned_version_id, pin_scope_type, customer_tier, is_eu.
pin_scope_type is 'actor', 'workspace', or 'organization' indicating which scope level the effective pin came from (NULL if not pinned).
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
source_definition_id |
string | null |
no | null |
Source connector definition ID (UUID) to search for. Exactly one of source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name is required. Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics. |
source_canonical_name |
string | null |
no | null |
Canonical source connector name to search for. Exactly one of source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name is required. Examples: 'source-youtube-analytics', 'YouTube Analytics'. |
destination_definition_id |
string | null |
no | null |
Destination connector definition ID (UUID) to search for. Exactly one of source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name is required. Example: 'e5c8e66c-a480-4a5e-9c0e-e8e5e4c5c5c5' for DuckDB. |
destination_canonical_name |
string | null |
no | null |
Canonical destination connector name to search for. Exactly one of source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name is required. Examples: 'destination-duckdb', 'DuckDB'. |
organization_id |
string | enum("664c690e-5263-49ba-b01f-4a6759b3330a") | null |
no | null |
Optional organization ID (UUID) or alias to filter results. If provided, only connections in this organization will be returned. Accepts '@airbyte-internal' as an alias for the Airbyte internal org. |
limit |
integer |
no | 1000 |
Maximum number of results (default: 1000) |
customer_tier_filter |
enum("TIER_0", "TIER_1", "TIER_2", "UNKNOWN", "ALL") |
no | "TIER_2" |
Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. Filters results to only include connections belonging to organizations in the specified tier. Use 'ALL' to include all tiers. |
exclude_pinned |
boolean |
no | false |
If True, exclude connections whose connector is already pinned to a specific version (at any scope level: actor, workspace, or organization). Useful for 'prove fix' workflows where you want to find unpinned connections for live testing. Default: False (include all connections). |
enabled_schedules_only |
boolean |
no | false |
If True, only return connections that are both active (not paused/inactive) and on an automated sync schedule (not manual-trigger-only). Useful for canary workflows where you need connections that will produce organic syncs during a monitoring window. Default: False (include all connections). |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"source_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Source connector definition ID (UUID) to search for. Exactly one of source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name is required. Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics."
},
"source_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical source connector name to search for. Exactly one of source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name is required. Examples: 'source-youtube-analytics', 'YouTube Analytics'."
},
"destination_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Destination connector definition ID (UUID) to search for. Exactly one of source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name is required. Example: 'e5c8e66c-a480-4a5e-9c0e-e8e5e4c5c5c5' for DuckDB."
},
"destination_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical destination connector name to search for. Exactly one of source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name is required. Examples: 'destination-duckdb', 'DuckDB'."
},
"organization_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Organization ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@airbyte-internal\") and its value\nis the actual organization UUID. Use `OrganizationAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"664c690e-5263-49ba-b01f-4a6759b3330a"
],
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Optional organization ID (UUID) or alias to filter results. If provided, only connections in this organization will be returned. Accepts '@airbyte-internal' as an alias for the Airbyte internal org."
},
"limit": {
"default": 1000,
"description": "Maximum number of results (default: 1000)",
"type": "integer"
},
"customer_tier_filter": {
"default": "TIER_2",
"description": "Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. Filters results to only include connections belonging to organizations in the specified tier. Use 'ALL' to include all tiers.",
"enum": [
"TIER_0",
"TIER_1",
"TIER_2",
"UNKNOWN",
"ALL"
],
"type": "string"
},
"exclude_pinned": {
"default": false,
"description": "If True, exclude connections whose connector is already pinned to a specific version (at any scope level: actor, workspace, or organization). Useful for 'prove fix' workflows where you want to find unpinned connections for live testing. Default: False (include all connections).",
"type": "boolean"
},
"enabled_schedules_only": {
"default": false,
"description": "If True, only return connections that are both active (not paused/inactive) and on an automated sync schedule (not manual-trigger-only). Useful for canary workflows where you need connections that will produce organic syncs during a monitoring window. Default: False (include all connections).",
"type": "boolean"
}
},
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_connections_by_stream
Find connections that have a specific stream enabled in their catalog.
This tool searches the connection's configured catalog (JSONB) for streams matching the specified name. It's particularly useful when validating connector fixes that affect specific streams - you can quickly find customer connections that use the affected stream.
Results are always enriched with customer_tier and is_eu fields. The customer_tier_filter parameter is required to ensure tier-aware querying.
Use cases:
- Finding connections with a specific stream enabled for regression testing
- Validating connector fixes that affect particular streams
- Identifying which customers use rarely-enabled streams
Returns a list of connection dicts with workspace context and clickable Cloud UI URLs.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
stream_name |
string |
yes | — | Name of the stream to search for in connection catalogs. This must match the exact stream name as configured in the connection. Examples: 'global_exclusions', 'campaigns', 'users'. |
source_definition_id |
string | null |
no | null |
Source connector definition ID (UUID) to search for. Provide this OR source_canonical_name (exactly one required). Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics. |
source_canonical_name |
string | null |
no | null |
Canonical source connector name to search for. Provide this OR source_definition_id (exactly one required). Examples: 'source-klaviyo', 'Klaviyo', 'source-youtube-analytics'. |
organization_id |
string | enum("664c690e-5263-49ba-b01f-4a6759b3330a") | null |
no | null |
Optional organization ID (UUID) or alias to filter results. If provided, only connections in this organization will be returned. Accepts '@airbyte-internal' as an alias for the Airbyte internal org. |
limit |
integer |
no | 100 |
Maximum number of results (default: 100) |
customer_tier_filter |
enum("TIER_0", "TIER_1", "TIER_2", "UNKNOWN", "ALL") |
no | "TIER_2" |
Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. Filters results to only include connections belonging to organizations in the specified tier. Use 'ALL' to include all tiers. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"stream_name": {
"description": "Name of the stream to search for in connection catalogs. This must match the exact stream name as configured in the connection. Examples: 'global_exclusions', 'campaigns', 'users'.",
"type": "string"
},
"source_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Source connector definition ID (UUID) to search for. Provide this OR source_canonical_name (exactly one required). Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics."
},
"source_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical source connector name to search for. Provide this OR source_definition_id (exactly one required). Examples: 'source-klaviyo', 'Klaviyo', 'source-youtube-analytics'."
},
"organization_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Organization ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@airbyte-internal\") and its value\nis the actual organization UUID. Use `OrganizationAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"664c690e-5263-49ba-b01f-4a6759b3330a"
],
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Optional organization ID (UUID) or alias to filter results. If provided, only connections in this organization will be returned. Accepts '@airbyte-internal' as an alias for the Airbyte internal org."
},
"limit": {
"default": 100,
"description": "Maximum number of results (default: 100)",
"type": "integer"
},
"customer_tier_filter": {
"default": "TIER_2",
"description": "Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. Filters results to only include connections belonging to organizations in the specified tier. Use 'ALL' to include all tiers.",
"enum": [
"TIER_0",
"TIER_1",
"TIER_2",
"UNKNOWN",
"ALL"
],
"type": "string"
}
},
"required": [
"stream_name"
],
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_connector_connection_stats
Get aggregate connection stats for multiple connectors.
Returns counts of connections grouped by pinned version for each connector, including:
- Total, enabled, and active connection counts
- Pinned vs unpinned breakdown
- Latest attempt status breakdown (succeeded, failed, cancelled, running, unknown)
This tool is designed for release monitoring workflows. It allows you to:
- Query recently released connectors to identify which ones to monitor
- Get aggregate stats showing how many connections are using each version
- See health metrics (pass/fail) broken down by version
The lookback_days parameter controls the lookback window for:
- Counting 'active' connections (those with recent sync activity)
- Determining 'latest attempt status' (most recent attempt within the window)
Connections with no sync activity in the lookback window will have 'unknown' status in the latest_attempt breakdown.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
source_definition_ids |
array<string> | null |
no | null |
List of source connector definition IDs (UUIDs) to get stats for. Example: ['afa734e4-3571-11ec-991a-1e0031268139'] |
destination_definition_ids |
array<string> | null |
no | null |
List of destination connector definition IDs (UUIDs) to get stats for. Example: ['94bd199c-2ff0-4aa2-b98e-17f0acb72610'] |
lookback_days |
integer |
no | 7 |
Number of days to look back for 'active' connections (default: 7). Connections with sync activity within this window are counted as active. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"source_definition_ids": {
"anyOf": [
{
"items": {
"type": "string"
},
"type": "array"
},
{
"type": "null"
}
],
"default": null,
"description": "List of source connector definition IDs (UUIDs) to get stats for. Example: ['afa734e4-3571-11ec-991a-1e0031268139']"
},
"destination_definition_ids": {
"anyOf": [
{
"items": {
"type": "string"
},
"type": "array"
},
{
"type": "null"
}
],
"default": null,
"description": "List of destination connector definition IDs (UUIDs) to get stats for. Example: ['94bd199c-2ff0-4aa2-b98e-17f0acb72610']"
},
"lookback_days": {
"default": 7,
"description": "Number of days to look back for 'active' connections (default: 7). Connections with sync activity within this window are counted as active.",
"type": "integer"
}
},
"type": "object"
}
Show output JSON schema
{
"description": "Response containing connection stats for multiple connectors.",
"properties": {
"sources": {
"description": "Stats for source connectors",
"items": {
"description": "Aggregate connection stats for a connector.",
"properties": {
"connector_definition_id": {
"description": "The connector definition UUID",
"type": "string"
},
"connector_type": {
"description": "'source' or 'destination'",
"type": "string"
},
"canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The canonical connector name if resolved"
},
"total_connections": {
"description": "Total number of non-deprecated connections",
"type": "integer"
},
"enabled_connections": {
"description": "Number of enabled (active status) connections",
"type": "integer"
},
"active_connections": {
"description": "Number of connections with recent sync activity",
"type": "integer"
},
"pinned_connections": {
"description": "Number of connections with explicit version pins",
"type": "integer"
},
"unpinned_connections": {
"description": "Number of connections on default version",
"type": "integer"
},
"latest_attempt": {
"description": "Overall breakdown by latest attempt status",
"properties": {
"succeeded": {
"default": 0,
"description": "Connections where latest attempt succeeded",
"type": "integer"
},
"failed": {
"default": 0,
"description": "Connections where latest attempt failed",
"type": "integer"
},
"cancelled": {
"default": 0,
"description": "Connections where latest attempt was cancelled",
"type": "integer"
},
"running": {
"default": 0,
"description": "Connections where latest attempt is still running",
"type": "integer"
},
"unknown": {
"default": 0,
"description": "Connections with no recent attempts in the lookback window",
"type": "integer"
}
},
"type": "object"
},
"by_version": {
"description": "Stats broken down by pinned version",
"items": {
"description": "Stats for connections pinned to a specific version.",
"properties": {
"pinned_version_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"description": "The connector version UUID (None for unpinned connections)"
},
"docker_image_tag": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The docker image tag for this version"
},
"total_connections": {
"description": "Total number of connections",
"type": "integer"
},
"enabled_connections": {
"description": "Number of enabled (active status) connections",
"type": "integer"
},
"active_connections": {
"description": "Number of connections with recent sync activity",
"type": "integer"
},
"latest_attempt": {
"description": "Breakdown by latest attempt status",
"properties": {
"succeeded": {
"default": 0,
"description": "Connections where latest attempt succeeded",
"type": "integer"
},
"failed": {
"default": 0,
"description": "Connections where latest attempt failed",
"type": "integer"
},
"cancelled": {
"default": 0,
"description": "Connections where latest attempt was cancelled",
"type": "integer"
},
"running": {
"default": 0,
"description": "Connections where latest attempt is still running",
"type": "integer"
},
"unknown": {
"default": 0,
"description": "Connections with no recent attempts in the lookback window",
"type": "integer"
}
},
"type": "object"
}
},
"required": [
"pinned_version_id",
"total_connections",
"enabled_connections",
"active_connections",
"latest_attempt"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"connector_definition_id",
"connector_type",
"total_connections",
"enabled_connections",
"active_connections",
"pinned_connections",
"unpinned_connections",
"latest_attempt",
"by_version"
],
"type": "object"
},
"type": "array"
},
"destinations": {
"description": "Stats for destination connectors",
"items": {
"description": "Aggregate connection stats for a connector.",
"properties": {
"connector_definition_id": {
"description": "The connector definition UUID",
"type": "string"
},
"connector_type": {
"description": "'source' or 'destination'",
"type": "string"
},
"canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The canonical connector name if resolved"
},
"total_connections": {
"description": "Total number of non-deprecated connections",
"type": "integer"
},
"enabled_connections": {
"description": "Number of enabled (active status) connections",
"type": "integer"
},
"active_connections": {
"description": "Number of connections with recent sync activity",
"type": "integer"
},
"pinned_connections": {
"description": "Number of connections with explicit version pins",
"type": "integer"
},
"unpinned_connections": {
"description": "Number of connections on default version",
"type": "integer"
},
"latest_attempt": {
"description": "Overall breakdown by latest attempt status",
"properties": {
"succeeded": {
"default": 0,
"description": "Connections where latest attempt succeeded",
"type": "integer"
},
"failed": {
"default": 0,
"description": "Connections where latest attempt failed",
"type": "integer"
},
"cancelled": {
"default": 0,
"description": "Connections where latest attempt was cancelled",
"type": "integer"
},
"running": {
"default": 0,
"description": "Connections where latest attempt is still running",
"type": "integer"
},
"unknown": {
"default": 0,
"description": "Connections with no recent attempts in the lookback window",
"type": "integer"
}
},
"type": "object"
},
"by_version": {
"description": "Stats broken down by pinned version",
"items": {
"description": "Stats for connections pinned to a specific version.",
"properties": {
"pinned_version_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"description": "The connector version UUID (None for unpinned connections)"
},
"docker_image_tag": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The docker image tag for this version"
},
"total_connections": {
"description": "Total number of connections",
"type": "integer"
},
"enabled_connections": {
"description": "Number of enabled (active status) connections",
"type": "integer"
},
"active_connections": {
"description": "Number of connections with recent sync activity",
"type": "integer"
},
"latest_attempt": {
"description": "Breakdown by latest attempt status",
"properties": {
"succeeded": {
"default": 0,
"description": "Connections where latest attempt succeeded",
"type": "integer"
},
"failed": {
"default": 0,
"description": "Connections where latest attempt failed",
"type": "integer"
},
"cancelled": {
"default": 0,
"description": "Connections where latest attempt was cancelled",
"type": "integer"
},
"running": {
"default": 0,
"description": "Connections where latest attempt is still running",
"type": "integer"
},
"unknown": {
"default": 0,
"description": "Connections with no recent attempts in the lookback window",
"type": "integer"
}
},
"type": "object"
}
},
"required": [
"pinned_version_id",
"total_connections",
"enabled_connections",
"active_connections",
"latest_attempt"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"connector_definition_id",
"connector_type",
"total_connections",
"enabled_connections",
"active_connections",
"pinned_connections",
"unpinned_connections",
"latest_attempt",
"by_version"
],
"type": "object"
},
"type": "array"
},
"lookback_days": {
"description": "Lookback window used for 'active' connections",
"type": "integer"
},
"generated_at": {
"description": "When this response was generated",
"format": "date-time",
"type": "string"
}
},
"required": [
"lookback_days",
"generated_at"
],
"type": "object"
}
query_prod_connector_rollouts
Query connector rollouts with flexible filtering.
Returns rollouts based on the provided filters. If no filters are specified, returns all active rollouts. Useful for monitoring rollout status and history.
Filter behavior:
- rollout_id: Returns that specific rollout (ignores other filters)
- active_only: Returns only active (non-terminal) rollouts
- actor_definition_id: Returns rollouts for that specific connector
- No filters: Returns all active rollouts (same as active_only=True)
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
actor_definition_id |
string | null |
no | null |
Connector definition UUID to filter by (optional) |
rollout_id |
string | null |
no | null |
Specific rollout UUID to look up (optional) |
active_only |
boolean |
no | false |
If true, only return active (non-terminal) rollouts |
limit |
integer |
no | 100 |
Maximum number of results (default: 100) |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"actor_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector definition UUID to filter by (optional)"
},
"rollout_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Specific rollout UUID to look up (optional)"
},
"active_only": {
"default": false,
"description": "If true, only return active (non-terminal) rollouts",
"type": "boolean"
},
"limit": {
"default": 100,
"description": "Maximum number of results (default: 100)",
"type": "integer"
}
},
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"description": "Information about a connector rollout.",
"properties": {
"rollout_id": {
"description": "The rollout UUID",
"type": "string"
},
"actor_definition_id": {
"description": "The connector definition UUID",
"type": "string"
},
"state": {
"description": "Rollout state: initialized, workflow_started, in_progress, paused, finalizing, succeeded, errored, failed_rolled_back, canceled",
"type": "string"
},
"initial_rollout_pct": {
"anyOf": [
{
"type": "integer"
},
{
"type": "null"
}
],
"default": null,
"description": "Initial rollout percentage"
},
"current_target_rollout_pct": {
"anyOf": [
{
"type": "integer"
},
{
"type": "null"
}
],
"default": null,
"description": "Current target rollout percentage"
},
"final_target_rollout_pct": {
"anyOf": [
{
"type": "integer"
},
{
"type": "null"
}
],
"default": null,
"description": "Final target rollout percentage"
},
"has_breaking_changes": {
"description": "Whether the RC has breaking changes",
"type": "boolean"
},
"max_step_wait_time_mins": {
"anyOf": [
{
"type": "integer"
},
{
"type": "null"
}
],
"default": null,
"description": "Maximum wait time between rollout steps in minutes"
},
"rollout_strategy": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Rollout strategy: manual, automated, overridden"
},
"updated_by_user_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "User ID recorded as last updating the rollout"
},
"updated_by_user_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Name recorded as last updating the rollout"
},
"updated_by_user_email": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Email recorded as last updating the rollout"
},
"workflow_run_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Temporal workflow run ID"
},
"error_msg": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Error message if errored"
},
"failed_reason": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Reason for failure if failed"
},
"paused_reason": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Reason for pause if paused"
},
"tag": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Optional tag for the rollout"
},
"created_at": {
"anyOf": [
{
"format": "date-time",
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "When the rollout was created"
},
"updated_at": {
"anyOf": [
{
"format": "date-time",
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "When the rollout was last updated"
},
"completed_at": {
"anyOf": [
{
"format": "date-time",
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "When the rollout completed (if terminal)"
},
"expires_at": {
"anyOf": [
{
"format": "date-time",
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "When the rollout expires"
},
"rc_docker_image_tag": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Docker image tag of the release candidate"
},
"rc_docker_repository": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Docker repository of the release candidate"
},
"initial_docker_image_tag": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Docker image tag of the initial version"
},
"initial_docker_repository": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Docker repository of the initial version"
},
"filters": {
"anyOf": [
{
"additionalProperties": true,
"type": "object"
},
{
"type": "null"
}
],
"default": null,
"description": "Raw rollout filters JSON (e.g., {'tierFilter': {'tier': 'TIER_0'}})"
},
"customer_tier": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Customer tier targeted by this rollout (extracted from filters), e.g., 'TIER_0', 'TIER_1'. None if no tier filter is set."
}
},
"required": [
"rollout_id",
"actor_definition_id",
"state",
"has_breaking_changes"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_connector_versions
List all versions for a connector definition.
Returns all published versions of a connector, ordered by last_published date descending. Useful for understanding version history and finding specific version IDs for pinning or rollout monitoring.
Returns list of dicts with keys: version_id, docker_image_tag, docker_repository, release_stage, support_level, cdk_version, language, last_published, release_date
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
connector_definition_id |
string |
yes | — | Connector definition UUID to list versions for |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"connector_definition_id": {
"description": "Connector definition UUID to list versions for",
"type": "string"
}
},
"required": [
"connector_definition_id"
],
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_dataplanes
List all dataplane groups with workspace counts.
Returns information about all active dataplane groups in Airbyte Cloud, including the number of workspaces in each. Useful for understanding the distribution of workspaces across regions (US, US-Central, EU).
Returns list of dicts with keys: dataplane_group_id, dataplane_name, organization_id, enabled, tombstone, created_at, workspace_count
Parameters:
_No parameters._
Show input JSON schema
{
"additionalProperties": false,
"properties": {},
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_failed_sync_attempts_for_connector
List failed sync attempts for ALL actors using a connector type.
This tool finds all actors with the given connector definition and returns their failed sync attempts, regardless of whether they have explicit version pins.
Results are always enriched with customer_tier and is_eu fields. The customer_tier_filter parameter is required to ensure tier-aware querying.
This is useful for investigating connector issues across all users. Use this when you want to find failures for a connector type regardless of which version users are on.
Supports both SOURCE and DESTINATION connectors. Provide exactly one of: source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name.
Key fields in results:
- failure_summary: JSON containing failure details including failureType and messages
- customer_tier: TIER_0, TIER_1, TIER_2, or UNKNOWN
- is_eu: Whether the workspace is in the EU region
- pin_origin_type, pin_origin, pinned_version_id: Version pin context (NULL if not pinned)
- pin_scope_type: 'actor', 'workspace', or 'organization' (NULL if not pinned)
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
source_definition_id |
string | null |
no | null |
Source connector definition ID (UUID) to search for. Provide this OR source_canonical_name OR destination_definition_id OR destination_canonical_name (exactly one required). Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics. |
source_canonical_name |
string | null |
no | null |
Canonical source connector name to search for. Provide this OR source_definition_id OR destination_definition_id OR destination_canonical_name (exactly one required). Examples: 'source-youtube-analytics', 'YouTube Analytics'. |
destination_definition_id |
string | null |
no | null |
Destination connector definition ID (UUID) to search for. Provide this OR destination_canonical_name OR source_definition_id OR source_canonical_name (exactly one required). Example: '94bd199c-2ff0-4aa2-b98e-17f0acb72610' for DuckDB. |
destination_canonical_name |
string | null |
no | null |
Canonical destination connector name to search for. Provide this OR destination_definition_id OR source_definition_id OR source_canonical_name (exactly one required). Examples: 'destination-duckdb', 'DuckDB'. |
organization_id |
string | enum("664c690e-5263-49ba-b01f-4a6759b3330a") | null |
no | null |
Optional organization ID (UUID) or alias to filter results. If provided, only failed attempts from this organization will be returned. Accepts '@airbyte-internal' as an alias for the Airbyte internal org. |
lookback_days |
integer |
no | 7 |
Number of days to look back (default: 7) |
limit |
integer |
no | 100 |
Maximum number of results (default: 100) |
customer_tier_filter |
enum("TIER_0", "TIER_1", "TIER_2", "UNKNOWN", "ALL") |
no | "TIER_2" |
Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. Filters results to only include connections belonging to organizations in the specified tier. Use 'ALL' to include all tiers. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"source_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Source connector definition ID (UUID) to search for. Provide this OR source_canonical_name OR destination_definition_id OR destination_canonical_name (exactly one required). Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics."
},
"source_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical source connector name to search for. Provide this OR source_definition_id OR destination_definition_id OR destination_canonical_name (exactly one required). Examples: 'source-youtube-analytics', 'YouTube Analytics'."
},
"destination_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Destination connector definition ID (UUID) to search for. Provide this OR destination_canonical_name OR source_definition_id OR source_canonical_name (exactly one required). Example: '94bd199c-2ff0-4aa2-b98e-17f0acb72610' for DuckDB."
},
"destination_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical destination connector name to search for. Provide this OR destination_definition_id OR source_definition_id OR source_canonical_name (exactly one required). Examples: 'destination-duckdb', 'DuckDB'."
},
"organization_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Organization ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@airbyte-internal\") and its value\nis the actual organization UUID. Use `OrganizationAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"664c690e-5263-49ba-b01f-4a6759b3330a"
],
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Optional organization ID (UUID) or alias to filter results. If provided, only failed attempts from this organization will be returned. Accepts '@airbyte-internal' as an alias for the Airbyte internal org."
},
"lookback_days": {
"default": 7,
"description": "Number of days to look back (default: 7)",
"type": "integer"
},
"limit": {
"default": 100,
"description": "Maximum number of results (default: 100)",
"type": "integer"
},
"customer_tier_filter": {
"default": "TIER_2",
"description": "Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. Filters results to only include connections belonging to organizations in the specified tier. Use 'ALL' to include all tiers.",
"enum": [
"TIER_0",
"TIER_1",
"TIER_2",
"UNKNOWN",
"ALL"
],
"type": "string"
}
},
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_new_connector_releases
List recently published connector versions.
Returns connector versions published within the specified number of days. Uses last_published timestamp which reflects when the version was actually deployed to the registry (not the changelog date).
Returns list of dicts with keys: version_id, connector_definition_id, docker_repository, docker_image_tag, last_published, release_date, release_stage, support_level, cdk_version, language, created_at
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
days |
integer |
no | 7 |
Number of days to look back (default: 7) |
limit |
integer |
no | 100 |
Maximum number of results (default: 100) |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"days": {
"default": 7,
"description": "Number of days to look back (default: 7)",
"type": "integer"
},
"limit": {
"default": 100,
"description": "Maximum number of results (default: 100)",
"type": "integer"
}
},
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_org_admin_contacts
Find Org Admins and Workspace Admins of the given workspaces.
Activity is each user's latest user-attributed connection timeline event in
this organization within the lookback window. Emails are returned because
callers need them to CC customers; never log them. Only direct user grants
are included; group-granted permissions are not queried. The org's
customer_tier is included for awareness only; outreach must still reach
every affected organization regardless of tier.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
organization_id |
string |
yes | — | Organization UUID whose direct Org Admins to find. |
workspace_ids |
array<string> |
no | [] |
Workspace UUIDs whose direct Workspace Admins to include. Only these workspaces contribute Workspace Admins. |
activity_lookback_days |
integer |
no | 90 |
Days to look back for user-attributed connection activity. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"organization_id": {
"description": "Organization UUID whose direct Org Admins to find.",
"type": "string"
},
"workspace_ids": {
"default": [],
"description": "Workspace UUIDs whose direct Workspace Admins to include. Only these workspaces contribute Workspace Admins.",
"items": {
"type": "string"
},
"type": "array"
},
"activity_lookback_days": {
"default": 90,
"description": "Days to look back for user-attributed connection activity.",
"maximum": 365,
"minimum": 1,
"type": "integer"
}
},
"required": [
"organization_id"
],
"type": "object"
}
Show output JSON schema
{
"description": "Organization administrator contacts and recent user activity.",
"properties": {
"organization_id": {
"description": "The organization UUID.",
"type": "string"
},
"customer_tier": {
"description": "Organization customer tier, included for awareness only.",
"type": "string"
},
"tier_warnings": {
"description": "Warnings raised while resolving the organization customer tier.",
"items": {
"type": "string"
},
"type": "array"
},
"activity_lookback_days": {
"description": "Number of days used to look back for connection activity.",
"type": "integer"
},
"admins": {
"description": "Administrators ordered by Airbyte user UUID.",
"items": {
"description": "An organization or selected-workspace administrator contact.",
"properties": {
"user_id": {
"description": "The Airbyte user UUID.",
"type": "string"
},
"name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The administrator's name."
},
"email": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The administrator's email address."
},
"status": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The user's status."
},
"is_org_admin": {
"description": "Whether the user is an Org Admin.",
"type": "boolean"
},
"admin_workspace_ids": {
"description": "Sorted workspace UUIDs where the user is a Workspace Admin.",
"items": {
"type": "string"
},
"type": "array"
},
"user_created_at": {
"anyOf": [
{
"format": "date-time",
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "When the user account was created."
},
"user_updated_at": {
"anyOf": [
{
"format": "date-time",
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "When the user account was last updated."
},
"last_connection_event_at": {
"anyOf": [
{
"format": "date-time",
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Latest user-attributed connection timeline event in the activity window."
}
},
"required": [
"user_id",
"is_org_admin",
"admin_workspace_ids"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"organization_id",
"customer_tier",
"activity_lookback_days",
"admins"
],
"type": "object"
}
query_prod_organizations
Search organizations by name or email substring.
Performs a case-insensitive substring match on organization name and email.
Use the returned organization_id values with other tools like
query_prod_connections_by_connector or lookup_customer_tiers.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
name_contains |
string |
yes | — | Case-insensitive substring to search for in organization name or email. For example, 'acme' will match organizations named 'Acme Corp' or with email 'admin@acme.io'. |
limit |
integer |
no | 20 |
Maximum number of organizations to return (default: 20) |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"name_contains": {
"description": "Case-insensitive substring to search for in organization name or email. For example, 'acme' will match organizations named 'Acme Corp' or with email 'admin@acme.io'.",
"type": "string"
},
"limit": {
"default": 20,
"description": "Maximum number of organizations to return (default: 20)",
"type": "integer"
}
},
"required": [
"name_contains"
],
"type": "object"
}
Show output JSON schema
{
"description": "Result of searching organizations by name substring.",
"properties": {
"name_contains": {
"description": "The search substring that was used",
"type": "string"
},
"total_found": {
"description": "Total number of organizations matching",
"type": "integer"
},
"organizations": {
"description": "List of matching organizations",
"items": {
"description": "A single organization returned by a name/email search.",
"properties": {
"organization_id": {
"description": "The organization UUID",
"type": "string"
},
"organization_name": {
"description": "The name of the organization",
"type": "string"
},
"email": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The email address associated with the organization"
},
"customer_tier": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Customer tier (TIER_0, TIER_1, TIER_2, or UNKNOWN). Enriched from the GCS tier cache."
}
},
"required": [
"organization_id",
"organization_name"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"name_contains",
"total_found",
"organizations"
],
"type": "object"
}
query_prod_pin_stats_for_organization
Query connector versions pinned anywhere under an organization.
Returns one row per pinned version, aggregating every connector_version
pin whose scope belongs to the organization — the org itself, one of its
workspaces, or an actor within one of those workspaces (actor, workspace,
and organization scopes). Each row carries the per-scope pin breakdown, the
manual/rollout/breaking-change split, and a has_active_rollout flag.
This powers the first step of the Organization Pins view (pick an org, then
see the versions pinned under it). Use query_prod_pins_for_organization
for the individual pins behind a selected version.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
organization_id |
string | enum("664c690e-5263-49ba-b01f-4a6759b3330a") |
yes | — | Organization UUID (or @airbyte-internal alias) to scope pins to. Resolve organization names to an ID first via search_organizations. |
connector_definition_id |
string | null |
no | null |
Connector definition UUID to filter by (optional). Mutually exclusive with connector_canonical_name. |
connector_canonical_name |
string | null |
no | null |
Connector canonical name (e.g. source-postgres) to filter by. Resolved to a definition ID via the registry. Mutually exclusive with connector_definition_id. |
limit |
integer |
no | 1000 |
Maximum number of versions to return (default: 1000). |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"organization_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Organization ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@airbyte-internal\") and its value\nis the actual organization UUID. Use `OrganizationAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"664c690e-5263-49ba-b01f-4a6759b3330a"
],
"type": "string"
}
],
"description": "Organization UUID (or `@airbyte-internal` alias) to scope pins to. Resolve organization names to an ID first via `search_organizations`."
},
"connector_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector definition UUID to filter by (optional). Mutually exclusive with `connector_canonical_name`."
},
"connector_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector canonical name (e.g. `source-postgres`) to filter by. Resolved to a definition ID via the registry. Mutually exclusive with `connector_definition_id`."
},
"limit": {
"default": 1000,
"description": "Maximum number of versions to return (default: 1000).",
"type": "integer"
}
},
"required": [
"organization_id"
],
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"description": "A connector version pinned somewhere under an organization, with counts.",
"properties": {
"version_id": {
"description": "The actor_definition_version UUID",
"type": "string"
},
"connector_definition_id": {
"description": "The connector definition UUID",
"type": "string"
},
"connector_name": {
"description": "Human-readable connector name",
"type": "string"
},
"docker_repository": {
"description": "Docker repository path",
"type": "string"
},
"docker_image_tag": {
"description": "Docker image tag for this version",
"type": "string"
},
"last_published": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "ISO timestamp when this version was last published"
},
"pin_count": {
"description": "Total pins under the org targeting this version (all scopes)",
"type": "integer"
},
"manual_pins": {
"default": 0,
"description": "Pins with no system origin (user-created manual pins), any scope",
"type": "integer"
},
"rollout_pins": {
"default": 0,
"description": "Pins created by connector rollouts",
"type": "integer"
},
"breaking_change_pins": {
"default": 0,
"description": "Pins created by breaking changes",
"type": "integer"
},
"actor_pins": {
"description": "Manual actor-scoped pins (excludes rollout and breaking-change)",
"type": "integer"
},
"workspace_pins": {
"description": "Workspace-scoped pins under the org",
"type": "integer"
},
"org_pins": {
"description": "Organization-scoped pins",
"type": "integer"
},
"has_active_rollout": {
"default": false,
"description": "`True` if at least one rollout pin is backed by a non-terminal `connector_rollout`",
"type": "boolean"
}
},
"required": [
"version_id",
"connector_definition_id",
"connector_name",
"docker_repository",
"docker_image_tag",
"pin_count",
"actor_pins",
"workspace_pins",
"org_pins"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_pins_for_organization
List the individual connector-version pins discovered under an organization.
Returns one row per scoped_configuration pin whose scope belongs to the
organization (org/workspace/actor), resolving the pinned connector and
version, the scope's display name, the manual author's email, and — for
rollout-origin pins — the backing connector_rollout id and state. This
directly answers whether each pin is manual or caused by an active rollout.
This powers the second step of the Organization Pins view: after picking a
version from query_prod_pin_stats_for_organization, pass its
pinned_version_id here to list the pins behind it.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
organization_id |
string | enum("664c690e-5263-49ba-b01f-4a6759b3330a") |
yes | — | Organization UUID (or @airbyte-internal alias) to scope pins to. Resolve organization names to an ID first via search_organizations. |
connector_definition_id |
string | null |
no | null |
Connector definition UUID to filter by (optional). Mutually exclusive with connector_canonical_name. |
connector_canonical_name |
string | null |
no | null |
Connector canonical name (e.g. source-postgres) to filter by. Resolved to a definition ID via the registry. Mutually exclusive with connector_definition_id. |
pinned_version_id |
string | null |
no | null |
Actor_definition_version UUID to return only pins targeting that version. This is the post-selection filter for the org pins tab. |
origin_filter |
enum("all", "manual", "connector_rollout", "breaking_change") |
no | "all" |
Restrict by how the pin was created: all (default), manual, connector_rollout, or breaking_change. |
limit |
integer |
no | 1000 |
Maximum number of pins to return (default: 1000). |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"organization_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Organization ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@airbyte-internal\") and its value\nis the actual organization UUID. Use `OrganizationAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"664c690e-5263-49ba-b01f-4a6759b3330a"
],
"type": "string"
}
],
"description": "Organization UUID (or `@airbyte-internal` alias) to scope pins to. Resolve organization names to an ID first via `search_organizations`."
},
"connector_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector definition UUID to filter by (optional). Mutually exclusive with `connector_canonical_name`."
},
"connector_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector canonical name (e.g. `source-postgres`) to filter by. Resolved to a definition ID via the registry. Mutually exclusive with `connector_definition_id`."
},
"pinned_version_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Actor_definition_version UUID to return only pins targeting that version. This is the post-selection filter for the org pins tab."
},
"origin_filter": {
"description": "Restrict by how the pin was created: `all` (default), `manual`, `connector_rollout`, or `breaking_change`.",
"enum": [
"all",
"manual",
"connector_rollout",
"breaking_change"
],
"type": "string",
"default": "all"
},
"limit": {
"default": 1000,
"description": "Maximum number of pins to return (default: 1000).",
"type": "integer"
}
},
"required": [
"organization_id"
],
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"description": "A single `scoped_configuration` pin discovered under an organization.",
"properties": {
"connector_definition_id": {
"description": "The connector definition UUID",
"type": "string"
},
"connector_name": {
"description": "Human-readable connector name",
"type": "string"
},
"docker_repository": {
"description": "Docker repository path",
"type": "string"
},
"pinned_version_id": {
"description": "The pinned actor_definition_version UUID",
"type": "string"
},
"pinned_version_tag": {
"description": "Docker image tag of the pinned version",
"type": "string"
},
"pin_scope_type": {
"description": "Scope of the pin: `organization`, `workspace`, or `actor`",
"type": "string"
},
"scope_id": {
"description": "UUID of the scoped entity",
"type": "string"
},
"scope_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Display name of the scoped entity, when resolvable"
},
"pin_category": {
"description": "Derived pin type: `manual`, `rollout`, or `breaking_change`",
"type": "string"
},
"set_by": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Email (or name) of the user who set a manual pin, when known"
},
"rollout_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Backing connector_rollout UUID for rollout pins"
},
"rollout_state": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "State of the backing rollout, for rollout pins"
},
"is_active_rollout": {
"default": false,
"description": "`True` when `rollout_state` is a non-terminal (active) state",
"type": "boolean"
},
"description": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Free-text pin reason"
},
"reference_url": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Reference URL attached to the pin, when present"
},
"created_at": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "ISO timestamp when the pin was created"
},
"expires_at": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "ISO timestamp when the pin expires, when set"
}
},
"required": [
"connector_definition_id",
"connector_name",
"docker_repository",
"pinned_version_id",
"pinned_version_tag",
"pin_scope_type",
"scope_id",
"pin_category"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_recent_syncs_for_connector
List recent sync jobs for ALL actors using a connector type.
This tool finds all actors with the given connector definition and returns their recent sync jobs, regardless of whether they have explicit version pins. It filters out deleted actors, deleted workspaces, and deprecated connections.
Results are always enriched with customer_tier and is_eu fields. The customer_tier_filter parameter is required to ensure tier-aware querying.
Use this tool to:
- Find healthy connections with recent successful syncs (status_filter='succeeded')
- Investigate connector issues across all users (status_filter='failed')
- Get an overview of all recent sync activity (status_filter='all')
Set exclude_pinned=True to filter out syncs for actors that are already pinned to a
specific version. This is useful for 'prove fix' live connection testing workflows
where you want to find unpinned connections to test against.
Set enabled_schedules_only=True to restrict results to connections that are both
enabled (status='active') and on an automated schedule (not manual-trigger-only).
This is useful for canary prerelease workflows where you need connections that
will run organically during the monitoring window.
Supports both SOURCE and DESTINATION connectors. Provide exactly one of: source_definition_id, source_canonical_name, destination_definition_id, or destination_canonical_name.
Key fields in results:
- job_status: 'succeeded', 'failed', 'cancelled', etc.
- connection_id, connection_name: The connection that ran the sync
- actor_id, actor_name: The source or destination actor
- customer_tier: TIER_0, TIER_1, TIER_2, or UNKNOWN
- is_eu: Whether the workspace is in the EU region
- pin_origin_type, pin_origin, pinned_version_id: Version pin context (NULL if not pinned)
- pin_scope_type: 'actor', 'workspace', or 'organization' (NULL if not pinned)
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
source_definition_id |
string | null |
no | null |
Source connector definition ID (UUID) to search for. Provide this OR source_canonical_name OR destination_definition_id OR destination_canonical_name (exactly one required). Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics. |
source_canonical_name |
string | null |
no | null |
Canonical source connector name to search for. Provide this OR source_definition_id OR destination_definition_id OR destination_canonical_name (exactly one required). Examples: 'source-youtube-analytics', 'YouTube Analytics'. |
destination_definition_id |
string | null |
no | null |
Destination connector definition ID (UUID) to search for. Provide this OR destination_canonical_name OR source_definition_id OR source_canonical_name (exactly one required). Example: '94bd199c-2ff0-4aa2-b98e-17f0acb72610' for DuckDB. |
destination_canonical_name |
string | null |
no | null |
Canonical destination connector name to search for. Provide this OR destination_definition_id OR source_definition_id OR source_canonical_name (exactly one required). Examples: 'destination-duckdb', 'DuckDB'. |
status_filter |
enum("all", "succeeded", "failed") |
no | "all" |
Filter by job status: 'all' (default), 'succeeded', or 'failed'. Use 'succeeded' to find healthy connections with recent successful syncs. Use 'failed' to find connections with recent failures. |
organization_id |
string | enum("664c690e-5263-49ba-b01f-4a6759b3330a") | null |
no | null |
Optional organization ID (UUID) or alias to filter results. If provided, only syncs from this organization will be returned. Accepts '@airbyte-internal' as an alias for the Airbyte internal org. |
lookback_days |
integer |
no | 7 |
Number of days to look back (default: 7) |
limit |
integer |
no | 100 |
Maximum number of results (default: 100) |
customer_tier_filter |
enum("TIER_0", "TIER_1", "TIER_2", "UNKNOWN", "ALL") |
no | "TIER_2" |
Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. Filters results to only include connections belonging to organizations in the specified tier. Use 'ALL' to include all tiers. |
exclude_pinned |
boolean |
no | false |
If True, exclude syncs for actors that are already pinned to a specific version (at any scope level: actor, workspace, or organization). Useful for 'prove fix' workflows where you want to find unpinned connections for live testing. Default: False (include all syncs). |
enabled_schedules_only |
boolean |
no | false |
If True, only return syncs for connections that are both active (not paused/inactive) and on an automated sync schedule (not manual-trigger-only). Useful for canary workflows where you need connections that will produce organic syncs during a monitoring window. Default: False (include all connections). |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"source_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Source connector definition ID (UUID) to search for. Provide this OR source_canonical_name OR destination_definition_id OR destination_canonical_name (exactly one required). Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics."
},
"source_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical source connector name to search for. Provide this OR source_definition_id OR destination_definition_id OR destination_canonical_name (exactly one required). Examples: 'source-youtube-analytics', 'YouTube Analytics'."
},
"destination_definition_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Destination connector definition ID (UUID) to search for. Provide this OR destination_canonical_name OR source_definition_id OR source_canonical_name (exactly one required). Example: '94bd199c-2ff0-4aa2-b98e-17f0acb72610' for DuckDB."
},
"destination_canonical_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical destination connector name to search for. Provide this OR destination_definition_id OR source_definition_id OR source_canonical_name (exactly one required). Examples: 'destination-duckdb', 'DuckDB'."
},
"status_filter": {
"description": "Filter by job status: 'all' (default), 'succeeded', or 'failed'. Use 'succeeded' to find healthy connections with recent successful syncs. Use 'failed' to find connections with recent failures.",
"enum": [
"all",
"succeeded",
"failed"
],
"type": "string",
"default": "all"
},
"organization_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Organization ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@airbyte-internal\") and its value\nis the actual organization UUID. Use `OrganizationAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"664c690e-5263-49ba-b01f-4a6759b3330a"
],
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Optional organization ID (UUID) or alias to filter results. If provided, only syncs from this organization will be returned. Accepts '@airbyte-internal' as an alias for the Airbyte internal org."
},
"lookback_days": {
"default": 7,
"description": "Number of days to look back (default: 7)",
"type": "integer"
},
"limit": {
"default": 100,
"description": "Maximum number of results (default: 100)",
"type": "integer"
},
"customer_tier_filter": {
"default": "TIER_2",
"description": "Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. Filters results to only include connections belonging to organizations in the specified tier. Use 'ALL' to include all tiers.",
"enum": [
"TIER_0",
"TIER_1",
"TIER_2",
"UNKNOWN",
"ALL"
],
"type": "string"
},
"exclude_pinned": {
"default": false,
"description": "If True, exclude syncs for actors that are already pinned to a specific version (at any scope level: actor, workspace, or organization). Useful for 'prove fix' workflows where you want to find unpinned connections for live testing. Default: False (include all syncs).",
"type": "boolean"
},
"enabled_schedules_only": {
"default": false,
"description": "If True, only return syncs for connections that are both active (not paused/inactive) and on an automated sync schedule (not manual-trigger-only). Useful for canary workflows where you need connections that will produce organic syncs during a monitoring window. Default: False (include all connections).",
"type": "boolean"
}
},
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_recent_syncs_for_connector_version
List sync jobs that were run with a specific connector version.
Works for both source and destination connectors. Automatically detects the connector type from the version metadata and uses the appropriate query variant.
Accepts either connector_version_id (UUID) or connector_name +
connector_version (e.g. source-pokeapi + 0.3.59). When using
name + version, the docker_repository is derived from the canonical
name (e.g. source-pokeapi → airbyte/source-pokeapi).
Filters on the version stamped into jobs.config at job-creation time,
not the current pin state. This avoids false positives (pre-pin syncs
counted as RC) and false negatives (post-unpin syncs missed).
Pin columns (pin_origin_type, pin_origin, pin_scope_type) are
still included as informational output but are not used for filtering.
Returns list of dicts with keys: job_id, connection_id, job_status,
started_at, job_updated_at, connection_name, actor_id, actor_name,
actor_definition_id, source_definition_version_id,
destination_definition_version_id, pin_origin_type,
pin_origin, pin_scope_type, workspace_id, workspace_name,
organization_id, dataplane_group_id, dataplane_name.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
connector_version_id |
string | null |
no | null |
Connector version UUID. Provide this OR connector_name + connector_version. |
connector_name |
string | null |
no | null |
Canonical connector name (e.g. source-pokeapi, destination-duckdb). Used with connector_version to resolve the version UUID. |
connector_version |
string | null |
no | null |
Semver version tag (e.g. 0.3.59). Used with connector_name to resolve the version UUID. |
days |
integer |
no | 7 |
Number of days to look back (default: 7) |
limit |
integer |
no | 100 |
Maximum number of results (default: 100) |
successful_only |
boolean |
no | false |
If True, only return successful syncs (default: False) |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"connector_version_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Connector version UUID. Provide this OR connector_name + connector_version."
},
"connector_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Canonical connector name (e.g. `source-pokeapi`, `destination-duckdb`). Used with `connector_version` to resolve the version UUID."
},
"connector_version": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Semver version tag (e.g. `0.3.59`). Used with `connector_name` to resolve the version UUID."
},
"days": {
"default": 7,
"description": "Number of days to look back (default: 7)",
"type": "integer"
},
"limit": {
"default": 100,
"description": "Maximum number of results (default: 100)",
"type": "integer"
},
"successful_only": {
"default": false,
"description": "If `True`, only return successful syncs (default: `False`)",
"type": "boolean"
}
},
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"items": {
"additionalProperties": true,
"type": "object"
},
"type": "array"
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_workspace_info
Get workspace information including dataplane group.
Returns details about a specific workspace, including which dataplane (region) it belongs to. Useful for determining if a workspace is in the EU region for filtering purposes.
Returns dict with keys: workspace_id, workspace_name, slug, organization_id, dataplane_group_id, dataplane_name, created_at, tombstone Or None if workspace not found.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
workspace_id |
string | enum("266ebdfe-0d7b-4540-9817-de7e4505ba61") |
yes | — | Workspace UUID or alias to look up. Accepts '@devin-ai-sandbox' as an alias for the Devin AI sandbox workspace. |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"workspace_id": {
"anyOf": [
{
"type": "string"
},
{
"description": "Workspace ID aliases that can be used in place of UUIDs.\n\nEach member's name is the alias (e.g., \"@devin-ai-sandbox\") and its value\nis the actual workspace UUID. Use `WorkspaceAliasEnum.resolve()` to\nresolve aliases to actual IDs.",
"enum": [
"266ebdfe-0d7b-4540-9817-de7e4505ba61"
],
"type": "string"
}
],
"description": "Workspace UUID or alias to look up. Accepts '@devin-ai-sandbox' as an alias for the Devin AI sandbox workspace."
}
},
"required": [
"workspace_id"
],
"type": "object"
}
Show output JSON schema
{
"properties": {
"result": {
"anyOf": [
{
"additionalProperties": true,
"type": "object"
},
{
"type": "null"
}
]
}
},
"required": [
"result"
],
"type": "object",
"x-fastmcp-wrap-result": true
}
query_prod_workspaces
Search workspaces by name substring or email domain.
At least one of name_contains or email_domain must be provided.
When name_contains is given, performs a case-insensitive substring match
on workspace name and slug. When email_domain is given, matches
workspaces by user email domain.
The returned organization IDs can be used with other tools like
query_prod_connections_by_connector to find connections within
those organizations for safe testing.
Parameters:
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
name_contains |
string | null |
no | null |
Case-insensitive substring to search for in workspace name or slug. For example, 'acme' will match workspaces named 'Acme Staging'. |
email_domain |
string | null |
no | null |
Email domain to search for (e.g., 'motherduck.com'). Do not include the '@' symbol. |
limit |
integer |
no | 100 |
Maximum number of workspaces to return (default: 100) |
Show input JSON schema
{
"additionalProperties": false,
"properties": {
"name_contains": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Case-insensitive substring to search for in workspace name or slug. For example, 'acme' will match workspaces named 'Acme Staging'."
},
"email_domain": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Email domain to search for (e.g., 'motherduck.com'). Do not include the '@' symbol."
},
"limit": {
"default": 100,
"description": "Maximum number of workspaces to return (default: 100)",
"type": "integer"
}
},
"type": "object"
}
Show output JSON schema
{
"description": "Result of searching workspaces by name or email domain.",
"properties": {
"name_contains": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The name substring that was searched for"
},
"email_domain": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The email domain that was searched for (e.g., 'motherduck.com')"
},
"total_workspaces_found": {
"description": "Total number of workspaces matching",
"type": "integer"
},
"unique_organization_ids": {
"description": "List of unique organization IDs found",
"items": {
"type": "string"
},
"type": "array"
},
"workspaces": {
"description": "List of matching workspaces",
"items": {
"description": "Information about a workspace.",
"properties": {
"organization_id": {
"description": "The organization UUID",
"type": "string"
},
"workspace_id": {
"description": "The workspace UUID",
"type": "string"
},
"workspace_name": {
"description": "The name of the workspace",
"type": "string"
},
"slug": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The workspace slug (URL-friendly identifier)"
},
"email": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The email address associated with the workspace"
},
"dataplane_group_id": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The dataplane group UUID (region)"
},
"dataplane_name": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "The name of the dataplane (e.g., 'US', 'EU')"
},
"created_at": {
"anyOf": [
{
"format": "date-time",
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "When the workspace was created"
},
"customer_tier": {
"anyOf": [
{
"type": "string"
},
{
"type": "null"
}
],
"default": null,
"description": "Customer tier (TIER_0, TIER_1, TIER_2, or UNKNOWN). Enriched from the GCS tier cache."
},
"is_eu": {
"anyOf": [
{
"type": "boolean"
},
{
"type": "null"
}
],
"default": null,
"description": "Whether the workspace is in the EU region (derived from dataplane_name)."
}
},
"required": [
"organization_id",
"workspace_id",
"workspace_name"
],
"type": "object"
},
"type": "array"
}
},
"required": [
"total_workspaces_found",
"unique_organization_ids",
"workspaces"
],
"type": "object"
}
1# Copyright (c) 2025 Airbyte, Inc., all rights reserved. 2"""MCP tools for querying the Airbyte Cloud Prod DB Replica. 3 4This module provides MCP tools that wrap the query functions from 5airbyte_ops_mcp.prod_db_access.queries for use by AI agents. 6 7## MCP reference 8 9.. include:: ../../../docs/mcp-generated/prod_db_ops.md 10 :start-line: 2 11""" 12 13from __future__ import annotations 14 15__all__: list[str] = [] 16 17import json 18import uuid 19from datetime import datetime, timedelta, timezone 20from enum import StrEnum 21from typing import Annotated, Any 22 23from airbyte.exceptions import PyAirbyteInputError 24from fastmcp import FastMCP 25from fastmcp_extensions import mcp_tool, register_mcp_tools 26from pydantic import BaseModel, Field 27 28from airbyte_ops_mcp.cloud_admin.registry_lookup import ( 29 resolve_canonical_name_to_definition_id, 30) 31from airbyte_ops_mcp.constants import OrganizationAliasEnum, WorkspaceAliasEnum 32from airbyte_ops_mcp.prod_db_access.queries import ( 33 is_source_connector, 34 query_actor_population_by_org, 35 query_actors_pinned_to_version, 36 query_connection_sync_activity_from_prod, 37 query_connections_by_connector, 38 query_connections_by_destination_connector, 39 query_connections_by_stream, 40 query_connector_rollouts, 41 query_connector_versions, 42 query_dataplanes_list, 43 query_destination_connection_stats, 44 query_failed_sync_attempts_for_connector, 45 query_new_connector_releases, 46 query_org_admin_contacts, 47 query_org_connector_pins, 48 query_org_pin_stats, 49 query_org_user_last_connection_events, 50 query_recent_syncs_for_connector, 51 query_source_connection_stats, 52 query_syncs_for_connector_version, 53 query_version_actor_health, 54 query_versions_with_pins, 55 query_workspace_info, 56 query_workspaces_by_email_domain, 57 resolve_version_id_by_tag, 58 resolve_version_info, 59 search_organizations, 60 search_workspaces, 61) 62from airbyte_ops_mcp.tier_cache import ( 63 TierFilter, 64 TierSummary, 65 enrich_rows_by_org, 66 filter_rows_by_tier, 67 get_org_tiers, 68 tier_source_warnings, 69) 70from airbyte_ops_mcp.version_summaries import ( 71 summarize_population, 72 summarize_version_health, 73) 74 75 76class StatusFilter(StrEnum): 77 """Filter for job status in sync queries.""" 78 79 ALL = "all" 80 SUCCEEDED = "succeeded" 81 FAILED = "failed" 82 83 84# Cloud UI base URL for building connection URLs 85CLOUD_UI_BASE_URL = "https://cloud.airbyte.com" 86 87 88def _validate_sync_activity_scope( 89 *, 90 organization_id: str | None, 91 workspace_id: str | None, 92 connection_ids: list[str] | None, 93) -> None: 94 """Require at least one explicit scope filter for sync activity queries.""" 95 if organization_id or workspace_id or connection_ids: 96 return 97 raise PyAirbyteInputError( 98 message=( 99 "Provide at least one scope filter: `organization_id`, `workspace_id`, " 100 "or `connection_ids`." 101 ), 102 context={ 103 "organization_id": organization_id, 104 "workspace_id": workspace_id, 105 "connection_ids": connection_ids, 106 }, 107 ) 108 109 110def _validate_sync_activity_window( 111 *, 112 start_at: datetime, 113 end_at: datetime, 114) -> tuple[datetime, datetime]: 115 """Validate that `start_at` and `end_at` describe a usable window. 116 117 Returns the timestamps normalized to UTC. Raises `PyAirbyteInputError` for 118 naive timestamps or inverted ranges. No clock-relative caps are enforced 119 here; the caller is trusted to choose a sensible window. 120 """ 121 if start_at.tzinfo is None or end_at.tzinfo is None: 122 raise PyAirbyteInputError( 123 message="`start_at` and `end_at` must include timezone information.", 124 context={ 125 "start_at": start_at.isoformat(), 126 "end_at": end_at.isoformat(), 127 }, 128 ) 129 130 normalized_start = start_at.astimezone(timezone.utc) 131 normalized_end = end_at.astimezone(timezone.utc) 132 133 if normalized_start >= normalized_end: 134 raise PyAirbyteInputError( 135 message="`start_at` must be earlier than `end_at`.", 136 context={ 137 "start_at": normalized_start.isoformat(), 138 "end_at": normalized_end.isoformat(), 139 }, 140 ) 141 return normalized_start, normalized_end 142 143 144# ============================================================================= 145# Pydantic Models for MCP Tool Responses 146# ============================================================================= 147 148 149class OrganizationSearchHit(BaseModel): 150 """A single organization returned by a name/email search.""" 151 152 organization_id: str = Field(description="The organization UUID") 153 organization_name: str = Field(description="The name of the organization") 154 email: str | None = Field( 155 default=None, description="The email address associated with the organization" 156 ) 157 customer_tier: str | None = Field( 158 default=None, 159 description="Customer tier (TIER_0, TIER_1, TIER_2, or UNKNOWN). Enriched from the GCS tier cache.", 160 ) 161 162 163class OrganizationSearchResult(BaseModel): 164 """Result of searching organizations by name substring.""" 165 166 name_contains: str = Field(description="The search substring that was used") 167 total_found: int = Field(description="Total number of organizations matching") 168 organizations: list[OrganizationSearchHit] = Field( 169 description="List of matching organizations" 170 ) 171 172 173class WorkspaceInfo(BaseModel): 174 """Information about a workspace.""" 175 176 organization_id: str = Field(description="The organization UUID") 177 workspace_id: str = Field(description="The workspace UUID") 178 workspace_name: str = Field(description="The name of the workspace") 179 slug: str | None = Field( 180 default=None, description="The workspace slug (URL-friendly identifier)" 181 ) 182 email: str | None = Field( 183 default=None, description="The email address associated with the workspace" 184 ) 185 dataplane_group_id: str | None = Field( 186 default=None, description="The dataplane group UUID (region)" 187 ) 188 dataplane_name: str | None = Field( 189 default=None, description="The name of the dataplane (e.g., 'US', 'EU')" 190 ) 191 created_at: datetime | None = Field( 192 default=None, description="When the workspace was created" 193 ) 194 customer_tier: str | None = Field( 195 default=None, 196 description="Customer tier (TIER_0, TIER_1, TIER_2, or UNKNOWN). Enriched from the GCS tier cache.", 197 ) 198 is_eu: bool | None = Field( 199 default=None, 200 description="Whether the workspace is in the EU region (derived from dataplane_name).", 201 ) 202 203 204class WorkspaceSearchResult(BaseModel): 205 """Result of searching workspaces by name or email domain.""" 206 207 name_contains: str | None = Field( 208 default=None, description="The name substring that was searched for" 209 ) 210 email_domain: str | None = Field( 211 default=None, 212 description="The email domain that was searched for (e.g., 'motherduck.com')", 213 ) 214 total_workspaces_found: int = Field( 215 description="Total number of workspaces matching" 216 ) 217 unique_organization_ids: list[str] = Field( 218 description="List of unique organization IDs found" 219 ) 220 workspaces: list[WorkspaceInfo] = Field(description="List of matching workspaces") 221 222 223# Keep backward-compatible alias for any external references 224WorkspacesByEmailDomainResult = WorkspaceSearchResult 225 226 227class LatestAttemptBreakdown(BaseModel): 228 """Breakdown of connections by latest attempt status.""" 229 230 succeeded: int = Field( 231 default=0, description="Connections where latest attempt succeeded" 232 ) 233 failed: int = Field( 234 default=0, description="Connections where latest attempt failed" 235 ) 236 cancelled: int = Field( 237 default=0, description="Connections where latest attempt was cancelled" 238 ) 239 running: int = Field( 240 default=0, description="Connections where latest attempt is still running" 241 ) 242 unknown: int = Field( 243 default=0, 244 description="Connections with no recent attempts in the lookback window", 245 ) 246 247 248class VersionPinStats(BaseModel): 249 """Stats for connections pinned to a specific version.""" 250 251 pinned_version_id: str | None = Field( 252 description="The connector version UUID (None for unpinned connections)" 253 ) 254 docker_image_tag: str | None = Field( 255 default=None, description="The docker image tag for this version" 256 ) 257 total_connections: int = Field(description="Total number of connections") 258 enabled_connections: int = Field( 259 description="Number of enabled (active status) connections" 260 ) 261 active_connections: int = Field( 262 description="Number of connections with recent sync activity" 263 ) 264 latest_attempt: LatestAttemptBreakdown = Field( 265 description="Breakdown by latest attempt status" 266 ) 267 268 269class ConnectorConnectionStats(BaseModel): 270 """Aggregate connection stats for a connector.""" 271 272 connector_definition_id: str = Field(description="The connector definition UUID") 273 connector_type: str = Field(description="'source' or 'destination'") 274 canonical_name: str | None = Field( 275 default=None, description="The canonical connector name if resolved" 276 ) 277 total_connections: int = Field( 278 description="Total number of non-deprecated connections" 279 ) 280 enabled_connections: int = Field( 281 description="Number of enabled (active status) connections" 282 ) 283 active_connections: int = Field( 284 description="Number of connections with recent sync activity" 285 ) 286 pinned_connections: int = Field( 287 description="Number of connections with explicit version pins" 288 ) 289 unpinned_connections: int = Field( 290 description="Number of connections on default version" 291 ) 292 latest_attempt: LatestAttemptBreakdown = Field( 293 description="Overall breakdown by latest attempt status" 294 ) 295 by_version: list[VersionPinStats] = Field( 296 description="Stats broken down by pinned version" 297 ) 298 299 300class ConnectorConnectionStatsResponse(BaseModel): 301 """Response containing connection stats for multiple connectors.""" 302 303 sources: list[ConnectorConnectionStats] = Field( 304 default_factory=list, description="Stats for source connectors" 305 ) 306 destinations: list[ConnectorConnectionStats] = Field( 307 default_factory=list, description="Stats for destination connectors" 308 ) 309 lookback_days: int = Field( 310 description="Lookback window used for 'active' connections" 311 ) 312 generated_at: datetime = Field(description="When this response was generated") 313 314 315class OrgAdminContact(BaseModel): 316 """An organization or selected-workspace administrator contact.""" 317 318 user_id: str = Field(description="The Airbyte user UUID.") 319 name: str | None = Field(default=None, description="The administrator's name.") 320 email: str | None = Field( 321 default=None, description="The administrator's email address." 322 ) 323 status: str | None = Field(default=None, description="The user's status.") 324 is_org_admin: bool = Field(description="Whether the user is an Org Admin.") 325 admin_workspace_ids: list[str] = Field( 326 description="Sorted workspace UUIDs where the user is a Workspace Admin." 327 ) 328 user_created_at: datetime | None = Field( 329 default=None, description="When the user account was created." 330 ) 331 user_updated_at: datetime | None = Field( 332 default=None, description="When the user account was last updated." 333 ) 334 last_connection_event_at: datetime | None = Field( 335 default=None, 336 description="Latest user-attributed connection timeline event in the activity window.", 337 ) 338 339 340class OrgAdminContactsResponse(BaseModel): 341 """Organization administrator contacts and recent user activity.""" 342 343 organization_id: str = Field(description="The organization UUID.") 344 customer_tier: str = Field( 345 description="Organization customer tier, included for awareness only." 346 ) 347 tier_warnings: list[str] = Field( 348 default_factory=list, 349 description="Warnings raised while resolving the organization customer tier.", 350 ) 351 activity_lookback_days: int = Field( 352 description="Number of days used to look back for connection activity." 353 ) 354 admins: list[OrgAdminContact] = Field( 355 description="Administrators ordered by Airbyte user UUID." 356 ) 357 358 359def _opt_str(value: Any) -> str | None: 360 """Convert a nullable value to str, returning None if the value is None/falsy.""" 361 return str(value) if value else None 362 363 364@mcp_tool( 365 read_only=True, 366 idempotent=True, 367) 368def query_prod_org_admin_contacts( 369 organization_id: Annotated[ 370 str, 371 Field(description="Organization UUID whose direct Org Admins to find."), 372 ], 373 workspace_ids: Annotated[ 374 list[str], 375 Field( 376 description=( 377 "Workspace UUIDs whose direct Workspace Admins to include. " 378 "Only these workspaces contribute Workspace Admins." 379 ), 380 ), 381 ] = [], # noqa: B006 382 activity_lookback_days: Annotated[ 383 int, 384 Field( 385 description="Days to look back for user-attributed connection activity.", 386 ge=1, 387 le=365, 388 ), 389 ] = 90, 390) -> OrgAdminContactsResponse: 391 """Find Org Admins and Workspace Admins of the given workspaces. 392 393 Activity is each user's latest user-attributed connection timeline event in 394 this organization within the lookback window. Emails are returned because 395 callers need them to CC customers; never log them. Only direct user grants 396 are included; group-granted permissions are not queried. The org's 397 `customer_tier` is included for awareness only; outreach must still reach 398 every affected organization regardless of tier. 399 """ 400 normalized_organization_id = _require_organization_id(organization_id) 401 normalized_workspace_ids: list[str] = [] 402 for workspace_id in workspace_ids: 403 try: 404 normalized_workspace_ids.append(str(uuid.UUID(workspace_id.strip()))) 405 except (AttributeError, ValueError) as exc: 406 raise PyAirbyteInputError( 407 message=f"`workspace_ids` contains an invalid workspace UUID: {workspace_id!r}.", 408 ) from exc 409 410 tier_result = get_org_tiers( 411 organization_ids=[normalized_organization_id], 412 allow_degraded=True, 413 )[0] 414 admin_rows = query_org_admin_contacts( 415 normalized_organization_id, 416 normalized_workspace_ids, 417 ) 418 admins_by_user_id: dict[str, OrgAdminContact] = {} 419 for row in admin_rows: 420 user_id = str(row["user_id"]) 421 admin = admins_by_user_id.get(user_id) 422 if admin is None: 423 admin = OrgAdminContact( 424 user_id=user_id, 425 name=row["user_name"], 426 email=row["user_email"], 427 status=row["user_status"], 428 is_org_admin=False, 429 admin_workspace_ids=[], 430 user_created_at=row["user_created_at"], 431 user_updated_at=row["user_updated_at"], 432 ) 433 admins_by_user_id[user_id] = admin 434 435 if row["admin_role"] == "organization_admin": 436 admin.is_org_admin = True 437 elif row["admin_role"] == "workspace_admin" and row["workspace_id"] is not None: 438 workspace_id = str(row["workspace_id"]) 439 if workspace_id not in admin.admin_workspace_ids: 440 admin.admin_workspace_ids.append(workspace_id) 441 442 if admins_by_user_id: 443 activity_rows = query_org_user_last_connection_events( 444 normalized_organization_id, 445 sorted(admins_by_user_id), 446 datetime.now(timezone.utc) - timedelta(days=activity_lookback_days), 447 ) 448 for row in activity_rows: 449 user_id = str(row["user_id"]) 450 if user_id in admins_by_user_id: 451 admins_by_user_id[user_id].last_connection_event_at = row[ 452 "last_connection_event_at" 453 ] 454 455 admins = [admins_by_user_id[user_id] for user_id in sorted(admins_by_user_id)] 456 for admin in admins: 457 admin.admin_workspace_ids.sort() 458 return OrgAdminContactsResponse( 459 organization_id=normalized_organization_id, 460 customer_tier=str(tier_result.customer_tier), 461 tier_warnings=tier_source_warnings(tier_result.source_health), 462 activity_lookback_days=activity_lookback_days, 463 admins=admins, 464 ) 465 466 467@mcp_tool( 468 read_only=True, 469 idempotent=True, 470) 471def query_prod_dataplanes() -> list[dict[str, Any]]: 472 """List all dataplane groups with workspace counts. 473 474 Returns information about all active dataplane groups in Airbyte Cloud, 475 including the number of workspaces in each. Useful for understanding 476 the distribution of workspaces across regions (US, US-Central, EU). 477 478 Returns list of dicts with keys: dataplane_group_id, dataplane_name, 479 organization_id, enabled, tombstone, created_at, workspace_count 480 """ 481 return query_dataplanes_list() 482 483 484@mcp_tool( 485 read_only=True, 486 idempotent=True, 487) 488def query_prod_workspace_info( 489 workspace_id: Annotated[ 490 str | WorkspaceAliasEnum, 491 Field( 492 description="Workspace UUID or alias to look up. " 493 "Accepts '@devin-ai-sandbox' as an alias for the Devin AI sandbox workspace." 494 ), 495 ], 496) -> dict[str, Any] | None: 497 """Get workspace information including dataplane group. 498 499 Returns details about a specific workspace, including which dataplane 500 (region) it belongs to. Useful for determining if a workspace is in 501 the EU region for filtering purposes. 502 503 Returns dict with keys: workspace_id, workspace_name, slug, organization_id, 504 dataplane_group_id, dataplane_name, created_at, tombstone 505 Or None if workspace not found. 506 """ 507 # Resolve workspace ID alias (workspace_id is required, so resolved value is never None) 508 resolved_workspace_id = WorkspaceAliasEnum.resolve(workspace_id) 509 assert resolved_workspace_id is not None # Type narrowing: workspace_id is required 510 511 return query_workspace_info(resolved_workspace_id) 512 513 514@mcp_tool( 515 read_only=True, 516 idempotent=True, 517) 518def query_prod_connector_versions( 519 connector_definition_id: Annotated[ 520 str, 521 Field(description="Connector definition UUID to list versions for"), 522 ], 523) -> list[dict[str, Any]]: 524 """List all versions for a connector definition. 525 526 Returns all published versions of a connector, ordered by last_published 527 date descending. Useful for understanding version history and finding 528 specific version IDs for pinning or rollout monitoring. 529 530 Returns list of dicts with keys: version_id, docker_image_tag, docker_repository, 531 release_stage, support_level, cdk_version, language, last_published, release_date 532 """ 533 return query_connector_versions(connector_definition_id) 534 535 536@mcp_tool( 537 read_only=True, 538 idempotent=True, 539) 540def query_prod_new_connector_releases( 541 days: Annotated[ 542 int, 543 Field(description="Number of days to look back (default: 7)", default=7), 544 ] = 7, 545 limit: Annotated[ 546 int, 547 Field(description="Maximum number of results (default: 100)", default=100), 548 ] = 100, 549) -> list[dict[str, Any]]: 550 """List recently published connector versions. 551 552 Returns connector versions published within the specified number of days. 553 Uses last_published timestamp which reflects when the version was actually 554 deployed to the registry (not the changelog date). 555 556 Returns list of dicts with keys: version_id, connector_definition_id, docker_repository, 557 docker_image_tag, last_published, release_date, release_stage, support_level, 558 cdk_version, language, created_at 559 """ 560 return query_new_connector_releases(days=days, limit=limit) 561 562 563@mcp_tool( 564 read_only=True, 565 idempotent=True, 566) 567def query_prod_actors_by_pinned_connector_version( 568 connector_version_id: Annotated[ 569 str, 570 Field(description="Connector version UUID to find pinned instances for"), 571 ], 572) -> list[dict[str, Any]]: 573 """List actors (sources/destinations) effectively pinned to a specific connector version. 574 575 Returns all actors that are effectively pinned to a specific connector version, 576 considering all scope levels: actor-level pins, workspace-level pins, and 577 organization-level pins (with actor > workspace > organization precedence). 578 Useful for monitoring rollouts and understanding which customers are affected. 579 580 The actor_id field is the actor ID (superset of source_id/destination_id). 581 582 Returns list of dicts with keys: actor_id, connector_definition_id, origin_type, 583 origin, description, created_at, expires_at, pin_scope_type, actor_name, 584 workspace_id, workspace_name, organization_id, dataplane_group_id, dataplane_name 585 586 pin_scope_type is 'actor', 'workspace', or 'organization' indicating which scope 587 level the effective pin came from. 588 """ 589 return query_actors_pinned_to_version(connector_version_id) 590 591 592@mcp_tool( 593 read_only=True, 594 idempotent=True, 595) 596def query_prod_recent_syncs_for_connector_version( 597 connector_version_id: Annotated[ 598 str | None, 599 Field( 600 description=( 601 "Connector version UUID. Provide this OR " 602 "connector_name + connector_version." 603 ), 604 default=None, 605 ), 606 ] = None, 607 connector_name: Annotated[ 608 str | None, 609 Field( 610 description=( 611 "Canonical connector name (e.g. `source-pokeapi`, " 612 "`destination-duckdb`). Used with `connector_version` to " 613 "resolve the version UUID." 614 ), 615 default=None, 616 ), 617 ] = None, 618 connector_version: Annotated[ 619 str | None, 620 Field( 621 description=( 622 "Semver version tag (e.g. `0.3.59`). " 623 "Used with `connector_name` to resolve the version UUID." 624 ), 625 default=None, 626 ), 627 ] = None, 628 days: Annotated[ 629 int, 630 Field(description="Number of days to look back (default: 7)", default=7), 631 ] = 7, 632 limit: Annotated[ 633 int, 634 Field(description="Maximum number of results (default: 100)", default=100), 635 ] = 100, 636 successful_only: Annotated[ 637 bool, 638 Field( 639 description="If `True`, only return successful syncs (default: `False`)", 640 default=False, 641 ), 642 ] = False, 643) -> list[dict[str, Any]]: 644 """List sync jobs that were run with a specific connector version. 645 646 Works for both source and destination connectors. Automatically detects 647 the connector type from the version metadata and uses the appropriate 648 query variant. 649 650 Accepts either `connector_version_id` (UUID) or `connector_name` + 651 `connector_version` (e.g. `source-pokeapi` + `0.3.59`). When using 652 name + version, the `docker_repository` is derived from the canonical 653 name (e.g. `source-pokeapi` → `airbyte/source-pokeapi`). 654 655 Filters on the version stamped into `jobs.config` at job-creation time, 656 not the current pin state. This avoids false positives (pre-pin syncs 657 counted as RC) and false negatives (post-unpin syncs missed). 658 659 Pin columns (`pin_origin_type`, `pin_origin`, `pin_scope_type`) are 660 still included as informational output but are not used for filtering. 661 662 Returns list of dicts with keys: `job_id`, `connection_id`, `job_status`, 663 `started_at`, `job_updated_at`, `connection_name`, `actor_id`, `actor_name`, 664 `actor_definition_id`, `source_definition_version_id`, 665 `destination_definition_version_id`, `pin_origin_type`, 666 `pin_origin`, `pin_scope_type`, `workspace_id`, `workspace_name`, 667 `organization_id`, `dataplane_group_id`, `dataplane_name`. 668 """ 669 # Resolve inputs to a version UUID and connector type. 670 if connector_version_id is not None: 671 version_info = resolve_version_info(connector_version_id) 672 docker_repository = version_info["docker_repository"] 673 elif connector_name is not None and connector_version is not None: 674 # Derive docker_repository from canonical name. 675 docker_repository = f"airbyte/{connector_name}" 676 version_info = resolve_version_id_by_tag( 677 docker_repository=docker_repository, 678 docker_image_tag=connector_version, 679 ) 680 connector_version_id = version_info["version_id"] 681 else: 682 raise PyAirbyteInputError( 683 message=( 684 "Provide either `connector_version_id` or both " 685 "`connector_name` and `connector_version`." 686 ), 687 ) 688 689 is_destination = not is_source_connector(docker_repository) 690 return query_syncs_for_connector_version( 691 connector_version_id, 692 is_destination=is_destination, 693 days=days, 694 limit=limit, 695 successful_only=successful_only, 696 ) 697 698 699@mcp_tool( 700 read_only=True, 701 idempotent=True, 702 open_world=True, 703) 704def query_prod_recent_syncs_for_connector( 705 source_definition_id: Annotated[ 706 str | None, 707 Field( 708 description=( 709 "Source connector definition ID (UUID) to search for. " 710 "Provide this OR source_canonical_name OR destination_definition_id " 711 "OR destination_canonical_name (exactly one required). " 712 "Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics." 713 ), 714 default=None, 715 ), 716 ], 717 source_canonical_name: Annotated[ 718 str | None, 719 Field( 720 description=( 721 "Canonical source connector name to search for. " 722 "Provide this OR source_definition_id OR destination_definition_id " 723 "OR destination_canonical_name (exactly one required). " 724 "Examples: 'source-youtube-analytics', 'YouTube Analytics'." 725 ), 726 default=None, 727 ), 728 ], 729 destination_definition_id: Annotated[ 730 str | None, 731 Field( 732 description=( 733 "Destination connector definition ID (UUID) to search for. " 734 "Provide this OR destination_canonical_name OR source_definition_id " 735 "OR source_canonical_name (exactly one required). " 736 "Example: '94bd199c-2ff0-4aa2-b98e-17f0acb72610' for DuckDB." 737 ), 738 default=None, 739 ), 740 ], 741 destination_canonical_name: Annotated[ 742 str | None, 743 Field( 744 description=( 745 "Canonical destination connector name to search for. " 746 "Provide this OR destination_definition_id OR source_definition_id " 747 "OR source_canonical_name (exactly one required). " 748 "Examples: 'destination-duckdb', 'DuckDB'." 749 ), 750 default=None, 751 ), 752 ], 753 status_filter: Annotated[ 754 StatusFilter, 755 Field( 756 description=( 757 "Filter by job status: 'all' (default), 'succeeded', or 'failed'. " 758 "Use 'succeeded' to find healthy connections with recent successful syncs. " 759 "Use 'failed' to find connections with recent failures." 760 ), 761 default=StatusFilter.ALL, 762 ), 763 ], 764 organization_id: Annotated[ 765 str | OrganizationAliasEnum | None, 766 Field( 767 description=( 768 "Optional organization ID (UUID) or alias to filter results. " 769 "If provided, only syncs from this organization will be returned. " 770 "Accepts '@airbyte-internal' as an alias for the Airbyte internal org." 771 ), 772 default=None, 773 ), 774 ], 775 lookback_days: Annotated[ 776 int, 777 Field(description="Number of days to look back (default: 7)", default=7), 778 ], 779 limit: Annotated[ 780 int, 781 Field(description="Maximum number of results (default: 100)", default=100), 782 ], 783 customer_tier_filter: Annotated[ 784 TierFilter, 785 Field( 786 description=( 787 "Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. " 788 "Filters results to only include connections belonging to organizations " 789 "in the specified tier. Use 'ALL' to include all tiers." 790 ), 791 ), 792 ] = "TIER_2", 793 *, 794 exclude_pinned: Annotated[ 795 bool, 796 Field( 797 description=( 798 "If True, exclude syncs for actors that are already pinned to a " 799 "specific version (at any scope level: actor, workspace, or organization). " 800 "Useful for 'prove fix' workflows where you want to find unpinned " 801 "connections for live testing. Default: False (include all syncs)." 802 ), 803 default=False, 804 ), 805 ], 806 enabled_schedules_only: Annotated[ 807 bool, 808 Field( 809 description=( 810 "If True, only return syncs for connections that are both active " 811 "(not paused/inactive) and on an automated sync schedule " 812 "(not manual-trigger-only). Useful for canary workflows where " 813 "you need connections that will produce organic syncs during a " 814 "monitoring window. Default: False (include all connections)." 815 ), 816 default=False, 817 ), 818 ], 819) -> list[dict[str, Any]]: 820 """List recent sync jobs for ALL actors using a connector type. 821 822 This tool finds all actors with the given connector definition and returns their 823 recent sync jobs, regardless of whether they have explicit version pins. It filters 824 out deleted actors, deleted workspaces, and deprecated connections. 825 826 Results are always enriched with customer_tier and is_eu fields. 827 The customer_tier_filter parameter is required to ensure tier-aware querying. 828 829 Use this tool to: 830 - Find healthy connections with recent successful syncs (status_filter='succeeded') 831 - Investigate connector issues across all users (status_filter='failed') 832 - Get an overview of all recent sync activity (status_filter='all') 833 834 Set `exclude_pinned=True` to filter out syncs for actors that are already pinned to a 835 specific version. This is useful for 'prove fix' live connection testing workflows 836 where you want to find unpinned connections to test against. 837 838 Set `enabled_schedules_only=True` to restrict results to connections that are both 839 enabled (status='active') and on an automated schedule (not manual-trigger-only). 840 This is useful for canary prerelease workflows where you need connections that 841 will run organically during the monitoring window. 842 843 Supports both SOURCE and DESTINATION connectors. Provide exactly one of: 844 source_definition_id, source_canonical_name, destination_definition_id, 845 or destination_canonical_name. 846 847 Key fields in results: 848 - job_status: 'succeeded', 'failed', 'cancelled', etc. 849 - connection_id, connection_name: The connection that ran the sync 850 - actor_id, actor_name: The source or destination actor 851 - customer_tier: TIER_0, TIER_1, TIER_2, or UNKNOWN 852 - is_eu: Whether the workspace is in the EU region 853 - pin_origin_type, pin_origin, pinned_version_id: Version pin context (NULL if not pinned) 854 - pin_scope_type: 'actor', 'workspace', or 'organization' (NULL if not pinned) 855 """ 856 # Validate that exactly one connector parameter is provided 857 provided_params = [ 858 source_definition_id, 859 source_canonical_name, 860 destination_definition_id, 861 destination_canonical_name, 862 ] 863 num_provided = sum(p is not None for p in provided_params) 864 if num_provided != 1: 865 raise PyAirbyteInputError( 866 message=( 867 "Exactly one of source_definition_id, source_canonical_name, " 868 "destination_definition_id, or destination_canonical_name must be provided." 869 ), 870 ) 871 872 # Determine if this is a destination connector 873 is_destination = ( 874 destination_definition_id is not None or destination_canonical_name is not None 875 ) 876 877 # Resolve canonical name to definition ID if needed 878 resolved_definition_id: str 879 if source_canonical_name: 880 resolved_definition_id = resolve_canonical_name_to_definition_id( 881 canonical_name=source_canonical_name, 882 ) 883 elif destination_canonical_name: 884 resolved_definition_id = resolve_canonical_name_to_definition_id( 885 canonical_name=destination_canonical_name, 886 ) 887 elif source_definition_id: 888 resolved_definition_id = source_definition_id 889 else: 890 # We've validated exactly one param is provided, so this must be set 891 assert destination_definition_id is not None 892 resolved_definition_id = destination_definition_id 893 894 # Resolve organization ID alias 895 resolved_organization_id = OrganizationAliasEnum.resolve(organization_id) 896 897 rows = query_recent_syncs_for_connector( 898 connector_definition_id=resolved_definition_id, 899 is_destination=is_destination, 900 status_filter=status_filter, 901 organization_id=resolved_organization_id, 902 days=lookback_days, 903 limit=limit, 904 exclude_pinned=exclude_pinned, 905 enabled_schedules_only=enabled_schedules_only, 906 ) 907 908 enriched = enrich_rows_by_org( 909 rows=rows, 910 allow_degraded=True, 911 ) 912 return filter_rows_by_tier(enriched, customer_tier_filter) 913 914 915@mcp_tool( 916 read_only=True, 917 idempotent=True, 918 open_world=True, 919) 920def query_prod_failed_sync_attempts_for_connector( 921 source_definition_id: Annotated[ 922 str | None, 923 Field( 924 description=( 925 "Source connector definition ID (UUID) to search for. " 926 "Provide this OR source_canonical_name OR destination_definition_id " 927 "OR destination_canonical_name (exactly one required). " 928 "Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics." 929 ), 930 default=None, 931 ), 932 ] = None, 933 source_canonical_name: Annotated[ 934 str | None, 935 Field( 936 description=( 937 "Canonical source connector name to search for. " 938 "Provide this OR source_definition_id OR destination_definition_id " 939 "OR destination_canonical_name (exactly one required). " 940 "Examples: 'source-youtube-analytics', 'YouTube Analytics'." 941 ), 942 default=None, 943 ), 944 ] = None, 945 destination_definition_id: Annotated[ 946 str | None, 947 Field( 948 description=( 949 "Destination connector definition ID (UUID) to search for. " 950 "Provide this OR destination_canonical_name OR source_definition_id " 951 "OR source_canonical_name (exactly one required). " 952 "Example: '94bd199c-2ff0-4aa2-b98e-17f0acb72610' for DuckDB." 953 ), 954 default=None, 955 ), 956 ] = None, 957 destination_canonical_name: Annotated[ 958 str | None, 959 Field( 960 description=( 961 "Canonical destination connector name to search for. " 962 "Provide this OR destination_definition_id OR source_definition_id " 963 "OR source_canonical_name (exactly one required). " 964 "Examples: 'destination-duckdb', 'DuckDB'." 965 ), 966 default=None, 967 ), 968 ] = None, 969 organization_id: Annotated[ 970 str | OrganizationAliasEnum | None, 971 Field( 972 description=( 973 "Optional organization ID (UUID) or alias to filter results. " 974 "If provided, only failed attempts from this organization will be returned. " 975 "Accepts '@airbyte-internal' as an alias for the Airbyte internal org." 976 ), 977 default=None, 978 ), 979 ] = None, 980 lookback_days: Annotated[ 981 int, 982 Field(description="Number of days to look back (default: 7)", default=7), 983 ] = 7, 984 limit: Annotated[ 985 int, 986 Field(description="Maximum number of results (default: 100)", default=100), 987 ] = 100, 988 customer_tier_filter: Annotated[ 989 TierFilter, 990 Field( 991 description=( 992 "Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. " 993 "Filters results to only include connections belonging to organizations " 994 "in the specified tier. Use 'ALL' to include all tiers." 995 ), 996 ), 997 ] = "TIER_2", 998) -> list[dict[str, Any]]: 999 """List failed sync attempts for ALL actors using a connector type. 1000 1001 This tool finds all actors with the given connector definition and returns their 1002 failed sync attempts, regardless of whether they have explicit version pins. 1003 1004 Results are always enriched with customer_tier and is_eu fields. 1005 The customer_tier_filter parameter is required to ensure tier-aware querying. 1006 1007 This is useful for investigating connector issues across all users. Use this when 1008 you want to find failures for a connector type regardless of which version users 1009 are on. 1010 1011 Supports both SOURCE and DESTINATION connectors. Provide exactly one of: 1012 source_definition_id, source_canonical_name, destination_definition_id, 1013 or destination_canonical_name. 1014 1015 Key fields in results: 1016 - failure_summary: JSON containing failure details including failureType and messages 1017 - customer_tier: TIER_0, TIER_1, TIER_2, or UNKNOWN 1018 - is_eu: Whether the workspace is in the EU region 1019 - pin_origin_type, pin_origin, pinned_version_id: Version pin context (NULL if not pinned) 1020 - pin_scope_type: 'actor', 'workspace', or 'organization' (NULL if not pinned) 1021 """ 1022 provided_params = [ 1023 source_definition_id, 1024 source_canonical_name, 1025 destination_definition_id, 1026 destination_canonical_name, 1027 ] 1028 if sum(p is not None for p in provided_params) != 1: 1029 raise PyAirbyteInputError( 1030 message=( 1031 "Exactly one of source_definition_id, source_canonical_name, " 1032 "destination_definition_id, or destination_canonical_name must be provided." 1033 ), 1034 ) 1035 1036 is_destination = ( 1037 destination_definition_id is not None or destination_canonical_name is not None 1038 ) 1039 canonical_name = source_canonical_name or destination_canonical_name 1040 resolved_definition_id: str 1041 if canonical_name: 1042 resolved_definition_id = resolve_canonical_name_to_definition_id( 1043 canonical_name=canonical_name, 1044 ) 1045 else: 1046 definition_id = source_definition_id or destination_definition_id 1047 assert definition_id is not None 1048 resolved_definition_id = definition_id 1049 1050 # Resolve organization ID alias 1051 resolved_organization_id = OrganizationAliasEnum.resolve(organization_id) 1052 1053 rows = query_failed_sync_attempts_for_connector( 1054 connector_definition_id=resolved_definition_id, 1055 organization_id=resolved_organization_id, 1056 days=lookback_days, 1057 limit=limit, 1058 is_destination=is_destination, 1059 ) 1060 enriched = enrich_rows_by_org( 1061 rows=rows, 1062 allow_degraded=True, 1063 ) 1064 return filter_rows_by_tier(enriched, customer_tier_filter) 1065 1066 1067@mcp_tool( 1068 read_only=True, 1069 idempotent=True, 1070 open_world=True, 1071) 1072def query_prod_connections_by_connector( 1073 source_definition_id: Annotated[ 1074 str | None, 1075 Field( 1076 description=( 1077 "Source connector definition ID (UUID) to search for. " 1078 "Exactly one of source_definition_id, source_canonical_name, " 1079 "destination_definition_id, or destination_canonical_name is required. " 1080 "Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics." 1081 ), 1082 default=None, 1083 ), 1084 ] = None, 1085 source_canonical_name: Annotated[ 1086 str | None, 1087 Field( 1088 description=( 1089 "Canonical source connector name to search for. " 1090 "Exactly one of source_definition_id, source_canonical_name, " 1091 "destination_definition_id, or destination_canonical_name is required. " 1092 "Examples: 'source-youtube-analytics', 'YouTube Analytics'." 1093 ), 1094 default=None, 1095 ), 1096 ] = None, 1097 destination_definition_id: Annotated[ 1098 str | None, 1099 Field( 1100 description=( 1101 "Destination connector definition ID (UUID) to search for. " 1102 "Exactly one of source_definition_id, source_canonical_name, " 1103 "destination_definition_id, or destination_canonical_name is required. " 1104 "Example: 'e5c8e66c-a480-4a5e-9c0e-e8e5e4c5c5c5' for DuckDB." 1105 ), 1106 default=None, 1107 ), 1108 ] = None, 1109 destination_canonical_name: Annotated[ 1110 str | None, 1111 Field( 1112 description=( 1113 "Canonical destination connector name to search for. " 1114 "Exactly one of source_definition_id, source_canonical_name, " 1115 "destination_definition_id, or destination_canonical_name is required. " 1116 "Examples: 'destination-duckdb', 'DuckDB'." 1117 ), 1118 default=None, 1119 ), 1120 ] = None, 1121 organization_id: Annotated[ 1122 str | OrganizationAliasEnum | None, 1123 Field( 1124 description=( 1125 "Optional organization ID (UUID) or alias to filter results. " 1126 "If provided, only connections in this organization will be returned. " 1127 "Accepts '@airbyte-internal' as an alias for the Airbyte internal org." 1128 ), 1129 default=None, 1130 ), 1131 ] = None, 1132 limit: Annotated[ 1133 int, 1134 Field(description="Maximum number of results (default: 1000)", default=1000), 1135 ] = 1000, 1136 customer_tier_filter: Annotated[ 1137 TierFilter, 1138 Field( 1139 description=( 1140 "Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. " 1141 "Filters results to only include connections belonging to organizations " 1142 "in the specified tier. Use 'ALL' to include all tiers." 1143 ), 1144 ), 1145 ] = "TIER_2", 1146 *, 1147 exclude_pinned: Annotated[ 1148 bool, 1149 Field( 1150 description=( 1151 "If True, exclude connections whose connector is already pinned to a " 1152 "specific version (at any scope level: actor, workspace, or organization). " 1153 "Useful for 'prove fix' workflows where you want to find unpinned " 1154 "connections for live testing. Default: False (include all connections)." 1155 ), 1156 default=False, 1157 ), 1158 ], 1159 enabled_schedules_only: Annotated[ 1160 bool, 1161 Field( 1162 description=( 1163 "If True, only return connections that are both active " 1164 "(not paused/inactive) and on an automated sync schedule " 1165 "(not manual-trigger-only). Useful for canary workflows where " 1166 "you need connections that will produce organic syncs during a " 1167 "monitoring window. Default: False (include all connections)." 1168 ), 1169 default=False, 1170 ), 1171 ], 1172) -> list[dict[str, Any]]: 1173 """Search for all connections using a specific source or destination connector type. 1174 1175 This tool queries the Airbyte Cloud Prod DB Replica directly for fast results. 1176 It finds all connections where the source or destination connector matches the 1177 specified type, regardless of how the connector is named by users. 1178 1179 Results are always enriched with customer_tier and is_eu fields. 1180 The customer_tier_filter parameter is required to ensure tier-aware querying. 1181 1182 Optionally filter by organization_id to limit results to a specific organization. 1183 Use '@airbyte-internal' as an alias for the Airbyte internal organization. 1184 1185 Set `exclude_pinned=True` to filter out connections that are already pinned to a 1186 specific version. This is useful for 'prove fix' live connection testing workflows 1187 where you want to find unpinned connections to test against. 1188 1189 Set `enabled_schedules_only=True` to restrict results to connections that are both 1190 enabled (status='active') and on an automated schedule (not manual-trigger-only). 1191 This is useful for canary prerelease workflows where you need connections that 1192 will run organically during the monitoring window. 1193 1194 Returns a list of connection dicts with workspace context and clickable Cloud UI URLs. 1195 For source queries, returns: connection_id, connection_name, connection_url, source_id, 1196 source_name, source_definition_id, workspace_id, workspace_name, organization_id, 1197 dataplane_group_id, dataplane_name, pin_origin_type, pin_origin, pinned_version_id, 1198 pin_scope_type, customer_tier, is_eu. 1199 For destination queries, returns: connection_id, connection_name, connection_url, 1200 destination_id, destination_name, destination_definition_id, workspace_id, 1201 workspace_name, organization_id, dataplane_group_id, dataplane_name, pin_origin_type, 1202 pin_origin, pinned_version_id, pin_scope_type, customer_tier, is_eu. 1203 1204 pin_scope_type is 'actor', 'workspace', or 'organization' indicating which scope 1205 level the effective pin came from (NULL if not pinned). 1206 """ 1207 # Validate that exactly one of the four connector parameters is provided 1208 provided_params = [ 1209 source_definition_id, 1210 source_canonical_name, 1211 destination_definition_id, 1212 destination_canonical_name, 1213 ] 1214 num_provided = sum(p is not None for p in provided_params) 1215 if num_provided != 1: 1216 raise PyAirbyteInputError( 1217 message=( 1218 "Exactly one of source_definition_id, source_canonical_name, " 1219 "destination_definition_id, or destination_canonical_name must be provided." 1220 ), 1221 ) 1222 1223 # Determine if this is a source or destination query and resolve the definition ID 1224 is_source_query = ( 1225 source_definition_id is not None or source_canonical_name is not None 1226 ) 1227 resolved_definition_id: str 1228 1229 if source_canonical_name: 1230 resolved_definition_id = resolve_canonical_name_to_definition_id( 1231 canonical_name=source_canonical_name, 1232 ) 1233 elif source_definition_id: 1234 resolved_definition_id = source_definition_id 1235 elif destination_canonical_name: 1236 resolved_definition_id = resolve_canonical_name_to_definition_id( 1237 canonical_name=destination_canonical_name, 1238 ) 1239 else: 1240 resolved_definition_id = destination_definition_id # ty: ignore[invalid-assignment] 1241 1242 # Resolve organization ID alias 1243 resolved_organization_id = OrganizationAliasEnum.resolve(organization_id) 1244 1245 # Query the database based on connector type 1246 if is_source_query: 1247 rows = [ 1248 { 1249 "organization_id": str(row.get("organization_id", "")), 1250 "workspace_id": str(row["workspace_id"]), 1251 "workspace_name": row.get("workspace_name", ""), 1252 "connection_id": str(row["connection_id"]), 1253 "connection_name": row.get("connection_name", ""), 1254 "connection_url": ( 1255 f"{CLOUD_UI_BASE_URL}/workspaces/{row['workspace_id']}" 1256 f"/connections/{row['connection_id']}/status" 1257 ), 1258 "source_id": str(row["source_id"]), 1259 "source_name": row.get("source_name", ""), 1260 "source_definition_id": str(row["source_definition_id"]), 1261 "dataplane_group_id": str(row.get("dataplane_group_id", "")), 1262 "dataplane_name": row.get("dataplane_name", ""), 1263 "pin_origin_type": row.get("pin_origin_type"), 1264 "pin_origin": row.get("pin_origin"), 1265 "pinned_version_id": _opt_str(row.get("pinned_version_id")), 1266 "pin_scope_type": row.get("pin_scope_type"), 1267 } 1268 for row in query_connections_by_connector( 1269 connector_definition_id=resolved_definition_id, 1270 organization_id=resolved_organization_id, 1271 limit=limit, 1272 exclude_pinned=exclude_pinned, 1273 enabled_schedules_only=enabled_schedules_only, 1274 ) 1275 ] 1276 else: 1277 # Destination query 1278 rows = [ 1279 { 1280 "organization_id": str(row.get("organization_id", "")), 1281 "workspace_id": str(row["workspace_id"]), 1282 "workspace_name": row.get("workspace_name", ""), 1283 "connection_id": str(row["connection_id"]), 1284 "connection_name": row.get("connection_name", ""), 1285 "connection_url": ( 1286 f"{CLOUD_UI_BASE_URL}/workspaces/{row['workspace_id']}" 1287 f"/connections/{row['connection_id']}/status" 1288 ), 1289 "destination_id": str(row["destination_id"]), 1290 "destination_name": row.get("destination_name", ""), 1291 "destination_definition_id": str(row["destination_definition_id"]), 1292 "dataplane_group_id": str(row.get("dataplane_group_id", "")), 1293 "dataplane_name": row.get("dataplane_name", ""), 1294 "pin_origin_type": row.get("pin_origin_type"), 1295 "pin_origin": row.get("pin_origin"), 1296 "pinned_version_id": _opt_str(row.get("pinned_version_id")), 1297 "pin_scope_type": row.get("pin_scope_type"), 1298 } 1299 for row in query_connections_by_destination_connector( 1300 connector_definition_id=resolved_definition_id, 1301 organization_id=resolved_organization_id, 1302 limit=limit, 1303 exclude_pinned=exclude_pinned, 1304 enabled_schedules_only=enabled_schedules_only, 1305 ) 1306 ] 1307 1308 enriched = enrich_rows_by_org( 1309 rows=rows, 1310 allow_degraded=True, 1311 ) 1312 return filter_rows_by_tier(enriched, customer_tier_filter) 1313 1314 1315@mcp_tool( 1316 read_only=True, 1317 idempotent=True, 1318 open_world=True, 1319) 1320def query_prod_connections_by_stream( 1321 stream_name: Annotated[ 1322 str, 1323 Field( 1324 description=( 1325 "Name of the stream to search for in connection catalogs. " 1326 "This must match the exact stream name as configured in the connection. " 1327 "Examples: 'global_exclusions', 'campaigns', 'users'." 1328 ), 1329 ), 1330 ], 1331 source_definition_id: Annotated[ 1332 str | None, 1333 Field( 1334 description=( 1335 "Source connector definition ID (UUID) to search for. " 1336 "Provide this OR source_canonical_name (exactly one required). " 1337 "Example: 'afa734e4-3571-11ec-991a-1e0031268139' for YouTube Analytics." 1338 ), 1339 default=None, 1340 ), 1341 ], 1342 source_canonical_name: Annotated[ 1343 str | None, 1344 Field( 1345 description=( 1346 "Canonical source connector name to search for. " 1347 "Provide this OR source_definition_id (exactly one required). " 1348 "Examples: 'source-klaviyo', 'Klaviyo', 'source-youtube-analytics'." 1349 ), 1350 default=None, 1351 ), 1352 ], 1353 organization_id: Annotated[ 1354 str | OrganizationAliasEnum | None, 1355 Field( 1356 description=( 1357 "Optional organization ID (UUID) or alias to filter results. " 1358 "If provided, only connections in this organization will be returned. " 1359 "Accepts '@airbyte-internal' as an alias for the Airbyte internal org." 1360 ), 1361 default=None, 1362 ), 1363 ], 1364 limit: Annotated[ 1365 int, 1366 Field(description="Maximum number of results (default: 100)", default=100), 1367 ], 1368 customer_tier_filter: Annotated[ 1369 TierFilter, 1370 Field( 1371 description=( 1372 "Required tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. " 1373 "Filters results to only include connections belonging to organizations " 1374 "in the specified tier. Use 'ALL' to include all tiers." 1375 ), 1376 ), 1377 ] = "TIER_2", 1378) -> list[dict[str, Any]]: 1379 """Find connections that have a specific stream enabled in their catalog. 1380 1381 This tool searches the connection's configured catalog (JSONB) for streams 1382 matching the specified name. It's particularly useful when validating 1383 connector fixes that affect specific streams - you can quickly find 1384 customer connections that use the affected stream. 1385 1386 Results are always enriched with customer_tier and is_eu fields. 1387 The customer_tier_filter parameter is required to ensure tier-aware querying. 1388 1389 Use cases: 1390 - Finding connections with a specific stream enabled for regression testing 1391 - Validating connector fixes that affect particular streams 1392 - Identifying which customers use rarely-enabled streams 1393 1394 Returns a list of connection dicts with workspace context and clickable Cloud UI URLs. 1395 """ 1396 provided_params = [source_definition_id, source_canonical_name] 1397 num_provided = sum(p is not None for p in provided_params) 1398 if num_provided != 1: 1399 raise PyAirbyteInputError( 1400 message=( 1401 "Exactly one of source_definition_id or source_canonical_name " 1402 "must be provided." 1403 ), 1404 ) 1405 1406 resolved_definition_id: str 1407 if source_canonical_name: 1408 resolved_definition_id = resolve_canonical_name_to_definition_id( 1409 canonical_name=source_canonical_name, 1410 ) 1411 else: 1412 assert source_definition_id is not None 1413 resolved_definition_id = source_definition_id 1414 1415 resolved_organization_id = OrganizationAliasEnum.resolve(organization_id) 1416 1417 rows = [ 1418 { 1419 "organization_id": str(row.get("organization_id", "")), 1420 "workspace_id": str(row["workspace_id"]), 1421 "workspace_name": row.get("workspace_name", ""), 1422 "connection_id": str(row["connection_id"]), 1423 "connection_name": row.get("connection_name", ""), 1424 "connection_status": row.get("connection_status", ""), 1425 "connection_url": ( 1426 f"{CLOUD_UI_BASE_URL}/workspaces/{row['workspace_id']}" 1427 f"/connections/{row['connection_id']}/status" 1428 ), 1429 "source_id": str(row["source_id"]), 1430 "source_name": row.get("source_name", ""), 1431 "source_definition_id": str(row["source_definition_id"]), 1432 "dataplane_group_id": str(row.get("dataplane_group_id", "")), 1433 "dataplane_name": row.get("dataplane_name", ""), 1434 } 1435 for row in query_connections_by_stream( 1436 connector_definition_id=resolved_definition_id, 1437 stream_name=stream_name, 1438 organization_id=resolved_organization_id, 1439 limit=limit, 1440 ) 1441 ] 1442 enriched = enrich_rows_by_org( 1443 rows=rows, 1444 allow_degraded=True, 1445 ) 1446 return filter_rows_by_tier(enriched, customer_tier_filter) 1447 1448 1449@mcp_tool( 1450 read_only=True, 1451 idempotent=True, 1452) 1453def query_prod_organizations( 1454 name_contains: Annotated[ 1455 str, 1456 Field( 1457 description=( 1458 "Case-insensitive substring to search for in organization name or email. " 1459 "For example, 'acme' will match organizations named 'Acme Corp' or " 1460 "with email 'admin@acme.io'." 1461 ), 1462 ), 1463 ], 1464 limit: Annotated[ 1465 int, 1466 Field( 1467 description="Maximum number of organizations to return (default: 20)", 1468 default=20, 1469 ), 1470 ] = 20, 1471) -> OrganizationSearchResult: 1472 """Search organizations by name or email substring. 1473 1474 Performs a case-insensitive substring match on organization name and email. 1475 Use the returned `organization_id` values with other tools like 1476 `query_prod_connections_by_connector` or `lookup_customer_tiers`. 1477 """ 1478 rows = search_organizations(name_contains=name_contains, limit=limit) 1479 1480 orgs = [ 1481 OrganizationSearchHit( 1482 organization_id=str(row["organization_id"]), 1483 organization_name=row.get("organization_name", ""), 1484 email=row.get("email"), 1485 ) 1486 for row in rows 1487 ] 1488 1489 # Enrich with tier annotation 1490 org_ids = [o.organization_id for o in orgs] 1491 tier_results = { 1492 r.organization_id: r 1493 for r in get_org_tiers( 1494 organization_ids=org_ids, 1495 allow_degraded=True, 1496 ) 1497 } 1498 for org in orgs: 1499 tier_result = tier_results.get(org.organization_id) 1500 if tier_result: 1501 org.customer_tier = tier_result.customer_tier 1502 1503 return OrganizationSearchResult( 1504 name_contains=name_contains, 1505 total_found=len(orgs), 1506 organizations=orgs, 1507 ) 1508 1509 1510@mcp_tool( 1511 read_only=True, 1512 idempotent=True, 1513) 1514def query_prod_workspaces( 1515 name_contains: Annotated[ 1516 str | None, 1517 Field( 1518 description=( 1519 "Case-insensitive substring to search for in workspace name or slug. " 1520 "For example, 'acme' will match workspaces named 'Acme Staging'." 1521 ), 1522 default=None, 1523 ), 1524 ] = None, 1525 email_domain: Annotated[ 1526 str | None, 1527 Field( 1528 description=( 1529 "Email domain to search for (e.g., 'motherduck.com'). " 1530 "Do not include the '@' symbol." 1531 ), 1532 default=None, 1533 ), 1534 ] = None, 1535 limit: Annotated[ 1536 int, 1537 Field( 1538 description="Maximum number of workspaces to return (default: 100)", 1539 default=100, 1540 ), 1541 ] = 100, 1542) -> WorkspaceSearchResult: 1543 """Search workspaces by name substring or email domain. 1544 1545 At least one of `name_contains` or `email_domain` must be provided. 1546 When `name_contains` is given, performs a case-insensitive substring match 1547 on workspace name and slug. When `email_domain` is given, matches 1548 workspaces by user email domain. 1549 1550 The returned organization IDs can be used with other tools like 1551 `query_prod_connections_by_connector` to find connections within 1552 those organizations for safe testing. 1553 """ 1554 if not name_contains and not email_domain: 1555 raise PyAirbyteInputError( 1556 message="At least one of `name_contains` or `email_domain` must be provided.", 1557 ) 1558 1559 if name_contains: 1560 rows = search_workspaces(name_contains=name_contains, limit=limit) 1561 else: 1562 assert email_domain is not None 1563 clean_domain = email_domain.lstrip("@") 1564 rows = query_workspaces_by_email_domain(email_domain=clean_domain, limit=limit) 1565 1566 workspaces = [ 1567 WorkspaceInfo( 1568 organization_id=str(row["organization_id"]), 1569 workspace_id=str(row["workspace_id"]), 1570 workspace_name=row.get("workspace_name", ""), 1571 slug=row.get("slug"), 1572 email=row.get("email"), 1573 dataplane_group_id=_opt_str(row.get("dataplane_group_id")), 1574 dataplane_name=row.get("dataplane_name"), 1575 created_at=row.get("created_at"), 1576 ) 1577 for row in rows 1578 ] 1579 1580 # Enrich with tier annotation (annotation only, no filtering) 1581 unique_org_ids = list(dict.fromkeys(w.organization_id for w in workspaces)) 1582 tier_results = { 1583 r.organization_id: r 1584 for r in get_org_tiers( 1585 organization_ids=unique_org_ids, 1586 allow_degraded=True, 1587 ) 1588 } 1589 for ws in workspaces: 1590 tier_result = tier_results.get(ws.organization_id) 1591 if tier_result: 1592 ws.customer_tier = tier_result.customer_tier 1593 ws.is_eu = ws.dataplane_name == "EU" if ws.dataplane_name else False 1594 1595 return WorkspaceSearchResult( 1596 name_contains=name_contains, 1597 email_domain=email_domain.lstrip("@") if email_domain else None, 1598 total_workspaces_found=len(workspaces), 1599 unique_organization_ids=unique_org_ids, 1600 workspaces=workspaces, 1601 ) 1602 1603 1604# Backward-compatible alias 1605query_prod_workspaces_by_email_domain = query_prod_workspaces 1606 1607 1608def _build_connector_stats( 1609 connector_definition_id: str, 1610 connector_type: str, 1611 canonical_name: str | None, 1612 rows: list[dict[str, Any]], 1613 version_tags: dict[str, str | None], 1614) -> ConnectorConnectionStats: 1615 """Build ConnectorConnectionStats from query result rows.""" 1616 # Aggregate totals across all version groups 1617 total_connections = 0 1618 enabled_connections = 0 1619 active_connections = 0 1620 pinned_connections = 0 1621 unpinned_connections = 0 1622 total_succeeded = 0 1623 total_failed = 0 1624 total_cancelled = 0 1625 total_running = 0 1626 total_unknown = 0 1627 1628 by_version: list[VersionPinStats] = [] 1629 1630 for row in rows: 1631 version_id = row.get("pinned_version_id") 1632 row_total = int(row.get("total_connections", 0)) 1633 row_enabled = int(row.get("enabled_connections", 0)) 1634 row_active = int(row.get("active_connections", 0)) 1635 row_pinned = int(row.get("pinned_connections", 0)) 1636 row_unpinned = int(row.get("unpinned_connections", 0)) 1637 row_succeeded = int(row.get("succeeded_connections", 0)) 1638 row_failed = int(row.get("failed_connections", 0)) 1639 row_cancelled = int(row.get("cancelled_connections", 0)) 1640 row_running = int(row.get("running_connections", 0)) 1641 row_unknown = int(row.get("unknown_connections", 0)) 1642 1643 total_connections += row_total 1644 enabled_connections += row_enabled 1645 active_connections += row_active 1646 pinned_connections += row_pinned 1647 unpinned_connections += row_unpinned 1648 total_succeeded += row_succeeded 1649 total_failed += row_failed 1650 total_cancelled += row_cancelled 1651 total_running += row_running 1652 total_unknown += row_unknown 1653 1654 by_version.append( 1655 VersionPinStats( 1656 pinned_version_id=str(version_id) if version_id else None, 1657 docker_image_tag=version_tags.get(str(version_id)) 1658 if version_id 1659 else None, 1660 total_connections=row_total, 1661 enabled_connections=row_enabled, 1662 active_connections=row_active, 1663 latest_attempt=LatestAttemptBreakdown( 1664 succeeded=row_succeeded, 1665 failed=row_failed, 1666 cancelled=row_cancelled, 1667 running=row_running, 1668 unknown=row_unknown, 1669 ), 1670 ) 1671 ) 1672 1673 return ConnectorConnectionStats( 1674 connector_definition_id=connector_definition_id, 1675 connector_type=connector_type, 1676 canonical_name=canonical_name, 1677 total_connections=total_connections, 1678 enabled_connections=enabled_connections, 1679 active_connections=active_connections, 1680 pinned_connections=pinned_connections, 1681 unpinned_connections=unpinned_connections, 1682 latest_attempt=LatestAttemptBreakdown( 1683 succeeded=total_succeeded, 1684 failed=total_failed, 1685 cancelled=total_cancelled, 1686 running=total_running, 1687 unknown=total_unknown, 1688 ), 1689 by_version=by_version, 1690 ) 1691 1692 1693@mcp_tool( 1694 read_only=True, 1695 idempotent=True, 1696 open_world=True, 1697) 1698def query_prod_connector_connection_stats( 1699 source_definition_ids: Annotated[ 1700 list[str] | None, 1701 Field( 1702 description=( 1703 "List of source connector definition IDs (UUIDs) to get stats for. " 1704 "Example: ['afa734e4-3571-11ec-991a-1e0031268139']" 1705 ), 1706 default=None, 1707 ), 1708 ] = None, 1709 destination_definition_ids: Annotated[ 1710 list[str] | None, 1711 Field( 1712 description=( 1713 "List of destination connector definition IDs (UUIDs) to get stats for. " 1714 "Example: ['94bd199c-2ff0-4aa2-b98e-17f0acb72610']" 1715 ), 1716 default=None, 1717 ), 1718 ] = None, 1719 lookback_days: Annotated[ 1720 int, 1721 Field( 1722 description=( 1723 "Number of days to look back for 'active' connections (default: 7). " 1724 "Connections with sync activity within this window are counted as active." 1725 ), 1726 default=7, 1727 ), 1728 ] = 7, 1729) -> ConnectorConnectionStatsResponse: 1730 """Get aggregate connection stats for multiple connectors. 1731 1732 Returns counts of connections grouped by pinned version for each connector, 1733 including: 1734 - Total, enabled, and active connection counts 1735 - Pinned vs unpinned breakdown 1736 - Latest attempt status breakdown (succeeded, failed, cancelled, running, unknown) 1737 1738 This tool is designed for release monitoring workflows. It allows you to: 1739 1. Query recently released connectors to identify which ones to monitor 1740 2. Get aggregate stats showing how many connections are using each version 1741 3. See health metrics (pass/fail) broken down by version 1742 1743 The `lookback_days` parameter controls the lookback window for: 1744 - Counting 'active' connections (those with recent sync activity) 1745 - Determining 'latest attempt status' (most recent attempt within the window) 1746 1747 Connections with no sync activity in the lookback window will have 1748 'unknown' status in the latest_attempt breakdown. 1749 """ 1750 # Initialize empty lists if None 1751 source_ids = source_definition_ids or [] 1752 destination_ids = destination_definition_ids or [] 1753 1754 if not source_ids and not destination_ids: 1755 raise PyAirbyteInputError( 1756 message=( 1757 "At least one of source_definition_ids or destination_definition_ids " 1758 "must be provided." 1759 ), 1760 ) 1761 1762 sources: list[ConnectorConnectionStats] = [] 1763 destinations: list[ConnectorConnectionStats] = [] 1764 1765 # Process source connectors 1766 for source_def_id in source_ids: 1767 # Get version info for tag lookup 1768 versions = query_connector_versions(source_def_id) 1769 version_tags = { 1770 str(v["version_id"]): v.get("docker_image_tag") for v in versions 1771 } 1772 1773 # Get aggregate stats 1774 rows = query_source_connection_stats(source_def_id, days=lookback_days) 1775 1776 sources.append( 1777 _build_connector_stats( 1778 connector_definition_id=source_def_id, 1779 connector_type="source", 1780 canonical_name=None, 1781 rows=rows, 1782 version_tags=version_tags, 1783 ) 1784 ) 1785 1786 # Process destination connectors 1787 for dest_def_id in destination_ids: 1788 # Get version info for tag lookup 1789 versions = query_connector_versions(dest_def_id) 1790 version_tags = { 1791 str(v["version_id"]): v.get("docker_image_tag") for v in versions 1792 } 1793 1794 # Get aggregate stats 1795 rows = query_destination_connection_stats(dest_def_id, days=lookback_days) 1796 1797 destinations.append( 1798 _build_connector_stats( 1799 connector_definition_id=dest_def_id, 1800 connector_type="destination", 1801 canonical_name=None, 1802 rows=rows, 1803 version_tags=version_tags, 1804 ) 1805 ) 1806 1807 return ConnectorConnectionStatsResponse( 1808 sources=sources, 1809 destinations=destinations, 1810 lookback_days=lookback_days, 1811 generated_at=datetime.now(timezone.utc), 1812 ) 1813 1814 1815# ============================================================================= 1816# Connector Rollout Models and Tools 1817# ============================================================================= 1818 1819 1820class ConnectorRolloutInfo(BaseModel): 1821 """Information about a connector rollout.""" 1822 1823 rollout_id: str = Field(description="The rollout UUID") 1824 actor_definition_id: str = Field(description="The connector definition UUID") 1825 state: str = Field( 1826 description="Rollout state: initialized, workflow_started, in_progress, " 1827 "paused, finalizing, succeeded, errored, failed_rolled_back, canceled" 1828 ) 1829 initial_rollout_pct: int | None = Field( 1830 default=None, description="Initial rollout percentage" 1831 ) 1832 current_target_rollout_pct: int | None = Field( 1833 default=None, description="Current target rollout percentage" 1834 ) 1835 final_target_rollout_pct: int | None = Field( 1836 default=None, description="Final target rollout percentage" 1837 ) 1838 has_breaking_changes: bool = Field( 1839 description="Whether the RC has breaking changes" 1840 ) 1841 max_step_wait_time_mins: int | None = Field( 1842 default=None, description="Maximum wait time between rollout steps in minutes" 1843 ) 1844 rollout_strategy: str | None = Field( 1845 default=None, description="Rollout strategy: manual, automated, overridden" 1846 ) 1847 updated_by_user_id: str | None = Field( 1848 default=None, 1849 description="User ID recorded as last updating the rollout", 1850 ) 1851 updated_by_user_name: str | None = Field( 1852 default=None, 1853 description="Name recorded as last updating the rollout", 1854 ) 1855 updated_by_user_email: str | None = Field( 1856 default=None, 1857 description="Email recorded as last updating the rollout", 1858 ) 1859 workflow_run_id: str | None = Field( 1860 default=None, description="Temporal workflow run ID" 1861 ) 1862 error_msg: str | None = Field(default=None, description="Error message if errored") 1863 failed_reason: str | None = Field( 1864 default=None, description="Reason for failure if failed" 1865 ) 1866 paused_reason: str | None = Field( 1867 default=None, description="Reason for pause if paused" 1868 ) 1869 tag: str | None = Field(default=None, description="Optional tag for the rollout") 1870 created_at: datetime | None = Field( 1871 default=None, description="When the rollout was created" 1872 ) 1873 updated_at: datetime | None = Field( 1874 default=None, description="When the rollout was last updated" 1875 ) 1876 completed_at: datetime | None = Field( 1877 default=None, description="When the rollout completed (if terminal)" 1878 ) 1879 expires_at: datetime | None = Field( 1880 default=None, description="When the rollout expires" 1881 ) 1882 rc_docker_image_tag: str | None = Field( 1883 default=None, description="Docker image tag of the release candidate" 1884 ) 1885 rc_docker_repository: str | None = Field( 1886 default=None, description="Docker repository of the release candidate" 1887 ) 1888 initial_docker_image_tag: str | None = Field( 1889 default=None, description="Docker image tag of the initial version" 1890 ) 1891 initial_docker_repository: str | None = Field( 1892 default=None, description="Docker repository of the initial version" 1893 ) 1894 filters: dict[str, Any] | None = Field( 1895 default=None, 1896 description="Raw rollout filters JSON (e.g., {'tierFilter': {'tier': 'TIER_0'}})", 1897 ) 1898 customer_tier: str | None = Field( 1899 default=None, 1900 description="Customer tier targeted by this rollout (extracted from filters), " 1901 "e.g., 'TIER_0', 'TIER_1'. None if no tier filter is set.", 1902 ) 1903 1904 1905def _parse_rollout_filters(filters_raw: Any) -> dict[str, Any] | None: 1906 """Parse the rollout filters field from a database row. 1907 1908 The filters field may be a JSON string, a dict, or None. 1909 """ 1910 if filters_raw is None: 1911 return None 1912 if isinstance(filters_raw, dict): 1913 return filters_raw 1914 if isinstance(filters_raw, str): 1915 try: 1916 parsed = json.loads(filters_raw) 1917 if isinstance(parsed, dict): 1918 return parsed 1919 except (json.JSONDecodeError, TypeError): 1920 pass 1921 return None 1922 1923 1924def _extract_tier_from_filters(filters_raw: Any) -> str | None: 1925 """Extract customer tier from rollout filters JSON. 1926 1927 Supports two formats: 1928 - Legacy: `{"tierFilter": {"tier": "TIER_0"}}` 1929 - Current: `{"customerTierFilters": [{"name": "TIER", "value": ["TIER_1"], "operator": "IN"}]}` 1930 """ 1931 parsed = _parse_rollout_filters(filters_raw) 1932 if parsed is None: 1933 return None 1934 1935 # Current format: customerTierFilters list 1936 tier_filters = parsed.get("customerTierFilters") 1937 if isinstance(tier_filters, list): 1938 for entry in tier_filters: 1939 if isinstance(entry, dict) and entry.get("name") == "TIER": 1940 values = entry.get("value") 1941 if isinstance(values, list) and len(values) == 1: 1942 return str(values[0]) 1943 if isinstance(values, list) and len(values) > 1: 1944 return ", ".join(str(v) for v in values) 1945 1946 # Legacy format: tierFilter dict 1947 tier_filter = parsed.get("tierFilter") 1948 if isinstance(tier_filter, dict): 1949 tier = tier_filter.get("tier") 1950 if isinstance(tier, str): 1951 return tier 1952 1953 return None 1954 1955 1956def _row_to_connector_rollout_info(row: dict[str, Any]) -> ConnectorRolloutInfo: 1957 """Convert a database row to a ConnectorRolloutInfo model.""" 1958 return ConnectorRolloutInfo( 1959 rollout_id=str(row["rollout_id"]), 1960 actor_definition_id=str(row["actor_definition_id"]), 1961 state=row["state"], 1962 initial_rollout_pct=row.get("initial_rollout_pct"), 1963 current_target_rollout_pct=row.get("current_target_rollout_pct"), 1964 final_target_rollout_pct=row.get("final_target_rollout_pct"), 1965 has_breaking_changes=row["has_breaking_changes"], 1966 max_step_wait_time_mins=row.get("max_step_wait_time_mins"), 1967 rollout_strategy=row.get("rollout_strategy"), 1968 updated_by_user_id=str(row["updated_by_user_id"]) 1969 if row.get("updated_by_user_id") is not None 1970 else None, 1971 updated_by_user_name=row.get("updated_by_user_name"), 1972 updated_by_user_email=row.get("updated_by_user_email"), 1973 workflow_run_id=row.get("workflow_run_id"), 1974 error_msg=row.get("error_msg"), 1975 failed_reason=row.get("failed_reason"), 1976 paused_reason=row.get("paused_reason"), 1977 tag=row.get("tag"), 1978 created_at=row.get("created_at"), 1979 updated_at=row.get("updated_at"), 1980 completed_at=row.get("completed_at"), 1981 expires_at=row.get("expires_at"), 1982 rc_docker_image_tag=row.get("rc_docker_image_tag"), 1983 rc_docker_repository=row.get("rc_docker_repository"), 1984 initial_docker_image_tag=row.get("initial_docker_image_tag"), 1985 initial_docker_repository=row.get("initial_docker_repository"), 1986 filters=_parse_rollout_filters(row.get("filters")), 1987 customer_tier=_extract_tier_from_filters(row.get("filters")), 1988 ) 1989 1990 1991@mcp_tool( 1992 read_only=True, 1993 idempotent=True, 1994) 1995def query_prod_connector_rollouts( 1996 actor_definition_id: Annotated[ 1997 str | None, 1998 Field(description="Connector definition UUID to filter by (optional)"), 1999 ] = None, 2000 rollout_id: Annotated[ 2001 str | None, 2002 Field(description="Specific rollout UUID to look up (optional)"), 2003 ] = None, 2004 active_only: Annotated[ 2005 bool, 2006 Field(description="If true, only return active (non-terminal) rollouts"), 2007 ] = False, 2008 limit: Annotated[ 2009 int, 2010 Field(description="Maximum number of results (default: 100)"), 2011 ] = 100, 2012) -> list[ConnectorRolloutInfo]: 2013 """Query connector rollouts with flexible filtering. 2014 2015 Returns rollouts based on the provided filters. If no filters are specified, 2016 returns all active rollouts. Useful for monitoring rollout status and history. 2017 2018 Filter behavior: 2019 - rollout_id: Returns that specific rollout (ignores other filters) 2020 - active_only: Returns only active (non-terminal) rollouts 2021 - actor_definition_id: Returns rollouts for that specific connector 2022 - No filters: Returns all active rollouts (same as active_only=True) 2023 """ 2024 rows = query_connector_rollouts( 2025 actor_definition_id=actor_definition_id, 2026 rollout_id=rollout_id, 2027 active_only=active_only, 2028 limit=limit, 2029 ) 2030 return [_row_to_connector_rollout_info(row) for row in rows] 2031 2032 2033@mcp_tool( 2034 read_only=True, 2035 idempotent=True, 2036 open_world=True, 2037) 2038def query_prod_connection_sync_activity( 2039 start_at: Annotated[ 2040 datetime, 2041 Field( 2042 description=( 2043 "Inclusive start timestamp for the sync activity window. " 2044 "Must be timezone-aware (ISO 8601 with offset or `Z`)." 2045 ), 2046 ), 2047 ], 2048 end_at: Annotated[ 2049 datetime, 2050 Field( 2051 description=( 2052 "Exclusive end timestamp for the sync activity window. " 2053 "Must be timezone-aware and strictly after `start_at`." 2054 ), 2055 ), 2056 ], 2057 organization_id: Annotated[ 2058 str | OrganizationAliasEnum | None, 2059 Field( 2060 description=( 2061 "Optional organization UUID or alias. At least one of " 2062 "`organization_id`, `workspace_id`, or `connection_ids` is " 2063 "required. Accepts `@airbyte-internal` as an alias for the " 2064 "Airbyte internal org." 2065 ), 2066 default=None, 2067 ), 2068 ] = None, 2069 workspace_id: Annotated[ 2070 str | WorkspaceAliasEnum | None, 2071 Field( 2072 description=( 2073 "Optional workspace UUID or alias. At least one of " 2074 "`organization_id`, `workspace_id`, or `connection_ids` is " 2075 "required. Accepts `@devin-ai-sandbox` as an alias for the " 2076 "Devin AI sandbox workspace." 2077 ), 2078 default=None, 2079 ), 2080 ] = None, 2081 connection_ids: Annotated[ 2082 list[str] | None, 2083 Field( 2084 description=( 2085 "Optional list of connection UUIDs. At least one of " 2086 "`organization_id`, `workspace_id`, or `connection_ids` is " 2087 "required." 2088 ), 2089 default=None, 2090 ), 2091 ] = None, 2092 status_filter: Annotated[ 2093 StatusFilter, 2094 Field( 2095 description=( 2096 "Filter by job status: `all` (default), `succeeded`, or " 2097 "`failed`. Applied to `jobs.status` in the Prod DB Replica." 2098 ), 2099 default=StatusFilter.ALL, 2100 ), 2101 ] = StatusFilter.ALL, 2102 limit: Annotated[ 2103 int, 2104 Field( 2105 description="Maximum number of attempt rows to return.", 2106 default=1000, 2107 ), 2108 ] = 1000, 2109) -> list[dict[str, Any]]: 2110 """List recent sync jobs and attempts from the Prod DB Replica. 2111 2112 Returns one row per `(job, attempt)` pair for sync jobs whose `updated_at` 2113 falls in `[start_at, end_at)`, scoped to the provided organization, 2114 workspace, or connection IDs. Designed for live operational lookups — 2115 e.g. "what happened on this connection in the last hour" — not for 2116 historical analysis. 2117 2118 Each row is enriched with `customer_tier` and `is_eu` for the owning 2119 organization. Tier filtering is intentionally not applied — this is a 2120 read-only observability query. 2121 2122 Input requirements: 2123 - At least one of `organization_id`, `workspace_id`, or `connection_ids` 2124 must be provided (any combination is accepted). 2125 - `start_at` and `end_at` must be timezone-aware and `start_at < end_at`. 2126 2127 Key fields in each row: 2128 - `job_id`, `attempt_id`, `attempt_number` 2129 - `job_status`, `attempt_status` 2130 - `job_started_at`, `job_updated_at`, `attempt_ended_at` 2131 - `failure_summary` (JSON; populated when an attempt failed) 2132 - `connection_id`, `connection_name`, `connection_status` 2133 - `source_actor_id`, `source_actor_name`, `source_actor_definition_id` 2134 - `destination_actor_id`, `destination_actor_name`, 2135 `destination_actor_definition_id` 2136 - `workspace_id`, `workspace_name`, `organization_id` 2137 - `dataplane_group_id`, `dataplane_name` 2138 - `customer_tier`, `is_eu` (added by tier enrichment) 2139 """ 2140 resolved_organization_id = OrganizationAliasEnum.resolve(organization_id) 2141 resolved_workspace_id = WorkspaceAliasEnum.resolve(workspace_id) 2142 _validate_sync_activity_scope( 2143 organization_id=resolved_organization_id, 2144 workspace_id=resolved_workspace_id, 2145 connection_ids=connection_ids, 2146 ) 2147 normalized_start_at, normalized_end_at = _validate_sync_activity_window( 2148 start_at=start_at, 2149 end_at=end_at, 2150 ) 2151 2152 rows = query_connection_sync_activity_from_prod( 2153 start_at=normalized_start_at, 2154 end_at=normalized_end_at, 2155 organization_id=resolved_organization_id, 2156 workspace_id=resolved_workspace_id, 2157 connection_ids=connection_ids, 2158 status_filter=status_filter.value, 2159 limit=limit, 2160 ) 2161 return enrich_rows_by_org( 2162 rows=rows, 2163 allow_degraded=True, 2164 ) 2165 2166 2167# ============================================================================= 2168# Pinned Connector Versions Models and Tools 2169# ============================================================================= 2170 2171 2172class PinnedConnectorVersionInfo(BaseModel): 2173 """A connector version that has at least one scoped configuration pin.""" 2174 2175 version_id: str = Field(description="The actor_definition_version UUID") 2176 connector_definition_id: str = Field(description="The connector definition UUID") 2177 connector_name: str = Field(description="Human-readable connector name") 2178 docker_repository: str = Field(description="Docker repository path") 2179 docker_image_tag: str = Field(description="Docker image tag for this version") 2180 last_published: str | None = Field( 2181 default=None, description="ISO timestamp when this version was last published" 2182 ) 2183 pin_count: int = Field( 2184 description="Total number of scoped_configuration rows pinning to this version" 2185 ) 2186 breaking_change_pins: int = Field( 2187 default=0, 2188 description="Number of actor-scoped pins created by breaking changes", 2189 ) 2190 rollout_pins: int = Field( 2191 default=0, 2192 description="Number of pins created by connector rollouts", 2193 ) 2194 actor_pins: int = Field( 2195 description="Number of actor-scoped pins (excludes breaking change and rollout pins)" 2196 ) 2197 workspace_pins: int = Field(description="Number of workspace-scoped pins") 2198 org_pins: int = Field(description="Number of organization-scoped pins") 2199 2200 2201@mcp_tool( 2202 read_only=True, 2203 idempotent=True, 2204 open_world=True, 2205) 2206def query_connector_pin_stats( 2207 connector_definition_id: Annotated[ 2208 str | None, 2209 Field( 2210 description="Connector definition UUID to filter by (optional). " 2211 "Mutually exclusive with `connector_canonical_name`." 2212 ), 2213 ] = None, 2214 connector_canonical_name: Annotated[ 2215 str | None, 2216 Field( 2217 description="Connector canonical name (e.g. `source-postgres`) to filter by. " 2218 "Resolved to a definition ID via the registry. " 2219 "Mutually exclusive with `connector_definition_id`." 2220 ), 2221 ] = None, 2222) -> list[PinnedConnectorVersionInfo]: 2223 """Query connector versions that have at least one scoped configuration pin. 2224 2225 Returns versions from the prod DB that are referenced by at least one 2226 `scoped_configuration` pin (`key = 'connector_version'`). Each version 2227 appears exactly once with per-scope pin breakdown (actor, workspace, org). 2228 2229 If neither filter is provided, returns the global superset across all connectors. 2230 """ 2231 if connector_definition_id and connector_canonical_name: 2232 raise PyAirbyteInputError( 2233 message=( 2234 "Provide at most one of `connector_definition_id` or " 2235 "`connector_canonical_name`, not both." 2236 ), 2237 ) 2238 2239 resolved_id: str | None = None 2240 if connector_canonical_name: 2241 resolved_id = resolve_canonical_name_to_definition_id( 2242 canonical_name=connector_canonical_name, 2243 ) 2244 elif connector_definition_id: 2245 resolved_id = connector_definition_id 2246 2247 rows = query_versions_with_pins(actor_definition_id=resolved_id) 2248 return [ 2249 PinnedConnectorVersionInfo( 2250 version_id=str(row["version_id"]), 2251 connector_definition_id=str(row["connector_definition_id"]), 2252 connector_name=row["connector_name"], 2253 docker_repository=row["docker_repository"], 2254 docker_image_tag=row["docker_image_tag"], 2255 last_published=( 2256 row["last_published"].isoformat() if row.get("last_published") else None 2257 ), 2258 pin_count=row["pin_count"], 2259 breaking_change_pins=row.get("breaking_change_pins", 0), 2260 rollout_pins=row.get("rollout_pins", 0), 2261 actor_pins=row.get("actor_pins", 0), 2262 workspace_pins=row.get("workspace_pins", 0), 2263 org_pins=row.get("org_pins", 0), 2264 ) 2265 for row in rows 2266 ] 2267 2268 2269# ============================================================================= 2270# Organization-scoped Pin Models and Tools 2271# ============================================================================= 2272 2273 2274class PinOriginFilter(StrEnum): 2275 """How a `scoped_configuration` pin was created, used to filter pins.""" 2276 2277 ALL = "all" 2278 MANUAL = "manual" 2279 CONNECTOR_ROLLOUT = "connector_rollout" 2280 BREAKING_CHANGE = "breaking_change" 2281 2282 2283# Rollout states considered non-terminal ("active"). Kept in sync with the set 2284# used by the rollout-monitoring SQL in `airbyte_ops_mcp.prod_db_access.sql`. 2285_ACTIVE_ROLLOUT_STATES = frozenset( 2286 { 2287 "initialized", 2288 "workflow_started", 2289 "in_progress", 2290 "paused", 2291 "finalizing", 2292 "errored", 2293 } 2294) 2295 2296 2297def _pin_category(origin_type: str | None) -> str: 2298 """Classify a pin as `rollout`, `breaking_change`, or `manual`.""" 2299 if origin_type == "connector_rollout": 2300 return "rollout" 2301 if origin_type == "breaking_change": 2302 return "breaking_change" 2303 return "manual" 2304 2305 2306# Maps a `PinOriginFilter` to the `_pin_category` value it keeps. `origin_filter` 2307# is applied here in Python rather than in SQL — the fetched pin list is small, 2308# so an in-SQL `:origin_filter` OR-chain would only add scan cost. 2309_ORIGIN_FILTER_CATEGORIES: dict[str, str] = { 2310 "manual": "manual", 2311 "connector_rollout": "rollout", 2312 "breaking_change": "breaking_change", 2313} 2314 2315 2316def _resolve_connector_filter_id( 2317 *, 2318 connector_definition_id: str | None, 2319 connector_canonical_name: str | None, 2320) -> str | None: 2321 """Resolve mutually-exclusive connector inputs to a definition id or `None`. 2322 2323 Blank or whitespace-only inputs are treated as absent (`None`) so an 2324 optional filter passed as `""` does not become a real, zero-matching SQL 2325 filter. 2326 """ 2327 connector_definition_id = (connector_definition_id or "").strip() or None 2328 connector_canonical_name = (connector_canonical_name or "").strip() or None 2329 if connector_definition_id and connector_canonical_name: 2330 raise PyAirbyteInputError( 2331 message=( 2332 "Provide at most one of `connector_definition_id` or " 2333 "`connector_canonical_name`, not both." 2334 ), 2335 ) 2336 if connector_canonical_name: 2337 return resolve_canonical_name_to_definition_id( 2338 canonical_name=connector_canonical_name, 2339 ) 2340 return connector_definition_id 2341 2342 2343def _require_organization_id(organization_id: str | OrganizationAliasEnum) -> str: 2344 """Resolve a required organization id to a canonical UUID string. 2345 2346 Accepts an organization UUID or an `OrganizationAliasEnum` alias, resolving 2347 aliases to their UUID. The result is validated as a UUID and returned in 2348 canonical (lowercased) form for consistent logging and comparison. (The 2349 org-pin SQL casts this to native `uuid`, so matching is case-insensitive 2350 either way.) Raises `PyAirbyteInputError` on blank or malformed input rather 2351 than passing an invalid value to SQL and returning a confusing empty result. 2352 """ 2353 resolved = OrganizationAliasEnum.resolve(organization_id) 2354 if not resolved or not resolved.strip(): 2355 raise PyAirbyteInputError( 2356 message="`organization_id` is required (a non-empty organization UUID or alias).", 2357 ) 2358 try: 2359 return str(uuid.UUID(resolved.strip())) 2360 except ValueError as exc: 2361 raise PyAirbyteInputError( 2362 message="`organization_id` is not a valid organization UUID or alias.", 2363 ) from exc 2364 2365 2366def _normalize_optional_version_id(pinned_version_id: str | None) -> str | None: 2367 """Coerce a blank `pinned_version_id` to `None` and validate UUID format. 2368 2369 A blank or whitespace-only value is treated as "no filter" (`None`) rather 2370 than being passed to the SQL `uuid` cast, which would raise 2371 `invalid input syntax for type uuid`. 2372 """ 2373 normalized = (pinned_version_id or "").strip() or None 2374 if normalized is not None: 2375 try: 2376 uuid.UUID(normalized) 2377 except ValueError as exc: 2378 raise PyAirbyteInputError( 2379 message="`pinned_version_id` is not a valid UUID.", 2380 ) from exc 2381 return normalized 2382 2383 2384class OrgVersionPinStats(BaseModel): 2385 """A connector version pinned somewhere under an organization, with counts.""" 2386 2387 version_id: str = Field(description="The actor_definition_version UUID") 2388 connector_definition_id: str = Field(description="The connector definition UUID") 2389 connector_name: str = Field(description="Human-readable connector name") 2390 docker_repository: str = Field(description="Docker repository path") 2391 docker_image_tag: str = Field(description="Docker image tag for this version") 2392 last_published: str | None = Field( 2393 default=None, description="ISO timestamp when this version was last published" 2394 ) 2395 pin_count: int = Field( 2396 description="Total pins under the org targeting this version (all scopes)" 2397 ) 2398 manual_pins: int = Field( 2399 default=0, 2400 description="Pins with no system origin (user-created manual pins), any scope", 2401 ) 2402 rollout_pins: int = Field( 2403 default=0, description="Pins created by connector rollouts" 2404 ) 2405 breaking_change_pins: int = Field( 2406 default=0, description="Pins created by breaking changes" 2407 ) 2408 actor_pins: int = Field( 2409 description="Manual actor-scoped pins (excludes rollout and breaking-change)" 2410 ) 2411 workspace_pins: int = Field(description="Workspace-scoped pins under the org") 2412 org_pins: int = Field(description="Organization-scoped pins") 2413 has_active_rollout: bool = Field( 2414 default=False, 2415 description=( 2416 "`True` if at least one rollout pin is backed by a non-terminal " 2417 "`connector_rollout`" 2418 ), 2419 ) 2420 2421 2422class OrgConnectorPin(BaseModel): 2423 """A single `scoped_configuration` pin discovered under an organization.""" 2424 2425 connector_definition_id: str = Field(description="The connector definition UUID") 2426 connector_name: str = Field(description="Human-readable connector name") 2427 docker_repository: str = Field(description="Docker repository path") 2428 pinned_version_id: str = Field( 2429 description="The pinned actor_definition_version UUID" 2430 ) 2431 pinned_version_tag: str = Field( 2432 description="Docker image tag of the pinned version" 2433 ) 2434 pin_scope_type: str = Field( 2435 description="Scope of the pin: `organization`, `workspace`, or `actor`" 2436 ) 2437 scope_id: str = Field(description="UUID of the scoped entity") 2438 scope_name: str | None = Field( 2439 default=None, description="Display name of the scoped entity, when resolvable" 2440 ) 2441 pin_category: str = Field( 2442 description="Derived pin type: `manual`, `rollout`, or `breaking_change`" 2443 ) 2444 set_by: str | None = Field( 2445 default=None, 2446 description="Email (or name) of the user who set a manual pin, when known", 2447 ) 2448 rollout_id: str | None = Field( 2449 default=None, description="Backing connector_rollout UUID for rollout pins" 2450 ) 2451 rollout_state: str | None = Field( 2452 default=None, description="State of the backing rollout, for rollout pins" 2453 ) 2454 is_active_rollout: bool = Field( 2455 default=False, 2456 description="`True` when `rollout_state` is a non-terminal (active) state", 2457 ) 2458 description: str | None = Field(default=None, description="Free-text pin reason") 2459 reference_url: str | None = Field( 2460 default=None, description="Reference URL attached to the pin, when present" 2461 ) 2462 created_at: str | None = Field( 2463 default=None, description="ISO timestamp when the pin was created" 2464 ) 2465 expires_at: str | None = Field( 2466 default=None, description="ISO timestamp when the pin expires, when set" 2467 ) 2468 2469 2470@mcp_tool( 2471 read_only=True, 2472 idempotent=True, 2473 open_world=True, 2474) 2475def query_prod_pin_stats_for_organization( 2476 organization_id: Annotated[ 2477 str | OrganizationAliasEnum, 2478 Field( 2479 description="Organization UUID (or `@airbyte-internal` alias) to scope pins to. " 2480 "Resolve organization names to an ID first via `search_organizations`." 2481 ), 2482 ], 2483 connector_definition_id: Annotated[ 2484 str | None, 2485 Field( 2486 description="Connector definition UUID to filter by (optional). " 2487 "Mutually exclusive with `connector_canonical_name`." 2488 ), 2489 ] = None, 2490 connector_canonical_name: Annotated[ 2491 str | None, 2492 Field( 2493 description="Connector canonical name (e.g. `source-postgres`) to filter by. " 2494 "Resolved to a definition ID via the registry. " 2495 "Mutually exclusive with `connector_definition_id`." 2496 ), 2497 ] = None, 2498 limit: Annotated[ 2499 int, 2500 Field(description="Maximum number of versions to return (default: 1000)."), 2501 ] = 1000, 2502) -> list[OrgVersionPinStats]: 2503 """Query connector versions pinned anywhere under an organization. 2504 2505 Returns one row per pinned version, aggregating every `connector_version` 2506 pin whose scope belongs to the organization — the org itself, one of its 2507 workspaces, or an actor within one of those workspaces (actor, workspace, 2508 and organization scopes). Each row carries the per-scope pin breakdown, the 2509 manual/rollout/breaking-change split, and a `has_active_rollout` flag. 2510 2511 This powers the first step of the Organization Pins view (pick an org, then 2512 see the versions pinned under it). Use `query_prod_pins_for_organization` 2513 for the individual pins behind a selected version. 2514 """ 2515 resolved_org_id = _require_organization_id(organization_id) 2516 resolved_connector_id = _resolve_connector_filter_id( 2517 connector_definition_id=connector_definition_id, 2518 connector_canonical_name=connector_canonical_name, 2519 ) 2520 rows = query_org_pin_stats( 2521 resolved_org_id, 2522 connector_definition_id=resolved_connector_id, 2523 limit=limit, 2524 ) 2525 return [ 2526 OrgVersionPinStats( 2527 version_id=str(row["version_id"]), 2528 connector_definition_id=str(row["connector_definition_id"]), 2529 connector_name=row["connector_name"], 2530 docker_repository=row["docker_repository"], 2531 docker_image_tag=row["docker_image_tag"], 2532 last_published=( 2533 row["last_published"].isoformat() if row.get("last_published") else None 2534 ), 2535 pin_count=row["pin_count"], 2536 manual_pins=row.get("manual_pins", 0), 2537 rollout_pins=row.get("rollout_pins", 0), 2538 breaking_change_pins=row.get("breaking_change_pins", 0), 2539 actor_pins=row.get("actor_pins", 0), 2540 workspace_pins=row.get("workspace_pins", 0), 2541 org_pins=row.get("org_pins", 0), 2542 has_active_rollout=bool(row.get("has_active_rollout", False)), 2543 ) 2544 for row in rows 2545 ] 2546 2547 2548@mcp_tool( 2549 read_only=True, 2550 idempotent=True, 2551 open_world=True, 2552) 2553def query_prod_pins_for_organization( 2554 organization_id: Annotated[ 2555 str | OrganizationAliasEnum, 2556 Field( 2557 description="Organization UUID (or `@airbyte-internal` alias) to scope pins to. " 2558 "Resolve organization names to an ID first via `search_organizations`." 2559 ), 2560 ], 2561 connector_definition_id: Annotated[ 2562 str | None, 2563 Field( 2564 description="Connector definition UUID to filter by (optional). " 2565 "Mutually exclusive with `connector_canonical_name`." 2566 ), 2567 ] = None, 2568 connector_canonical_name: Annotated[ 2569 str | None, 2570 Field( 2571 description="Connector canonical name (e.g. `source-postgres`) to filter by. " 2572 "Resolved to a definition ID via the registry. " 2573 "Mutually exclusive with `connector_definition_id`." 2574 ), 2575 ] = None, 2576 pinned_version_id: Annotated[ 2577 str | None, 2578 Field( 2579 description="Actor_definition_version UUID to return only pins targeting " 2580 "that version. This is the post-selection filter for the org pins tab." 2581 ), 2582 ] = None, 2583 origin_filter: Annotated[ 2584 PinOriginFilter, 2585 Field( 2586 description="Restrict by how the pin was created: `all` (default), " 2587 "`manual`, `connector_rollout`, or `breaking_change`." 2588 ), 2589 ] = PinOriginFilter.ALL, 2590 limit: Annotated[ 2591 int, 2592 Field(description="Maximum number of pins to return (default: 1000)."), 2593 ] = 1000, 2594) -> list[OrgConnectorPin]: 2595 """List the individual connector-version pins discovered under an organization. 2596 2597 Returns one row per `scoped_configuration` pin whose scope belongs to the 2598 organization (org/workspace/actor), resolving the pinned connector and 2599 version, the scope's display name, the manual author's email, and — for 2600 rollout-origin pins — the backing `connector_rollout` id and state. This 2601 directly answers whether each pin is manual or caused by an active rollout. 2602 2603 This powers the second step of the Organization Pins view: after picking a 2604 version from `query_prod_pin_stats_for_organization`, pass its 2605 `pinned_version_id` here to list the pins behind it. 2606 """ 2607 resolved_org_id = _require_organization_id(organization_id) 2608 resolved_connector_id = _resolve_connector_filter_id( 2609 connector_definition_id=connector_definition_id, 2610 connector_canonical_name=connector_canonical_name, 2611 ) 2612 rows = query_org_connector_pins( 2613 resolved_org_id, 2614 connector_definition_id=resolved_connector_id, 2615 pinned_version_id=_normalize_optional_version_id(pinned_version_id), 2616 limit=limit, 2617 ) 2618 kept_category = _ORIGIN_FILTER_CATEGORIES.get(origin_filter.value) 2619 pins: list[OrgConnectorPin] = [] 2620 for row in rows: 2621 origin_type = row.get("origin_type") 2622 if kept_category is not None and _pin_category(origin_type) != kept_category: 2623 continue 2624 rollout_state = row.get("rollout_state") 2625 set_by = row.get("pinned_by_user_email") or row.get("pinned_by_user_name") 2626 created_at = row.get("created_at") 2627 expires_at = row.get("expires_at") 2628 pins.append( 2629 OrgConnectorPin( 2630 connector_definition_id=str(row["connector_definition_id"]), 2631 connector_name=row["connector_name"], 2632 docker_repository=row["docker_repository"], 2633 pinned_version_id=str(row["pinned_version_id"]), 2634 pinned_version_tag=row["pinned_version_tag"], 2635 pin_scope_type=str(row["pin_scope_type"]), 2636 scope_id=str(row["scope_id"]), 2637 scope_name=row.get("scope_name"), 2638 pin_category=_pin_category(origin_type), 2639 set_by=str(set_by) if set_by else None, 2640 rollout_id=str(row["rollout_id"]) if row.get("rollout_id") else None, 2641 rollout_state=rollout_state, 2642 is_active_rollout=rollout_state in _ACTIVE_ROLLOUT_STATES, 2643 description=row.get("description"), 2644 reference_url=row.get("reference_url"), 2645 created_at=created_at.isoformat() if created_at else None, 2646 expires_at=expires_at.isoformat() if expires_at else None, 2647 ) 2648 ) 2649 return pins 2650 2651 2652# ============================================================================= 2653# Health and Population Summary Models and Tools 2654# ============================================================================= 2655 2656 2657def _resolve_connector_target( 2658 *, 2659 connector_version_id: str | None, 2660 connector_name: str | None, 2661 connector_version: str | None, 2662 connector_definition_id: str | None, 2663 connector_canonical_name: str | None, 2664) -> tuple[str | None, str, str, str | None]: 2665 """Resolve mixed connector inputs to a common target tuple. 2666 2667 Returns `(version_id, definition_id, docker_repository, docker_image_tag)`. 2668 `version_id` and `docker_image_tag` are `None` when only a definition-level 2669 identifier was supplied. 2670 """ 2671 if connector_version_id is not None: 2672 info = resolve_version_info(connector_version_id) 2673 return ( 2674 connector_version_id, 2675 str(info["actor_definition_id"]), 2676 info["docker_repository"], 2677 info.get("docker_image_tag"), 2678 ) 2679 if connector_name is not None and connector_version is not None: 2680 docker_repository = f"airbyte/{connector_name}" 2681 info = resolve_version_id_by_tag( 2682 docker_repository=docker_repository, 2683 docker_image_tag=connector_version, 2684 ) 2685 return ( 2686 str(info["version_id"]), 2687 str(info["actor_definition_id"]), 2688 docker_repository, 2689 connector_version, 2690 ) 2691 2692 definition_id: str | None = None 2693 if connector_definition_id is not None: 2694 definition_id = connector_definition_id 2695 elif connector_canonical_name is not None: 2696 definition_id = resolve_canonical_name_to_definition_id( 2697 canonical_name=connector_canonical_name, 2698 ) 2699 if definition_id is None: 2700 raise PyAirbyteInputError( 2701 message=( 2702 "Provide one of: `connector_version_id`, " 2703 "`connector_name` + `connector_version`, " 2704 "`connector_definition_id`, or `connector_canonical_name`." 2705 ), 2706 ) 2707 versions = query_connector_versions(definition_id) 2708 if not versions: 2709 raise PyAirbyteInputError( 2710 message=f"No connector versions found for definition: {definition_id}", 2711 ) 2712 return (None, definition_id, versions[0]["docker_repository"], None) 2713 2714 2715class ConnectorPopulationSummary(BaseModel): 2716 """Applied vs potential pinning audience for a connector, split by tier.""" 2717 2718 connector_definition_id: str = Field(description="The connector definition UUID") 2719 connector_version_id: str | None = Field( 2720 default=None, 2721 description="The version UUID, when a specific version was requested", 2722 ) 2723 docker_repository: str = Field(description="Docker repository path") 2724 docker_image_tag: str | None = Field( 2725 default=None, description="Docker image tag, when a version was requested" 2726 ) 2727 customer_tier_filter: str = Field( 2728 description="Tier filter applied to the counts (`TIER_0`/`TIER_1`/`TIER_2`/`UNKNOWN`/`ALL`)" 2729 ) 2730 active: TierSummary = Field( 2731 description=( 2732 "Potential audience: enabled actors of the definition (those with at " 2733 "least one active connection, `status = 'active'`), by tier" 2734 ) 2735 ) 2736 pinned_any: TierSummary = Field( 2737 description="Active actors already pinned to any version, by tier" 2738 ) 2739 eligible: TierSummary = Field( 2740 description="Active actors not pinned to any version (available to pin), by tier" 2741 ) 2742 pinned_to_version: TierSummary | None = Field( 2743 default=None, 2744 description=( 2745 "Applied audience: actors pinned to the requested version, by tier. " 2746 "`None` when no specific version was requested." 2747 ), 2748 ) 2749 2750 2751class ConnectorVersionHealthSummary(BaseModel): 2752 """Four-bucket health rollup for the actors running a connector version.""" 2753 2754 connector_version_id: str = Field(description="The connector version UUID") 2755 connector_definition_id: str = Field(description="The connector definition UUID") 2756 docker_repository: str = Field(description="Docker repository path") 2757 docker_image_tag: str | None = Field( 2758 default=None, description="Docker image tag for this version" 2759 ) 2760 days: int = Field(description="Lookback window in days") 2761 customer_tier_filter: str = Field( 2762 description="Tier filter applied to the counts (`TIER_0`/`TIER_1`/`TIER_2`/`UNKNOWN`/`ALL`)" 2763 ) 2764 healthy: int = Field(description="Actors with at least one successful sync") 2765 unhealthy: int = Field( 2766 description="Actors with failures and no successes in the window" 2767 ) 2768 awaiting: int = Field( 2769 description="Actors that ran but produced only non-terminal jobs (no result yet)" 2770 ) 2771 disabled: int = Field( 2772 description=( 2773 "Actors pinned to the version that produced no jobs in the window — " 2774 "the dormant/inactive audience" 2775 ) 2776 ) 2777 total_actors: int = Field(description="Total actors counted across all states") 2778 healthy_by_tier: TierSummary = Field(description="Healthy actors by tier") 2779 unhealthy_by_tier: TierSummary = Field(description="Unhealthy actors by tier") 2780 awaiting_by_tier: TierSummary = Field(description="Awaiting-results actors by tier") 2781 disabled_by_tier: TierSummary = Field(description="Disabled actors by tier") 2782 2783 2784@mcp_tool( 2785 read_only=True, 2786 idempotent=True, 2787 open_world=True, 2788) 2789def query_connector_population_summary( 2790 connector_version_id: Annotated[ 2791 str | None, 2792 Field( 2793 description=( 2794 "Connector version UUID. When provided, the applied audience " 2795 "(`pinned_to_version`) is included. Provide this OR " 2796 "`connector_name` + `connector_version` OR a definition-level " 2797 "identifier." 2798 ), 2799 ), 2800 ] = None, 2801 connector_name: Annotated[ 2802 str | None, 2803 Field( 2804 description=( 2805 "Canonical connector name (e.g. `source-postgres`). Used with " 2806 "`connector_version` to resolve the version UUID." 2807 ), 2808 ), 2809 ] = None, 2810 connector_version: Annotated[ 2811 str | None, 2812 Field( 2813 description=( 2814 "Semver version tag (e.g. `0.3.59`). Used with `connector_name`." 2815 ), 2816 ), 2817 ] = None, 2818 connector_definition_id: Annotated[ 2819 str | None, 2820 Field( 2821 description=( 2822 "Connector definition UUID for a definition-level summary " 2823 "(no `pinned_to_version` breakdown)." 2824 ), 2825 ), 2826 ] = None, 2827 connector_canonical_name: Annotated[ 2828 str | None, 2829 Field( 2830 description=( 2831 "Canonical connector name resolved to a definition ID via the " 2832 "registry, for a definition-level summary." 2833 ), 2834 ), 2835 ] = None, 2836 customer_tier_filter: Annotated[ 2837 TierFilter, 2838 Field( 2839 description=( 2840 "Which customer tiers to count. Defaults to `TIER_2`; pass " 2841 "`ALL` to include TIER_0/TIER_1 (revenue-critical) customers." 2842 ), 2843 ), 2844 ] = "TIER_2", 2845) -> ConnectorPopulationSummary: 2846 """Summarize the applied vs potential pinning audience for a connector, by tier. 2847 2848 Answers "how many actors are pinned and how many are eligible for pinning, 2849 split by tier" — the population view analogous to a rollout's audience. 2850 2851 - `active`: the potential audience — non-tombstoned actors of the 2852 definition that have at least one *active* connection 2853 (`connection.status = 'active'`). Inactive/disabled and deprecated 2854 connections are excluded, so this reflects the enabled, rollout-touchable 2855 population rather than every actor ever created. 2856 - `pinned_any`: active actors that already have an effective 2857 `connector_version` pin at any scope (actor/workspace/org). 2858 - `eligible`: `active` minus `pinned_any` — actors available to pin. 2859 - `pinned_to_version`: the applied audience for the requested version 2860 (only when a version identifier was provided). 2861 2862 Backed by `scoped_configuration` + actor/connection tables, so it is cheap 2863 to compute. Accepts a version identifier (preferred, adds 2864 `pinned_to_version`) or a definition-level identifier. 2865 """ 2866 version_id, definition_id, docker_repository, docker_image_tag = ( 2867 _resolve_connector_target( 2868 connector_version_id=connector_version_id, 2869 connector_name=connector_name, 2870 connector_version=connector_version, 2871 connector_definition_id=connector_definition_id, 2872 connector_canonical_name=connector_canonical_name, 2873 ) 2874 ) 2875 is_destination = not is_source_connector(docker_repository) 2876 2877 population_rows = query_actor_population_by_org( 2878 actor_definition_id=definition_id, 2879 is_destination=is_destination, 2880 ) 2881 pinned_version_rows = ( 2882 query_actors_pinned_to_version(version_id) if version_id is not None else None 2883 ) 2884 summary = summarize_population( 2885 population_rows, 2886 pinned_version_rows=pinned_version_rows, 2887 tier_filter=customer_tier_filter, 2888 # The population query above is run without a `target_version_id`, so the 2889 # version-aware summaries do not apply (the per-version audience comes 2890 # from `pinned_version_rows` instead). 2891 has_target_version=False, 2892 ) 2893 return ConnectorPopulationSummary( 2894 connector_definition_id=definition_id, 2895 connector_version_id=version_id, 2896 docker_repository=docker_repository, 2897 docker_image_tag=docker_image_tag, 2898 customer_tier_filter=customer_tier_filter, 2899 active=summary.active_by_tier, 2900 pinned_any=summary.pinned_any_by_tier, 2901 eligible=summary.eligible_by_tier, 2902 pinned_to_version=summary.pinned_to_version_by_tier, 2903 ) 2904 2905 2906@mcp_tool( 2907 read_only=True, 2908 idempotent=True, 2909 open_world=True, 2910) 2911def query_connector_version_health_summary( 2912 connector_version_id: Annotated[ 2913 str | None, 2914 Field( 2915 description=( 2916 "Connector version UUID. Provide this OR " 2917 "`connector_name` + `connector_version`." 2918 ), 2919 ), 2920 ] = None, 2921 connector_name: Annotated[ 2922 str | None, 2923 Field( 2924 description=( 2925 "Canonical connector name (e.g. `source-postgres`). Used with " 2926 "`connector_version` to resolve the version UUID." 2927 ), 2928 ), 2929 ] = None, 2930 connector_version: Annotated[ 2931 str | None, 2932 Field( 2933 description=( 2934 "Semver version tag (e.g. `0.3.59`). Used with `connector_name`." 2935 ), 2936 ), 2937 ] = None, 2938 days: Annotated[ 2939 int, 2940 Field( 2941 description="Number of days to look back (default: 7, max: 30)", 2942 ge=1, 2943 le=30, 2944 ), 2945 ] = 7, 2946 include_pinned_disabled: Annotated[ 2947 bool, 2948 Field( 2949 description=( 2950 "If `True` (default), actors pinned to the version that ran no " 2951 "jobs in the window are counted as `disabled`." 2952 ), 2953 ), 2954 ] = True, 2955 customer_tier_filter: Annotated[ 2956 TierFilter, 2957 Field( 2958 description=( 2959 "Which customer tiers to count. Defaults to `TIER_2`; pass " 2960 "`ALL` to include TIER_0/TIER_1 (revenue-critical) customers." 2961 ), 2962 ), 2963 ] = "TIER_2", 2964) -> ConnectorVersionHealthSummary: 2965 """Summarize actor health for a connector version into four buckets. 2966 2967 Answers "how many actors on a version are healthy / unhealthy / awaiting / 2968 disabled". Classification per actor over the lookback window: 2969 2970 - `healthy`: at least one successful sync (the same success signal the 2971 autopilot health gate uses). 2972 - `unhealthy`: failures and no successes. 2973 - `awaiting`: ran but produced only non-terminal jobs (no result yet). 2974 - `disabled`: (when `include_pinned_disabled`) pinned to the version with no 2975 jobs at all in the window — the dormant/inactive audience. 2976 2977 Built on the attempt/version primitive — the version stamped into 2978 `jobs.config` at job-creation time — not the current pin state, so it 2979 reflects actors that actually ran this version. This scans jobs over the 2980 window and is more expensive than the population summary; keep `days` 2981 bounded and query per-version. 2982 """ 2983 if connector_version_id is None and not ( 2984 connector_name is not None and connector_version is not None 2985 ): 2986 raise PyAirbyteInputError( 2987 message=( 2988 "Provide either `connector_version_id` or both " 2989 "`connector_name` and `connector_version`." 2990 ), 2991 ) 2992 version_id, definition_id, docker_repository, docker_image_tag = ( 2993 _resolve_connector_target( 2994 connector_version_id=connector_version_id, 2995 connector_name=connector_name, 2996 connector_version=connector_version, 2997 connector_definition_id=None, 2998 connector_canonical_name=None, 2999 ) 3000 ) 3001 assert version_id is not None 3002 is_destination = not is_source_connector(docker_repository) 3003 3004 health_rows = query_version_actor_health( 3005 version_id, 3006 is_destination=is_destination, 3007 days=days, 3008 connector_definition_id=definition_id, 3009 ) 3010 pinned_actor_rows = ( 3011 query_actors_pinned_to_version(version_id) if include_pinned_disabled else None 3012 ) 3013 summary = summarize_version_health( 3014 health_rows, 3015 pinned_actor_rows=pinned_actor_rows, 3016 tier_filter=customer_tier_filter, 3017 ) 3018 return ConnectorVersionHealthSummary( 3019 connector_version_id=version_id, 3020 connector_definition_id=definition_id, 3021 docker_repository=docker_repository, 3022 docker_image_tag=docker_image_tag, 3023 days=days, 3024 customer_tier_filter=customer_tier_filter, 3025 healthy=summary.healthy, 3026 unhealthy=summary.unhealthy, 3027 awaiting=summary.awaiting, 3028 disabled=summary.disabled, 3029 total_actors=summary.total_actors, 3030 healthy_by_tier=summary.healthy_by_tier, 3031 unhealthy_by_tier=summary.unhealthy_by_tier, 3032 awaiting_by_tier=summary.awaiting_by_tier, 3033 disabled_by_tier=summary.disabled_by_tier, 3034 ) 3035 3036 3037def register_prod_db_ops_tools(app: FastMCP) -> None: 3038 """Register prod DB query tools with the FastMCP app.""" 3039 register_mcp_tools(app, mcp_module=__name__)