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 effective connector_version pin at any scope (actor/workspace/org).
  • eligible: active minus pinned_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: (when include_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, or connection_ids must be provided (any combination is accepted).
  • start_at and end_at must be timezone-aware and start_at < end_at.

Key fields in each row:

  • job_id, attempt_id, attempt_number
  • job_status, attempt_status
  • job_started_at, job_updated_at, attempt_ended_at
  • failure_summary (JSON; populated when an attempt failed)
  • connection_id, connection_name, connection_status
  • source_actor_id, source_actor_name, source_actor_definition_id
  • destination_actor_id, destination_actor_name, destination_actor_definition_id
  • workspace_id, workspace_name, organization_id
  • dataplane_group_id, dataplane_name
  • customer_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:

  1. Query recently released connectors to identify which ones to monitor
  2. Get aggregate stats showing how many connections are using each version
  3. 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__)