airbyte_ops_mcp.cli.cloud

CLI commands for Airbyte Cloud operations.

Commands:

airbyte-ops cloud connector get-version-info - Get connector version info airbyte-ops cloud connector set-version-override - Set connector version override airbyte-ops cloud connector clear-version-override - Clear connector version override airbyte-ops cloud connector regression-test - Run regression tests (single-version or comparison) airbyte-ops cloud connector fetch-connection-config - Fetch connection config to local file airbyte-ops cloud organization search - Search organizations by name/email airbyte-ops cloud organization info - Get organization info airbyte-ops cloud workspace search - Search workspaces by name/slug or email domain airbyte-ops cloud workspace info - Get workspace info

CLI reference

The commands below are regenerated by poe docs-generate via cyclopts's programmatic docs API; see docs/generate_cli.py.

airbyte-ops cloud COMMAND

Airbyte Cloud operations.

Commands:

  • connection: Connection operations in Airbyte Cloud.
  • connector: Deployed connector operations in Airbyte Cloud.
  • db: Database operations for Airbyte Cloud Prod DB Replica.
  • logs: GCP Cloud Logging operations for Airbyte Cloud.
  • organization: Organization operations in Airbyte Cloud.
  • workspace: Workspace operations in Airbyte Cloud.

airbyte-ops cloud connector

Deployed connector operations in Airbyte Cloud.

airbyte-ops cloud connector get-version-info

airbyte-ops cloud connector get-version-info WORKSPACE-ID CONNECTOR-ID CONNECTOR-TYPE

Get the current version information for a deployed connector.

Parameters:

  • WORKSPACE-ID, --workspace-id: The Airbyte Cloud workspace ID. [required]
  • CONNECTOR-ID, --connector-id: The ID of the deployed connector (source or destination). [required]
  • CONNECTOR-TYPE, --connector-type: The type of connector. [required] [choices: source, destination]

airbyte-ops cloud connector set-version-override

airbyte-ops cloud connector set-version-override VERSION REASON ISSUE-URL APPROVAL-COMMENT-URL [ARGS]

Set a version override for a deployed connector.

Requires admin authentication via AIRBYTE_INTERNAL_ADMIN_FLAG and AIRBYTE_INTERNAL_ADMIN_USER environment variables.

The --override-level flag selects the scope at which the pin is applied:

  • actor (default): pins a single deployed connector instance. Requires --workspace-id, --connector-id, and --connector-type.
  • workspace: pins ALL instances of a connector type within a workspace. Requires --workspace-id and --actor-definition-id.
  • organization: pins ALL instances of a connector type across an organization. Requires --organization-id and --actor-definition-id.

The customer_tier_filter gates the operation: the call fails if the actual tier of the target organization does not match. Use ALL to proceed regardless of tier (a warning is shown for sensitive tiers).

Parameters:

  • VERSION, --version: The semver version string to pin to (e.g., '2.1.5-preview.abc1234'). [required]
  • REASON, --reason: Explanation for the override (min 10 characters). [required]
  • ISSUE-URL, --issue-url: GitHub issue URL providing context for this operation. [required]
  • APPROVAL-COMMENT-URL, --approval-comment-url: Slack approval record URL where admin authorized this deployment. [required]
  • OVERRIDE-LEVEL, --override-level: Scope at which to apply the version override: 'actor' (single deployed connector instance), 'workspace' (all instances of a connector type within a workspace), or 'organization' (all instances across an organization). Defaults to 'actor'. [choices: actor, workspace, organization] [default: actor]
  • WORKSPACE-ID, --workspace-id: The Airbyte Cloud workspace ID. Required for 'actor' and 'workspace' override levels.
  • ORGANIZATION-ID, --organization-id: The Airbyte Cloud organization ID. Required for 'organization' override level.
  • CONNECTOR-ID, --connector-id: The ID of the deployed connector (source or destination). Required for 'actor' override level.
  • CONNECTOR-TYPE, --connector-type: The type of connector. Required for 'actor' override level. [choices: source, destination]
  • ACTOR-DEFINITION-ID, --actor-definition-id: The connector definition UUID. Required for 'workspace' and 'organization' override levels.
  • AI-AGENT-SESSION-URL, --ai-agent-session-url: URL to AI agent session driving this operation (for auditability).
  • REASON-URL, --reason-url: Optional URL with more context (e.g., issue link).
  • CUSTOMER-TIER-FILTER, --customer-tier-filter: Tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. The operation will be rejected if the actual customer tier does not match. Defaults to 'TIER_2' (non-sensitive customers). [choices: TIER_0, TIER_1, TIER_2, UNKNOWN, ALL] [default: TIER_2]

airbyte-ops cloud connector clear-version-override

airbyte-ops cloud connector clear-version-override ISSUE-URL APPROVAL-COMMENT-URL [ARGS]

Clear a version override from a deployed connector.

Requires admin authentication via AIRBYTE_INTERNAL_ADMIN_FLAG and AIRBYTE_INTERNAL_ADMIN_USER environment variables.

The --override-level flag selects the scope at which the pin is removed:

  • actor (default): clears a single deployed connector instance pin. Requires --workspace-id, --connector-id, and --connector-type.
  • workspace: clears the workspace-level pin for a connector type. Requires --workspace-id and --actor-definition-id.
  • organization: clears the organization-level pin for a connector type. Requires --organization-id and --actor-definition-id.

Parameters:

  • ISSUE-URL, --issue-url: GitHub issue URL providing context for this operation. [required]
  • APPROVAL-COMMENT-URL, --approval-comment-url: Slack approval record URL where admin authorized this deployment. [required]
  • OVERRIDE-LEVEL, --override-level: Scope at which to clear the version override: 'actor', 'workspace', or 'organization'. Defaults to 'actor'. [choices: actor, workspace, organization] [default: actor]
  • WORKSPACE-ID, --workspace-id: The Airbyte Cloud workspace ID. Required for 'actor' and 'workspace' override levels.
  • ORGANIZATION-ID, --organization-id: The Airbyte Cloud organization ID. Required for 'organization' override level.
  • CONNECTOR-ID, --connector-id: The ID of the deployed connector (source or destination). Required for 'actor' override level.
  • CONNECTOR-TYPE, --connector-type: The type of connector. Required for 'actor' override level. [choices: source, destination]
  • ACTOR-DEFINITION-ID, --actor-definition-id: The connector definition UUID. Required for 'workspace' and 'organization' override levels.
  • AI-AGENT-SESSION-URL, --ai-agent-session-url: URL to AI agent session driving this operation (for auditability).
  • CUSTOMER-TIER-FILTER, --customer-tier-filter: Tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. The operation will be rejected if the actual customer tier does not match. Defaults to 'TIER_2' (non-sensitive customers). [choices: TIER_0, TIER_1, TIER_2, UNKNOWN, ALL] [default: TIER_2]

airbyte-ops cloud connector regression-test

airbyte-ops cloud connector regression-test [ARGS]

Run regression tests on connectors.

This command supports two modes:

Comparison mode (skip_compare=False, default): Runs the specified Airbyte protocol command against both the target (new) and control (baseline) connector versions, then compares the results. This helps identify regressions between versions.

Single-version mode (skip_compare=True): Runs the specified Airbyte protocol command against a single connector and validates the output. No comparison is performed.

Results are written to the output directory and to GitHub Actions outputs if running in CI.

You can provide the test image in three ways:

  1. --test-image: Use a pre-built image from Docker registry
  2. --connector-name: Build the image locally from source code
  3. --connection-id: Auto-detect from an Airbyte Cloud connection

You can provide config/catalog either via file paths OR via a connection_id that fetches them from Airbyte Cloud.

Parameters:

  • SKIP-COMPARE, --skip-compare, --no-skip-compare: If True, skip comparison and run single-version tests only. If False (default), run comparison tests (target vs control). [default: False]
  • SKIP-RECORD-COMPARISON, --skip-record-comparison, --no-skip-record-comparison: If True, skip the record-level comparison for read commands (counts, PK presence, PK integrity, field values) and gate the verdict on command outcomes only. Escape hatch for sources whose data legitimately changes between the two live runs (rolling windows, feeds that age out records) until HTTP request caching makes strict comparison reliable. [default: False]
  • TEST-IMAGE, --test-image: Test connector image with tag (e.g., airbyte/source-github:1.0.0). This is the image under test - in comparison mode, it's compared against control_image.
  • CONTROL-IMAGE, --control-image: Control connector image (baseline version) with tag (e.g., airbyte/source-github:1.0.0). Ignored if skip_compare=True.
  • CONNECTOR-NAME, --connector-name: Connector name to build image from source (e.g., 'source-pokeapi'). If provided, builds the image locally with tag 'dev'. For comparison tests (default), this builds the target image. For single-version tests (skip_compare=True), this builds the test image.
  • REPO-ROOT, --repo-root: Path to the airbyte repo root. Required if connector_name is provided and the repo cannot be auto-detected.
  • COMMAND, --command: The Airbyte command to run. [choices: spec, check, discover, read] [default: check]
  • CONNECTION-ID, --connection-id: Airbyte Cloud connection ID to fetch config/catalog from. Mutually exclusive with config-path/catalog-path. If provided, test_image/control_image can be auto-detected.
  • CONFIG-PATH, --config-path: Path to the connector config JSON file.
  • CATALOG-PATH, --catalog-path: Path to the configured catalog JSON file (required for read).
  • STATE-PATH, --state-path: Path to the state JSON file (optional for read).
  • OUTPUT-DIR, --output-dir: Directory to store test artifacts. [default: /tmp/regression_test_artifacts]
  • ENABLE-HTTP-METRICS, --enable-http-metrics, --no-enable-http-metrics: Capture HTTP traffic metrics via mitmproxy. Requires mitmdump to be installed; without it the run proceeds without metrics. Only used in comparison mode, where HTTP replay turns it on anyway -- set this to record traffic while running both versions live (--no-http-replay). [default: False]
  • ENABLE-HTTP-REPLAY, --enable-http-replay, --no-http-replay: Replay the control run's recorded HTTP responses for the target run, so both versions see identical upstream data. Implies --enable-http-metrics. Defaults to True in comparison mode; has no effect with --skip-compare, which has no control run to replay from. Pass --no-http-replay to run both versions against the live API, which is the fallback whenever replay cannot be delivered anyway.
  • HTTP-REPLAY-MAX-BODY, --http-replay-max-body: Response bodies larger than this are streamed rather than recorded, and therefore re-fetched live instead of replayed. Raise it for a connector with large pages, at the cost of replay memory. Accepts a mitmproxy size such as '1m' or '10m' -- never a bare number, which means bytes. Unlike the other --http-replay-* options this one governs recording, so it applies with --enable-http-metrics alone. [default: 1m]
  • HTTP-REPLAY-MAX-DUMP-SIZE-MB, --http-replay-max-dump-size-mb: Skip replay and run the target live when the filtered replay corpus is larger than this, and discard a recorded dump this large once the run is done with it. mitmproxy holds the whole corpus in memory for the entire run, and a dump grows with everything the connector downloads, so this is what keeps a big connector from exhausting the CI runner's memory and disk. Dumps are never uploaded as artifacts, so on a local run -- where the dump is the only thing that can diagnose a coverage shortfall offline -- raise it to keep the dump. [default: 256]
  • HTTP-REPLAY-REUSE, --http-replay-reuse, --no-http-replay-reuse: Serve a recorded response more than once. By default each recorded flow is consumed on first match, which preserves recorded ordering and lets memory fall as the run proceeds; reuse keeps the whole corpus resident for the entire run. [default: False]
  • HTTP-REPLAY-IGNORE-PARAMS, --http-replay-ignore-params: Comma-separated query parameters to exclude from the replay match key, for connectors that put a timestamp or nonce in the URL.
  • STRICT-REPLAY, --strict-replay, --no-strict-replay: Kill requests that match nothing in the replay corpus instead of forwarding them upstream. Diagnostic only: it turns unmatched requests into connector errors, which is how you count them. mitmproxy stops intervening once its flowmap is empty, so combine this with --http-replay-reuse to cover the whole run rather than only the part before the corpus is consumed. [default: False]
  • SELECTED-STREAMS, --selected-streams: Comma-separated list of stream names to include in the read. Only these streams will be included in the configured catalog. This is useful to limit data volume by testing only specific streams.
  • ENABLE-DEBUG-LOGS, --enable-debug-logs, --no-enable-debug-logs: Enable debug-level logging for regression test output. Also passed as LOG_LEVEL=DEBUG to the connector Docker container. [default: False]
  • WITH-STATE, --with-state, --no-state: Fetch and pass the connection's current state to the read command, producing a warm read instead of a cold read. Defaults to True when --connection-id is provided, False otherwise. Has no effect unless the command is read. Ignored when --state-path is explicitly provided.

airbyte-ops cloud connector consolidate-regression-reports

airbyte-ops cloud connector consolidate-regression-reports ARTIFACTS-DIR [ARGS]

Fold one run's per-command regression reports into a single page.

The regression test runs the CLI once per Airbyte command, so each command writes its own report into its own artifact. This produces the run-level view -- a verdict table across commands plus every report inlined -- so a reviewer downloads one file instead of four.

Never fails the run: a command whose report is missing or unreadable is skipped, because this is a reporting convenience and the per-command artifacts remain the source of truth.

Parameters:

  • ARTIFACTS-DIR, --artifacts-dir: Directory holding one subdirectory per command, each with the report.html that regression-test --output-dir wrote. [required]
  • OUTPUT, --output: Where to write the consolidated report. Defaults to report.html in --artifacts-dir.

airbyte-ops cloud connector fetch-connection-config

airbyte-ops cloud connector fetch-connection-config CONNECTION-ID [ARGS]

Fetch connection configuration from Airbyte Cloud to a local file.

This command retrieves the source configuration for a given connection ID and writes it to a JSON file. When --output-path is omitted the file is written to the platform temp directory to avoid accidentally committing secrets to a git repository.

Requires authentication via AIRBYTE_CLOUD_CLIENT_ID and AIRBYTE_CLOUD_CLIENT_SECRET environment variables.

When --with-secrets is specified, the command fetches unmasked secrets from the internal database using the connection-retriever. This additionally requires:

  • An OC issue URL for audit logging (--oc-issue-url)
  • GCP credentials via GCP_PROD_DB_ACCESS_CREDENTIALS env var or gcloud auth application-default login
  • Cloud SQL Python Connector access to the Prod DB replica.

Parameters:

  • CONNECTION-ID, --connection-id: The UUID of the Airbyte Cloud connection. [required]
  • OUTPUT-PATH, --output-path: Path to output file or directory. If directory, writes connection--config.json inside it. Default: platform temp directory (e.g. /tmp/connection--config.json)
  • WITH-SECRETS, --with-secrets, --no-secrets: If set, fetches unmasked secrets from the internal database. Requires GCP_PROD_DB_ACCESS_CREDENTIALS env var or gcloud auth application-default login. Must be used with --oc-issue-url. [default: False]
  • OC-ISSUE-URL, --oc-issue-url: OC issue URL for audit logging. Required when using --with-secrets.

airbyte-ops cloud connector rollout

Connector rollout operations.

Commands:

  • autopilot: AutoPilot rollout orchestration commands.
  • list: List connector rollouts from the production database.
airbyte-ops cloud connector rollout autopilot

AutoPilot rollout orchestration commands.

airbyte-ops cloud connector rollout autopilot auto-start
airbyte-ops cloud connector rollout autopilot auto-start [ARGS]

Start INITIALIZED rollouts that have autopilot auto-start enabled.

Queries the prod DB for rollouts in initialized state, checks the connector's autopilotConfig.autoStart gate, and starts eligible rollouts via the Cloud Config API.

Parameters:

  • CONNECTOR, --connector: Filter to a specific connector by canonical name (e.g., 'source-github').
  • DRY-RUN, --dry-run, --no-dry-run: Preview actions without executing them. [default: False]
airbyte-ops cloud connector rollout autopilot auto-advance
airbyte-ops cloud connector rollout autopilot auto-advance [ARGS]

Advance IN_PROGRESS rollouts within their current tier.

Advances in_progress rollouts that haven't reached their target percentage based on the configured strategy pacing (fast/slow/default).

Also recovers workflow_started rollouts: a rollout confirmed to have zero eligible actors is finalized as succeeded (0-of-0 eligible is complete, not stuck — promoting the RC to GA) instead of being left wedged; one with a positive estimate is re-driven; an unavailable estimate is skipped for a later retry.

Parameters:

  • CONNECTOR, --connector: Filter to a specific connector by canonical name (e.g., 'source-github').
  • DRY-RUN, --dry-run, --no-dry-run: Preview actions without executing them. [default: False]
airbyte-ops cloud connector rollout autopilot auto-promote
airbyte-ops cloud connector rollout autopilot auto-promote [ARGS]

Promote rollouts when at 100% of current tier.

Checks the autopilotConfig.autoPromoteStages gate, then scans the tier order forward for the next tier that has eligible actors and starts a rollout there, skipping empty intermediate tiers. If no later customer tier has actors (or the current tier is ALL), the current tier is finalized as succeeded (GA promotion).

Parameters:

  • CONNECTOR, --connector: Filter to a specific connector by canonical name (e.g., 'source-github').
  • DRY-RUN, --dry-run, --no-dry-run: Preview actions without executing them. [default: False]
airbyte-ops cloud connector rollout autopilot auto-triage-failed
airbyte-ops cloud connector rollout autopilot auto-triage-failed [ARGS]

Triage failed rollouts: log all failures, check unpin eligibility.

Finds rollouts in errored or paused state. Logs every failure unconditionally. Checks safe-to-downgrade eligibility per unsafeDowngrades. Actor-level unpinning is not yet implemented.

Parameters:

  • CONNECTOR, --connector: Filter to a specific connector by canonical name (e.g., 'source-github').
  • DRY-RUN, --dry-run, --no-dry-run: Preview actions without executing them. [default: False]
airbyte-ops cloud connector rollout autopilot auto-close
airbyte-ops cloud connector rollout autopilot auto-close [ARGS]

Close obsolete rollouts so only a connector's active candidate remains.

Closes a rollout when a newer RC supersedes it, when its RC is already the registry GA default, or when its RC is not the highest advertised release candidate for the connector. Uses retain_pins_on_cancellation=True so pinned actors are left undisturbed.

Parameters:

  • CONNECTOR, --connector: Filter to a specific connector by canonical name (e.g., 'source-github').
  • DRY-RUN, --dry-run, --no-dry-run: Preview actions without executing them. [default: False]
airbyte-ops cloud connector rollout autopilot auto-rollback-failed
airbyte-ops cloud connector rollout autopilot auto-rollback-failed [ARGS]

Full rollback/cancel of failed rollouts when safe to downgrade.

Finds rollouts in errored state. Checks the unsafeDowngrades gate. If the version is safe to downgrade, calls finalize_connector_rollout with state failed_rolled_back to cancel the entire rollout.

Parameters:

  • CONNECTOR, --connector: Filter to a specific connector by canonical name (e.g., 'source-github').
  • DRY-RUN, --dry-run, --no-dry-run: Preview actions without executing them. [default: False]
airbyte-ops cloud connector rollout list
airbyte-ops cloud connector rollout list [OPTIONS]

List connector rollouts from the production database.

Parameters:

  • --include-inactive, --no-include-inactive: Include terminal (canceled/succeeded/errored) rollouts. Requires --limit. [default: False]
  • --limit: Maximum number of rollouts to return. Required when --include-inactive is set. [default: 0]

airbyte-ops cloud db

Database operations for Airbyte Cloud Prod DB Replica.

airbyte-ops cloud db start-proxy

airbyte-ops cloud db start-proxy [ARGS]

Start the Cloud SQL Proxy for database access.

This command starts the Cloud SQL Auth Proxy to enable connections to the Airbyte Cloud Prod DB Replica. The proxy is required for database query tools.

By default, runs as a daemon (background process). Use --no-daemon to run in foreground mode where you can see logs and stop with Ctrl+C.

Credentials are read from the GCP_PROD_DB_ACCESS_CREDENTIALS environment variable, which should contain the service account JSON credentials.

After starting the proxy, set these environment variables to use database tools: export USE_CLOUD_SQL_PROXY=1 export DB_PORT={port}

Parameters:

  • PORT, --port: Port for the Cloud SQL Proxy to listen on. [default: 15432]
  • DAEMON, --daemon, --no-daemon: Run as daemon in background (default). Use --no-daemon for foreground. [default: True]

airbyte-ops cloud db stop-proxy

airbyte-ops cloud db stop-proxy

Stop the Cloud SQL Proxy daemon.

This command stops a Cloud SQL Proxy that was started with 'start-proxy'. It reads the PID from the PID file and sends a SIGTERM signal to stop the process.

airbyte-ops cloud connection

Connection operations in Airbyte Cloud.

airbyte-ops cloud connection state

Connection state operations in Airbyte Cloud.

Commands:

  • get: Get the current state for an Airbyte Cloud connection.
  • reset: Reset a configured stream's state so the next sync full-refreshes it.
  • set: Set the state for an Airbyte Cloud connection.
airbyte-ops cloud connection state get
airbyte-ops cloud connection state get CONNECTION-ID [ARGS]

Get the current state for an Airbyte Cloud connection.

Parameters:

  • CONNECTION-ID, --connection-id: The connection ID (UUID) to fetch state for. [required]
  • STREAM, --stream: Optional stream name to filter state for a single stream.
  • NAMESPACE, --namespace: Optional stream namespace to narrow the stream filter.
  • OUTPUT, --output: Path to write the state JSON output to a file.
airbyte-ops cloud connection state set
airbyte-ops cloud connection state set CONNECTION-ID [ARGS]

Set the state for an Airbyte Cloud connection.

State JSON can be provided as a positional argument, via --input , or piped through STDIN.

Uses the safe variant that prevents updates while a sync is running. When --stream is provided, only that stream's state is updated within the existing connection state.

Parameters:

  • CONNECTION-ID, --connection-id: The connection ID (UUID) to update state for. [required]
  • STATE-JSON, --state-json: The connection state as a JSON string. When --stream is used, this is just the stream's state blob (e.g., '{"cursor": "2024-01-01"}'). Otherwise, must include 'stateType', 'connectionId', and the appropriate state field. Can also be provided via --input or STDIN.
  • STREAM, --stream: Optional stream name to update state for a single stream only.
  • NAMESPACE, --namespace: Optional stream namespace to identify the stream.
  • INPUT, --input: Path to a JSON file containing the state to set.
airbyte-ops cloud connection state reset
airbyte-ops cloud connection state reset CONNECTION-ID STREAM [ARGS]

Reset a configured stream's state so the next sync full-refreshes it.

Uses the safe variant that prevents updates while a sync is running. Returns a restorable previous_state_backup in raw Config API format.

Parameters:

  • CONNECTION-ID, --connection-id: The connection ID (UUID) to update state for. [required]
  • STREAM, --stream: The configured stream name whose state should be reset. [required]
  • NAMESPACE, --namespace: Optional stream namespace to identify the stream.

airbyte-ops cloud connection catalog

Connection catalog operations in Airbyte Cloud.

Commands:

  • get: Get the configured catalog for an Airbyte Cloud connection.
  • set: Set the configured catalog for an Airbyte Cloud connection.
airbyte-ops cloud connection catalog get
airbyte-ops cloud connection catalog get CONNECTION-ID [ARGS]

Get the configured catalog for an Airbyte Cloud connection.

Parameters:

  • CONNECTION-ID, --connection-id: The connection ID (UUID) to fetch catalog for. [required]
  • OUTPUT, --output: Path to write the catalog JSON output to a file.
airbyte-ops cloud connection catalog set
airbyte-ops cloud connection catalog set CONNECTION-ID [ARGS]

Set the configured catalog for an Airbyte Cloud connection.

Catalog JSON can be provided as a positional argument, via --input , or piped through STDIN.

WARNING: This replaces the entire configured catalog.

Parameters:

  • CONNECTION-ID, --connection-id: The connection ID (UUID) to update catalog for. [required]
  • CATALOG-JSON, --catalog-json: The configured catalog as a JSON string. Can also be provided via --input or STDIN.
  • INPUT, --input: Path to a JSON file containing the catalog to set.

airbyte-ops cloud logs

GCP Cloud Logging operations for Airbyte Cloud.

airbyte-ops cloud logs lookup-cloud-backend-error

airbyte-ops cloud logs lookup-cloud-backend-error ERROR-ID [ARGS]

Look up error details from GCP Cloud Logging by error ID.

When an Airbyte Cloud API returns an error response with only an error ID (e.g., {"errorId": "3173452e-8f22-4286-a1ec-b0f16c1e078a"}), this command fetches the full stack trace and error details from GCP Cloud Logging.

Requires GCP credentials with Logs Viewer role on the target project. Set up credentials with: gcloud auth application-default login

Parameters:

  • ERROR-ID, --error-id: The error ID (UUID) to search for. This is typically returned in API error responses as {'errorId': '...'} [required]
  • LOOKBACK-DAYS, --lookback-days: Number of days to look back in logs. [default: 7]
  • MIN-SEVERITY-FILTER, --min-severity-filter: Optional minimum severity level to filter logs. [choices: debug, info, notice, warning, error, critical, alert, emergency]
  • RAW, --raw, --no-raw: Output raw JSON instead of formatted text. [default: False]

airbyte-ops cloud organization

Organization operations in Airbyte Cloud.

airbyte-ops cloud organization search NAME-CONTAINS [ARGS]

Search organizations by name or email substring.

Parameters:

  • NAME-CONTAINS, --name-contains: Case-insensitive substring to match against organization name or email. [required]
  • LIMIT, --limit: Maximum number of results to return. [default: 20]

airbyte-ops cloud organization info

airbyte-ops cloud organization info ORGANIZATION-ID

Get information about an organization (stub — not yet implemented).

Parameters:

  • ORGANIZATION-ID, --organization-id: The organization UUID. [required]

airbyte-ops cloud workspace

Workspace operations in Airbyte Cloud.

airbyte-ops cloud workspace search [ARGS]

Search workspaces by name substring or email domain.

Parameters:

  • NAME-CONTAINS, --name-contains: Case-insensitive substring to match against workspace name or slug.
  • EMAIL-DOMAIN, --email-domain: Email domain to search for (e.g. 'motherduck.com').
  • LIMIT, --limit: Maximum number of results to return. [default: 100]

airbyte-ops cloud workspace info

airbyte-ops cloud workspace info WORKSPACE-ID

Get information about a workspace.

Parameters:

  • WORKSPACE-ID, --workspace-id: The workspace UUID. [required]
   1# Copyright (c) 2025 Airbyte, Inc., all rights reserved.
   2"""CLI commands for Airbyte Cloud operations.
   3
   4Commands:
   5    airbyte-ops cloud connector get-version-info - Get connector version info
   6    airbyte-ops cloud connector set-version-override - Set connector version override
   7    airbyte-ops cloud connector clear-version-override - Clear connector version override
   8    airbyte-ops cloud connector regression-test - Run regression tests (single-version or comparison)
   9    airbyte-ops cloud connector fetch-connection-config - Fetch connection config to local file
  10    airbyte-ops cloud organization search - Search organizations by name/email
  11    airbyte-ops cloud organization info - Get organization info
  12    airbyte-ops cloud workspace search - Search workspaces by name/slug or email domain
  13    airbyte-ops cloud workspace info - Get workspace info
  14
  15## CLI reference
  16
  17The commands below are regenerated by `poe docs-generate` via cyclopts's
  18programmatic docs API; see `docs/generate_cli.py`.
  19
  20.. include:: ../../../docs/generated/cli/cloud.md
  21   :start-line: 2
  22"""
  23
  24from __future__ import annotations
  25
  26# Hide Python-level members from the pdoc page for this module; the rendered
  27# docs for this CLI group come entirely from the grafted `.. include::` in
  28# the module docstring above.
  29__all__: list[str] = []
  30
  31import json
  32import logging
  33import os
  34import shutil
  35import signal
  36import socket
  37import subprocess
  38import sys
  39import tempfile
  40import time
  41from pathlib import Path
  42from typing import Annotated, Any, Literal
  43
  44import requests
  45import yaml
  46from airbyte.cloud.workspaces import CloudWorkspace
  47from airbyte.exceptions import PyAirbyteInputError
  48from airbyte_cdk.models.connector_metadata import MetadataFile
  49from airbyte_cdk.utils.connector_paths import find_connector_root_from_name
  50from airbyte_cdk.utils.docker import build_connector_image, verify_docker_installation
  51from airbyte_protocol.models import ConfiguredAirbyteCatalog
  52from cyclopts import Parameter
  53from fastmcp_extensions.cli import (
  54    exit_with_error,
  55    print_error,
  56    print_json,
  57    print_success,
  58    print_warning,
  59)
  60
  61from airbyte_ops_mcp.cli._base import App, app
  62from airbyte_ops_mcp.cloud_admin.auth import CloudAuthError
  63from airbyte_ops_mcp.cloud_admin.connection_config import fetch_connection_config
  64from airbyte_ops_mcp.cloud_admin.connection_state import (
  65    reset_stream_state as reset_cloud_stream_state,
  66)
  67from airbyte_ops_mcp.cloud_admin.registry_lookup import (
  68    resolve_definition_id_to_canonical_info,
  69)
  70from airbyte_ops_mcp.cloud_admin.version_overrides import (
  71    ResolvedCloudAuth,
  72    VersionOverrideTarget,
  73    get_connector_version_info,
  74)
  75from airbyte_ops_mcp.cloud_admin.version_overrides import (
  76    set_version_override as set_cloud_version_override,
  77)
  78from airbyte_ops_mcp.constants import (
  79    CLOUD_SQL_INSTANCE,
  80    CLOUD_SQL_PROXY_PID_FILE,
  81    DEFAULT_CLOUD_SQL_PROXY_PORT,
  82    ENV_GCP_PROD_DB_ACCESS_CREDENTIALS,
  83)
  84from airbyte_ops_mcp.gcp_logs import GCPSeverity, fetch_error_logs
  85from airbyte_ops_mcp.prod_db_access.queries import (
  86    query_workspace_info,
  87    query_workspaces_by_email_domain,
  88)
  89from airbyte_ops_mcp.prod_db_access.queries import (
  90    search_organizations as _search_organizations,
  91)
  92from airbyte_ops_mcp.prod_db_access.queries import (
  93    search_workspaces as _search_workspaces,
  94)
  95from airbyte_ops_mcp.regression_tests.cdk_secrets import get_first_config_from_secrets
  96from airbyte_ops_mcp.regression_tests.ci_output import (
  97    ComparisonOutcome,
  98    command_result_succeeded,
  99    evaluate_comparison_outcome,
 100    generate_regression_report,
 101    generate_single_version_report,
 102    write_github_output,
 103    write_github_outputs,
 104    write_github_summary,
 105    write_json_output,
 106)
 107from airbyte_ops_mcp.regression_tests.config_overrides import (
 108    filter_configured_catalog_file,
 109)
 110from airbyte_ops_mcp.regression_tests.connection_fetcher import (
 111    enrich_catalog_schemas_from_cloud,
 112    fetch_connection_data,
 113    fetch_connection_state,
 114    save_connection_data_to_files,
 115    streams_missing_schemas,
 116)
 117from airbyte_ops_mcp.regression_tests.connection_secret_retriever import (
 118    SecretRetrievalError,
 119    enrich_catalog_schemas_from_platform_db,
 120    enrich_config_with_secrets,
 121    should_use_secret_retriever,
 122)
 123from airbyte_ops_mcp.regression_tests.connector_runner import (
 124    ConnectorRunner,
 125    ensure_image_available,
 126)
 127from airbyte_ops_mcp.regression_tests.http_metrics import (
 128    DEFAULT_MAX_REPLAY_DUMP_SIZE_MB,
 129    DEFAULT_REPLAY_EXTRA,
 130    DEFAULT_STREAM_LARGE_BODIES,
 131    MitmproxyManager,
 132    ReplayOptions,
 133    discard_oversized_dump,
 134    parse_http_dump,
 135)
 136from airbyte_ops_mcp.regression_tests.models import (
 137    Command,
 138    ComparableOutputs,
 139    ConnectorUnderTest,
 140    ExecutionInputs,
 141    TargetOrControl,
 142    duplicate_stream_names,
 143    get_primary_keys_per_stream,
 144    split_records_per_stream,
 145)
 146from airbyte_ops_mcp.regression_tests.regression import (
 147    ComparisonResult,
 148    RecordComparisonSummary,
 149    compare_catalog_schemas,
 150    compare_final_states,
 151    compare_specs,
 152    run_record_comparisons_from_files,
 153)
 154from airbyte_ops_mcp.regression_tests.report import (
 155    HTML_REPORT_FILENAME,
 156    CheckResult,
 157    DiffBlock,
 158    DiffLine,
 159    build_comparison_report_model,
 160    build_diff_block,
 161    build_json_block,
 162    build_single_version_report_model,
 163    consolidate_reports,
 164    failed_checks,
 165    render_inline_summary,
 166    write_html_report,
 167)
 168from airbyte_ops_mcp.telemetry import track_regression_test
 169from airbyte_ops_mcp.tier_cache import resolve_workspace
 170
 171# Path to connectors directory within the airbyte repo
 172CONNECTORS_SUBDIR = Path("airbyte-integrations") / "connectors"
 173
 174# Create the cloud sub-app
 175cloud_app = App(name="cloud", help="Airbyte Cloud operations.")
 176app.command(cloud_app)
 177
 178# Create the connector sub-app under cloud
 179connector_app = App(
 180    name="connector", help="Deployed connector operations in Airbyte Cloud."
 181)
 182cloud_app.command(connector_app)
 183
 184# Create the db sub-app under cloud
 185db_app = App(name="db", help="Database operations for Airbyte Cloud Prod DB Replica.")
 186cloud_app.command(db_app)
 187
 188# Create the connection sub-app under cloud
 189connection_app = App(name="connection", help="Connection operations in Airbyte Cloud.")
 190cloud_app.command(connection_app)
 191
 192# Create the state sub-app under connection
 193state_app = App(name="state", help="Connection state operations in Airbyte Cloud.")
 194connection_app.command(state_app)
 195
 196# Create the catalog sub-app under connection
 197catalog_app = App(
 198    name="catalog", help="Connection catalog operations in Airbyte Cloud."
 199)
 200connection_app.command(catalog_app)
 201
 202# Create the logs sub-app under cloud
 203logs_app = App(name="logs", help="GCP Cloud Logging operations for Airbyte Cloud.")
 204cloud_app.command(logs_app)
 205
 206# Create the organization sub-app under cloud
 207organization_app = App(
 208    name="organization", help="Organization operations in Airbyte Cloud."
 209)
 210cloud_app.command(organization_app)
 211
 212# Create the workspace sub-app under cloud
 213workspace_app = App(name="workspace", help="Workspace operations in Airbyte Cloud.")
 214cloud_app.command(workspace_app)
 215
 216
 217@db_app.command(name="start-proxy")
 218def start_proxy(
 219    port: Annotated[
 220        int,
 221        Parameter(help="Port for the Cloud SQL Proxy to listen on."),
 222    ] = DEFAULT_CLOUD_SQL_PROXY_PORT,
 223    daemon: Annotated[
 224        bool,
 225        Parameter(
 226            help="Run as daemon in background (default). Use --no-daemon for foreground."
 227        ),
 228    ] = True,
 229) -> None:
 230    """Start the Cloud SQL Proxy for database access.
 231
 232    This command starts the Cloud SQL Auth Proxy to enable connections to the
 233    Airbyte Cloud Prod DB Replica. The proxy is required for database query tools.
 234
 235    By default, runs as a daemon (background process). Use --no-daemon to run in
 236    foreground mode where you can see logs and stop with Ctrl+C.
 237
 238    Credentials are read from the GCP_PROD_DB_ACCESS_CREDENTIALS environment variable,
 239    which should contain the service account JSON credentials.
 240
 241    After starting the proxy, set these environment variables to use database tools:
 242        export USE_CLOUD_SQL_PROXY=1
 243        export DB_PORT={port}
 244
 245    Example:
 246        airbyte-ops cloud db start-proxy
 247        airbyte-ops cloud db start-proxy --port 15432
 248        airbyte-ops cloud db start-proxy --no-daemon
 249    """
 250    # Check if proxy is already running on the requested port (idempotency)
 251    try:
 252        with socket.create_connection(("127.0.0.1", port), timeout=0.5):
 253            # Something is already listening on this port
 254            pid_file = Path(CLOUD_SQL_PROXY_PID_FILE)
 255            pid_info = ""
 256            if pid_file.exists():
 257                pid_info = f" (PID: {pid_file.read_text().strip()})"
 258            print_success(
 259                f"Cloud SQL Proxy is already running on port {port}{pid_info}"
 260            )
 261            print_success("")
 262            print_success("To use database tools, set these environment variables:")
 263            print_success("  export USE_CLOUD_SQL_PROXY=1")
 264            print_success(f"  export DB_PORT={port}")
 265            return
 266    except (OSError, TimeoutError, ConnectionRefusedError):
 267        pass  # Port not in use, proceed with starting proxy
 268
 269    # Check if cloud-sql-proxy is installed
 270    proxy_path = shutil.which("cloud-sql-proxy")
 271    if not proxy_path:
 272        exit_with_error(
 273            "cloud-sql-proxy not found in PATH. "
 274            "Install it from: https://cloud.google.com/sql/docs/mysql/sql-proxy"
 275        )
 276
 277    # Get credentials from environment
 278    creds_json = os.getenv(ENV_GCP_PROD_DB_ACCESS_CREDENTIALS)
 279    if not creds_json:
 280        exit_with_error(
 281            f"{ENV_GCP_PROD_DB_ACCESS_CREDENTIALS} environment variable is not set. "
 282            "This should contain the GCP service account JSON credentials."
 283        )
 284
 285    # Build the command using --json-credentials to avoid writing to disk
 286    cmd = [
 287        proxy_path,
 288        CLOUD_SQL_INSTANCE,
 289        f"--port={port}",
 290        f"--json-credentials={creds_json}",
 291    ]
 292
 293    print_success(f"Starting Cloud SQL Proxy on port {port}...")
 294    print_success(f"Instance: {CLOUD_SQL_INSTANCE}")
 295    print_success("")
 296    print_success("To use database tools, set these environment variables:")
 297    print_success("  export USE_CLOUD_SQL_PROXY=1")
 298    print_success(f"  export DB_PORT={port}")
 299    print_success("")
 300
 301    if daemon:
 302        # Run in background (daemon mode) with log file for diagnostics
 303        log_file_path = Path("/tmp/airbyte-cloud-sql-proxy.log")
 304        log_file = log_file_path.open("ab")
 305        process = subprocess.Popen(
 306            cmd,
 307            stdout=subprocess.DEVNULL,
 308            stderr=log_file,
 309            start_new_session=True,
 310        )
 311
 312        # Brief wait to verify the process started successfully
 313        time.sleep(0.5)
 314        if process.poll() is not None:
 315            # Process exited immediately - read any error output
 316            log_file.close()
 317            error_output = ""
 318            if log_file_path.exists():
 319                error_output = log_file_path.read_text()[-1000:]  # Last 1000 chars
 320            exit_with_error(
 321                f"Cloud SQL Proxy failed to start (exit code: {process.returncode}).\n"
 322                f"Check logs at {log_file_path}\n"
 323                f"Recent output: {error_output}"
 324            )
 325
 326        # Write PID to file for stop-proxy command
 327        pid_file = Path(CLOUD_SQL_PROXY_PID_FILE)
 328        pid_file.write_text(str(process.pid))
 329        print_success(f"Cloud SQL Proxy started as daemon (PID: {process.pid})")
 330        print_success(f"Logs: {log_file_path}")
 331        print_success("To stop: airbyte-ops cloud db stop-proxy")
 332    else:
 333        # Run in foreground - replace current process
 334        # Signals (Ctrl+C) will be handled directly by the cloud-sql-proxy process
 335        print_success("Running in foreground. Press Ctrl+C to stop the proxy.")
 336        print_success("")
 337        os.execv(proxy_path, cmd)
 338
 339
 340@db_app.command(name="stop-proxy")
 341def stop_proxy() -> None:
 342    """Stop the Cloud SQL Proxy daemon.
 343
 344    This command stops a Cloud SQL Proxy that was started with 'start-proxy'.
 345    It reads the PID from the PID file and sends a SIGTERM signal to stop the process.
 346
 347    Example:
 348        airbyte-ops cloud db stop-proxy
 349    """
 350    pid_file = Path(CLOUD_SQL_PROXY_PID_FILE)
 351
 352    if not pid_file.exists():
 353        exit_with_error(
 354            f"PID file not found at {CLOUD_SQL_PROXY_PID_FILE}. "
 355            "No Cloud SQL Proxy daemon appears to be running."
 356        )
 357
 358    pid_str = pid_file.read_text().strip()
 359    if not pid_str.isdigit():
 360        pid_file.unlink()
 361        exit_with_error(f"Invalid PID in {CLOUD_SQL_PROXY_PID_FILE}: {pid_str}")
 362
 363    pid = int(pid_str)
 364
 365    # Check if process is still running
 366    try:
 367        os.kill(pid, 0)  # Signal 0 just checks if process exists
 368    except ProcessLookupError:
 369        pid_file.unlink()
 370        print_success(
 371            f"Cloud SQL Proxy (PID: {pid}) is not running. Cleaned up PID file."
 372        )
 373        return
 374    except PermissionError:
 375        exit_with_error(f"Permission denied to check process {pid}.")
 376
 377    # Send SIGTERM to stop the process
 378    try:
 379        os.kill(pid, signal.SIGTERM)
 380        print_success(f"Sent SIGTERM to Cloud SQL Proxy (PID: {pid}).")
 381    except ProcessLookupError:
 382        print_success(f"Cloud SQL Proxy (PID: {pid}) already stopped.")
 383    except PermissionError:
 384        exit_with_error(f"Permission denied to stop process {pid}.")
 385
 386    # Clean up PID file
 387    pid_file.unlink(missing_ok=True)
 388    print_success("Cloud SQL Proxy stopped.")
 389
 390
 391@connector_app.command(name="get-version-info")
 392def get_version_info(
 393    workspace_id: Annotated[
 394        str,
 395        Parameter(help="The Airbyte Cloud workspace ID."),
 396    ],
 397    connector_id: Annotated[
 398        str,
 399        Parameter(help="The ID of the deployed connector (source or destination)."),
 400    ],
 401    connector_type: Annotated[
 402        Literal["source", "destination"],
 403        Parameter(help="The type of connector."),
 404    ],
 405) -> None:
 406    """Get the current version information for a deployed connector."""
 407    result = get_connector_version_info(
 408        auth=_resolve_cli_cloud_auth(),
 409        workspace_id=workspace_id,
 410        actor_id=connector_id,
 411        actor_type=connector_type,
 412    )
 413    print_json(result.model_dump())
 414
 415
 416_OverrideLevel = Literal["actor", "workspace", "organization"]
 417_TierFilter = Literal["TIER_0", "TIER_1", "TIER_2", "UNKNOWN", "ALL"]
 418
 419
 420def _resolve_cli_cloud_auth() -> ResolvedCloudAuth:
 421    """Resolve Airbyte Cloud auth from environment variables for CLI use.
 422
 423    Mirrors the priority order the MCP server uses:
 424
 425    1. `AIRBYTE_CLOUD_BEARER_TOKEN`
 426    2. `AIRBYTE_CLOUD_CLIENT_ID` + `AIRBYTE_CLOUD_CLIENT_SECRET`
 427
 428    Raises a CLI error via `exit_with_error` if no usable credentials are present.
 429    """
 430    bearer_token = os.environ.get("AIRBYTE_CLOUD_BEARER_TOKEN")
 431    if bearer_token:
 432        return ResolvedCloudAuth(bearer_token=bearer_token)
 433
 434    client_id = os.environ.get("AIRBYTE_CLOUD_CLIENT_ID")
 435    client_secret = os.environ.get("AIRBYTE_CLOUD_CLIENT_SECRET")
 436    if client_id and client_secret:
 437        return ResolvedCloudAuth(client_id=client_id, client_secret=client_secret)
 438
 439    exit_with_error(
 440        "Missing Airbyte Cloud credentials. Set AIRBYTE_CLOUD_BEARER_TOKEN, or "
 441        "both AIRBYTE_CLOUD_CLIENT_ID and AIRBYTE_CLOUD_CLIENT_SECRET."
 442    )
 443
 444
 445def _validate_override_scope_args(
 446    *,
 447    override_level: _OverrideLevel,
 448    workspace_id: str | None,
 449    organization_id: str | None,
 450    connector_id: str | None,
 451    connector_type: Literal["source", "destination"] | None,
 452    actor_definition_id: str | None,
 453) -> None:
 454    """Validate scope-specific arguments for `set/clear-version-override`.
 455
 456    Calls `exit_with_error` (which terminates the process) if the supplied
 457    arguments are inconsistent with the requested `override_level`.
 458    """
 459    if override_level == "actor":
 460        missing = [
 461            name
 462            for name, value in (
 463                ("--workspace-id", workspace_id),
 464                ("--connector-id", connector_id),
 465                ("--connector-type", connector_type),
 466            )
 467            if not value
 468        ]
 469        if missing:
 470            exit_with_error(
 471                "Override level 'actor' requires " + ", ".join(missing) + "."
 472            )
 473        if actor_definition_id is not None:
 474            exit_with_error(
 475                "Override level 'actor' must not be combined with "
 476                "--actor-definition-id."
 477            )
 478        if organization_id is not None:
 479            exit_with_error(
 480                "Override level 'actor' must not be combined with --organization-id."
 481            )
 482        return
 483
 484    if override_level == "workspace":
 485        missing = [
 486            name
 487            for name, value in (
 488                ("--workspace-id", workspace_id),
 489                ("--actor-definition-id", actor_definition_id),
 490            )
 491            if not value
 492        ]
 493        if missing:
 494            exit_with_error(
 495                "Override level 'workspace' requires " + ", ".join(missing) + "."
 496            )
 497        if connector_id is not None or connector_type is not None:
 498            exit_with_error(
 499                "Override level 'workspace' must not be combined with "
 500                "--connector-id or --connector-type."
 501            )
 502        if organization_id is not None:
 503            exit_with_error(
 504                "Override level 'workspace' must not be combined with "
 505                "--organization-id."
 506            )
 507        return
 508
 509    # override_level == "organization"
 510    missing = [
 511        name
 512        for name, value in (
 513            ("--organization-id", organization_id),
 514            ("--actor-definition-id", actor_definition_id),
 515        )
 516        if not value
 517    ]
 518    if missing:
 519        exit_with_error(
 520            "Override level 'organization' requires " + ", ".join(missing) + "."
 521        )
 522    if connector_id is not None or connector_type is not None:
 523        exit_with_error(
 524            "Override level 'organization' must not be combined with "
 525            "--connector-id or --connector-type."
 526        )
 527    if workspace_id is not None:
 528        exit_with_error(
 529            "Override level 'organization' must not be combined with --workspace-id."
 530        )
 531
 532
 533def _dispatch_version_override(
 534    *,
 535    override_level: _OverrideLevel,
 536    workspace_id: str | None,
 537    organization_id: str | None,
 538    connector_id: str | None,
 539    connector_type: Literal["source", "destination"] | None,
 540    actor_definition_id: str | None,
 541    version: str | None,
 542    unset: bool,
 543    reason: str | None,
 544    reason_url: str | None,
 545    issue_url: str,
 546    approval_comment_url: str,
 547    ai_agent_session_url: str | None,
 548    customer_tier_filter: _TierFilter,
 549) -> None:
 550    """Route a version-override request through the normalized core helper.
 551
 552    Resolves `--actor-definition-id` to its canonical name and connector type
 553    via the cloud registry for `workspace`/`organization` scopes. Prints the
 554    operation result as JSON.
 555    """
 556    _validate_override_scope_args(
 557        override_level=override_level,
 558        workspace_id=workspace_id,
 559        organization_id=organization_id,
 560        connector_id=connector_id,
 561        connector_type=connector_type,
 562        actor_definition_id=actor_definition_id,
 563    )
 564
 565    auth = _resolve_cli_cloud_auth()
 566
 567    if override_level == "actor":
 568        assert workspace_id is not None
 569        assert connector_id is not None
 570        assert connector_type is not None
 571        assert organization_id is None
 572        try:
 573            ws_resolution = resolve_workspace(workspace_id)
 574            if not ws_resolution.organization_id:
 575                exit_with_error("Could not resolve organization for workspace.")
 576            result = set_cloud_version_override(
 577                auth=auth,
 578                target=VersionOverrideTarget(
 579                    scope="actor",
 580                    organization_id=ws_resolution.organization_id,
 581                    workspace_id=workspace_id,
 582                    actor_id=connector_id,
 583                    connector_type=connector_type,
 584                ),
 585                approval_comment_url=approval_comment_url,
 586                version=version,
 587                unset=unset,
 588                override_reason=reason,
 589                override_reason_reference_url=reason_url,
 590                issue_url=issue_url,
 591                ai_agent_session_url=ai_agent_session_url,
 592                customer_tier_filter=customer_tier_filter,
 593            )
 594        except (PyAirbyteInputError, CloudAuthError) as e:
 595            exit_with_error(str(e))
 596        _print_override_result(result.success, result.message)
 597        print_json(result.model_dump())
 598        return
 599
 600    try:
 601        assert actor_definition_id is not None
 602        canonical_name, resolved_connector_type = (
 603            resolve_definition_id_to_canonical_info(actor_definition_id)
 604        )
 605    except (PyAirbyteInputError, requests.RequestException) as e:
 606        exit_with_error(f"Failed to resolve --actor-definition-id: {e}")
 607    typed_connector_type: Literal["source", "destination"] = (
 608        "source" if resolved_connector_type == "source" else "destination"
 609    )
 610
 611    if override_level == "workspace":
 612        assert workspace_id is not None
 613        try:
 614            ws_resolution = resolve_workspace(workspace_id)
 615            if not ws_resolution.organization_id:
 616                exit_with_error("Could not resolve organization for workspace.")
 617            target = VersionOverrideTarget(
 618                scope="workspace",
 619                organization_id=ws_resolution.organization_id,
 620                workspace_id=workspace_id,
 621                connector_name=canonical_name,
 622                connector_type=typed_connector_type,
 623            )
 624        except PyAirbyteInputError as e:
 625            exit_with_error(str(e))
 626    else:
 627        assert organization_id is not None
 628        target = VersionOverrideTarget(
 629            scope="organization",
 630            organization_id=organization_id,
 631            connector_name=canonical_name,
 632            connector_type=typed_connector_type,
 633        )
 634
 635    try:
 636        result = set_cloud_version_override(
 637            auth=auth,
 638            target=target,
 639            approval_comment_url=approval_comment_url,
 640            version=version,
 641            unset=unset,
 642            override_reason=reason,
 643            override_reason_reference_url=reason_url,
 644            issue_url=issue_url,
 645            ai_agent_session_url=ai_agent_session_url,
 646            customer_tier_filter=customer_tier_filter,
 647        )
 648    except (PyAirbyteInputError, CloudAuthError) as e:
 649        exit_with_error(str(e))
 650    _print_override_result(result.success, result.message)
 651    print_json(result.model_dump())
 652
 653
 654def _print_override_result(success: bool, message: str) -> None:
 655    """Emit an override-operation status message via the shared CLI helpers."""
 656    if success:
 657        print_success(message)
 658    else:
 659        print_error(message)
 660
 661
 662@connector_app.command(name="set-version-override")
 663def set_version_override(
 664    version: Annotated[
 665        str,
 666        Parameter(
 667            help="The semver version string to pin to (e.g., '2.1.5-preview.abc1234')."
 668        ),
 669    ],
 670    reason: Annotated[
 671        str,
 672        Parameter(help="Explanation for the override (min 10 characters)."),
 673    ],
 674    issue_url: Annotated[
 675        str,
 676        Parameter(help="GitHub issue URL providing context for this operation."),
 677    ],
 678    approval_comment_url: Annotated[
 679        str,
 680        Parameter(
 681            help="Slack approval record URL where admin authorized this deployment."
 682        ),
 683    ],
 684    override_level: Annotated[
 685        _OverrideLevel,
 686        Parameter(
 687            help=(
 688                "Scope at which to apply the version override: 'actor' (single "
 689                "deployed connector instance), 'workspace' (all instances of a "
 690                "connector type within a workspace), or 'organization' (all "
 691                "instances across an organization). Defaults to 'actor'."
 692            ),
 693        ),
 694    ] = "actor",
 695    workspace_id: Annotated[
 696        str | None,
 697        Parameter(
 698            help=(
 699                "The Airbyte Cloud workspace ID. Required for 'actor' and "
 700                "'workspace' override levels."
 701            )
 702        ),
 703    ] = None,
 704    organization_id: Annotated[
 705        str | None,
 706        Parameter(
 707            help=(
 708                "The Airbyte Cloud organization ID. Required for "
 709                "'organization' override level."
 710            )
 711        ),
 712    ] = None,
 713    connector_id: Annotated[
 714        str | None,
 715        Parameter(
 716            help=(
 717                "The ID of the deployed connector (source or destination). "
 718                "Required for 'actor' override level."
 719            )
 720        ),
 721    ] = None,
 722    connector_type: Annotated[
 723        Literal["source", "destination"] | None,
 724        Parameter(help="The type of connector. Required for 'actor' override level."),
 725    ] = None,
 726    actor_definition_id: Annotated[
 727        str | None,
 728        Parameter(
 729            help=(
 730                "The connector definition UUID. Required for 'workspace' and "
 731                "'organization' override levels."
 732            )
 733        ),
 734    ] = None,
 735    ai_agent_session_url: Annotated[
 736        str | None,
 737        Parameter(
 738            help="URL to AI agent session driving this operation (for auditability)."
 739        ),
 740    ] = None,
 741    reason_url: Annotated[
 742        str | None,
 743        Parameter(help="Optional URL with more context (e.g., issue link)."),
 744    ] = None,
 745    customer_tier_filter: Annotated[
 746        _TierFilter,
 747        Parameter(
 748            help=(
 749                "Tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. "
 750                "The operation will be rejected if the actual customer tier does not match. "
 751                "Defaults to 'TIER_2' (non-sensitive customers)."
 752            ),
 753        ),
 754    ] = "TIER_2",
 755) -> None:
 756    """Set a version override for a deployed connector.
 757
 758    Requires admin authentication via `AIRBYTE_INTERNAL_ADMIN_FLAG` and
 759    `AIRBYTE_INTERNAL_ADMIN_USER` environment variables.
 760
 761    The `--override-level` flag selects the scope at which the pin is applied:
 762
 763    - `actor` (default): pins a single deployed connector instance. Requires
 764      `--workspace-id`, `--connector-id`, and `--connector-type`.
 765    - `workspace`: pins ALL instances of a connector type within a workspace.
 766      Requires `--workspace-id` and `--actor-definition-id`.
 767    - `organization`: pins ALL instances of a connector type across an
 768      organization. Requires `--organization-id` and `--actor-definition-id`.
 769
 770    The `customer_tier_filter` gates the operation: the call fails if the
 771    actual tier of the target organization does not match. Use `ALL` to
 772    proceed regardless of tier (a warning is shown for sensitive tiers).
 773    """
 774    _dispatch_version_override(
 775        override_level=override_level,
 776        workspace_id=workspace_id,
 777        organization_id=organization_id,
 778        connector_id=connector_id,
 779        connector_type=connector_type,
 780        actor_definition_id=actor_definition_id,
 781        version=version,
 782        unset=False,
 783        reason=reason,
 784        reason_url=reason_url,
 785        issue_url=issue_url,
 786        approval_comment_url=approval_comment_url,
 787        ai_agent_session_url=ai_agent_session_url,
 788        customer_tier_filter=customer_tier_filter,
 789    )
 790
 791
 792@connector_app.command(name="clear-version-override")
 793def clear_version_override(
 794    issue_url: Annotated[
 795        str,
 796        Parameter(help="GitHub issue URL providing context for this operation."),
 797    ],
 798    approval_comment_url: Annotated[
 799        str,
 800        Parameter(
 801            help="Slack approval record URL where admin authorized this deployment."
 802        ),
 803    ],
 804    override_level: Annotated[
 805        _OverrideLevel,
 806        Parameter(
 807            help=(
 808                "Scope at which to clear the version override: 'actor', "
 809                "'workspace', or 'organization'. Defaults to 'actor'."
 810            ),
 811        ),
 812    ] = "actor",
 813    workspace_id: Annotated[
 814        str | None,
 815        Parameter(
 816            help=(
 817                "The Airbyte Cloud workspace ID. Required for 'actor' and "
 818                "'workspace' override levels."
 819            )
 820        ),
 821    ] = None,
 822    organization_id: Annotated[
 823        str | None,
 824        Parameter(
 825            help=(
 826                "The Airbyte Cloud organization ID. Required for "
 827                "'organization' override level."
 828            )
 829        ),
 830    ] = None,
 831    connector_id: Annotated[
 832        str | None,
 833        Parameter(
 834            help=(
 835                "The ID of the deployed connector (source or destination). "
 836                "Required for 'actor' override level."
 837            )
 838        ),
 839    ] = None,
 840    connector_type: Annotated[
 841        Literal["source", "destination"] | None,
 842        Parameter(help="The type of connector. Required for 'actor' override level."),
 843    ] = None,
 844    actor_definition_id: Annotated[
 845        str | None,
 846        Parameter(
 847            help=(
 848                "The connector definition UUID. Required for 'workspace' and "
 849                "'organization' override levels."
 850            )
 851        ),
 852    ] = None,
 853    ai_agent_session_url: Annotated[
 854        str | None,
 855        Parameter(
 856            help="URL to AI agent session driving this operation (for auditability)."
 857        ),
 858    ] = None,
 859    customer_tier_filter: Annotated[
 860        _TierFilter,
 861        Parameter(
 862            help=(
 863                "Tier filter: 'TIER_0', 'TIER_1', 'TIER_2', 'UNKNOWN', or 'ALL'. "
 864                "The operation will be rejected if the actual customer tier does not match. "
 865                "Defaults to 'TIER_2' (non-sensitive customers)."
 866            ),
 867        ),
 868    ] = "TIER_2",
 869) -> None:
 870    """Clear a version override from a deployed connector.
 871
 872    Requires admin authentication via `AIRBYTE_INTERNAL_ADMIN_FLAG` and
 873    `AIRBYTE_INTERNAL_ADMIN_USER` environment variables.
 874
 875    The `--override-level` flag selects the scope at which the pin is removed:
 876
 877    - `actor` (default): clears a single deployed connector instance pin.
 878      Requires `--workspace-id`, `--connector-id`, and `--connector-type`.
 879    - `workspace`: clears the workspace-level pin for a connector type.
 880      Requires `--workspace-id` and `--actor-definition-id`.
 881    - `organization`: clears the organization-level pin for a connector type.
 882      Requires `--organization-id` and `--actor-definition-id`.
 883    """
 884    _dispatch_version_override(
 885        override_level=override_level,
 886        workspace_id=workspace_id,
 887        organization_id=organization_id,
 888        connector_id=connector_id,
 889        connector_type=connector_type,
 890        actor_definition_id=actor_definition_id,
 891        version=None,
 892        unset=True,
 893        reason=None,
 894        reason_url=None,
 895        issue_url=issue_url,
 896        approval_comment_url=approval_comment_url,
 897        ai_agent_session_url=ai_agent_session_url,
 898        customer_tier_filter=customer_tier_filter,
 899    )
 900
 901
 902def _load_json_file(file_path: Path) -> dict | None:
 903    """Load a JSON file and return its contents.
 904
 905    Returns None if the file doesn't exist or contains invalid JSON.
 906    """
 907    if not file_path.exists():
 908        return None
 909    try:
 910        return json.loads(file_path.read_text())
 911    except json.JSONDecodeError as e:
 912        print_error(f"Failed to parse JSON in file: {file_path}\nError: {e}")
 913        return None
 914
 915
 916def _run_connector_command(
 917    connector_image: str,
 918    command: Command,
 919    output_dir: Path,
 920    target_or_control: TargetOrControl,
 921    config_path: Path | None = None,
 922    catalog_path: Path | None = None,
 923    state_path: Path | None = None,
 924    proxy_url: str | None = None,
 925    ca_cert_path: Path | None = None,
 926    enable_debug_logs: bool = False,
 927) -> tuple[dict, ComparableOutputs]:
 928    """Run a connector command and return its results and its declared objects.
 929
 930    The dict is a flattened view of `ExecutionResult`: counts, exit status and
 931    artifact paths, which is all the reports need. What a comparison needs
 932    instead -- the spec, the discovered catalog -- travels beside it as
 933    `ComparableOutputs` rather than growing a dict that several consumers read.
 934    The `ExecutionResult` itself is still dropped here: it holds every parsed
 935    message, and a comparison keeps both sides alive at once.
 936
 937    Args:
 938        connector_image: Full connector image name with tag.
 939        command: The Airbyte command to run.
 940        output_dir: Directory to store output files.
 941        target_or_control: Whether this is target or control version.
 942        config_path: Path to connector config JSON file.
 943        catalog_path: Path to configured catalog JSON file.
 944        state_path: Path to state JSON file.
 945        proxy_url: Optional HTTP proxy URL for traffic capture.
 946        ca_cert_path: Path to the proxy's CA certificate, so the container trusts
 947            it. Without it every HTTPS request through the proxy fails TLS
 948            verification and the capture comes back near-empty.
 949        enable_debug_logs: If `True`, pass `LOG_LEVEL=DEBUG` to the connector container.
 950
 951    Returns:
 952        The flattened results, and the objects the connector declared.
 953    """
 954    connector = ConnectorUnderTest.from_image_name(connector_image, target_or_control)
 955
 956    config = _load_json_file(config_path) if config_path else None
 957    state = _load_json_file(state_path) if state_path else None
 958
 959    configured_catalog = None
 960    if catalog_path and catalog_path.exists():
 961        catalog_json = catalog_path.read_text()
 962        configured_catalog = ConfiguredAirbyteCatalog.parse_raw(catalog_json)
 963
 964    # Pass log level as an environment variable to the connector container
 965    container_env: dict[str, str] = {}
 966    if enable_debug_logs:
 967        container_env["LOG_LEVEL"] = "DEBUG"
 968
 969    execution_inputs = ExecutionInputs(
 970        connector_under_test=connector,
 971        command=command,
 972        output_dir=output_dir,
 973        config=config,
 974        configured_catalog=configured_catalog,
 975        state=state,
 976        environment_variables=container_env or None,
 977    )
 978
 979    runner = ConnectorRunner(
 980        execution_inputs, proxy_url=proxy_url, ca_cert_path=ca_cert_path
 981    )
 982    result = runner.run()
 983
 984    result.save_artifacts(output_dir)
 985
 986    result_dict: dict[str, object] = {
 987        "connector": connector_image,
 988        "command": command.value,
 989        "success": result.success,
 990        "exit_code": result.exit_code,
 991        "stdout_file": str(result.stdout_file_path),
 992        "stderr_file": str(result.stderr_file_path),
 993        "message_counts": {
 994            k.value: v for k, v in result.get_message_count_per_type().items()
 995        },
 996        "record_counts_per_stream": result.get_record_count_per_stream(),
 997    }
 998
 999    # For CHECK, include the parsed connectionStatus from the Airbyte protocol.
1000    # The CDK exits 0 regardless of check outcome; the real pass/fail signal
1001    # lives in the CONNECTION_STATUS message (status=SUCCEEDED or FAILED).
1002    if command == Command.CHECK:
1003        conn_status = result.get_connection_status()
1004        if conn_status is not None:
1005            result_dict["connection_status"] = conn_status.status.value
1006            result_dict["connection_status_message"] = conn_status.message or ""
1007        else:
1008            result_dict["connection_status"] = None
1009            result_dict["connection_status_message"] = ""
1010
1011    return result_dict, ComparableOutputs.from_execution_result(result)
1012
1013
1014def _annotate_run_verdict(command: str, result: dict[str, Any]) -> dict[str, Any]:
1015    """Copy a run result with `success` as the verdict, exit status kept separately.
1016
1017    `ExecutionResult.success` is the container exit status, which for `check` says
1018    nothing about whether the connector could connect -- the CDK exits 0 even when a
1019    check fails. Any payload read as the verdict must therefore not expose the exit
1020    status under `success`, at any nesting level. The raw value stays available as
1021    `exited_cleanly`; for every command but `check` the two are equal.
1022
1023    Returns a copy: callers pass the unannotated originals to the report generator,
1024    which reads `success` as the exit status.
1025    """
1026    return {
1027        **result,
1028        "success": command_result_succeeded(command, result),
1029        "exited_cleanly": result["success"],
1030    }
1031
1032
1033# The checks the output comparisons contribute, named for the row a reviewer
1034# reads in the report and in the step summary.
1035_SPEC_CHECK = "Spec compatibility"
1036_CATALOG_SCHEMA_CHECK = "Catalog schema"
1037_FINAL_STATE_CHECK = "Final state per stream"
1038
1039# How many individual changes a check's one-line summary names before it says
1040# how many more there are. The full list is in the report and in the run's JSON
1041# payload, and the comparators order their findings worst-first, so the ones
1042# named here are the ones worth naming.
1043_MAX_LISTED_CHANGES = 3
1044
1045# How many changes the JSON payload lists per comparison. Same reason as
1046# `_MAX_DIFF_PAYLOAD_CHARS`: the payload is written to `GITHUB_OUTPUT`, and a
1047# reordered `oneOf` on a wide spec produces hundreds of sentences.
1048_MAX_PAYLOAD_CHANGES = 100
1049
1050# How many per-stream entries it lists. A wide source has thousands of streams,
1051# and one line each is enough to exceed the same output budget without a single
1052# diff being attached.
1053_MAX_PAYLOAD_STREAMS = 100
1054
1055# How much of a comparison's detail the run's JSON payload carries. That payload
1056# is written to `GITHUB_OUTPUT`, which is size-limited, and a wide connector's
1057# catalog diff has no natural bound -- so the payload keeps as many stream diffs
1058# as fit and says how many it dropped. The catalogs themselves are saved
1059# artifacts, which is where an unbounded diff belongs.
1060_MAX_DIFF_PAYLOAD_CHARS = 50_000
1061
1062
1063def _compare_outputs(
1064    *,
1065    command: str,
1066    outcome: ComparisonOutcome,
1067    control_outputs: ComparableOutputs,
1068    target_outputs: ComparableOutputs,
1069) -> tuple[tuple[CheckResult, ...], tuple[DiffBlock, ...], dict[str, Any]]:
1070    """Compare what the two versions produced, for the commands that produce it.
1071
1072    This is the content check for `spec`, `discover` and `read`: without it their
1073    verdict is the exit code, and a connector that drops a config field, changes
1074    a stream's schema or corrupts the state the next sync resumes from exits 0
1075    all the way to a release. `check` reports its own status, so it is the one
1076    command with nothing to diff here.
1077
1078    Only runs when both commands succeeded. A failed run has nothing to compare,
1079    and reporting "no spec to compare" beside the command failure that caused it
1080    would state the same problem twice, in weaker terms the second time.
1081
1082    A comparator that raises is left to propagate. The comparators normalise what
1083    they are given, so the only ways here are a programming error in a future
1084    caller or a schema nested past the recursion limit -- and the workflow
1085    already has the right vocabulary for both: a step that dies is an
1086    `internal_failure`, reported as infrastructure rather than as a regression,
1087    with the connector's artifacts uploaded either way. Catching it here would
1088    only let a tooling bug read as "the new version broke its spec".
1089
1090    Args:
1091        command: The Airbyte command that was run, as the CLI names it. This is
1092            always `"read"` for a read, including one given a state file: the
1093            `READ` -> `READ_WITH_STATE` upgrade only changes the `Command` the
1094            container is run with.
1095        outcome: The already-evaluated command outcome.
1096        control_outputs: What the control version produced.
1097        target_outputs: What the target version produced.
1098
1099    Returns:
1100        The checks to fold into the verdict, the diff blocks that show what they
1101        found, and the comparison detail for the run's JSON payload -- all empty
1102        when there was nothing to compare.
1103    """
1104    if not outcome.success:
1105        return (), (), {}
1106
1107    if command == "spec":
1108        return _spec_comparison(control_outputs.spec, target_outputs.spec)
1109
1110    if command == "discover":
1111        return _catalog_comparison(control_outputs.catalog, target_outputs.catalog)
1112
1113    if command == "read":
1114        return _final_state_comparison(
1115            control_outputs.final_states, target_outputs.final_states
1116        )
1117
1118    return (), (), {}
1119
1120
1121def _spec_comparison(
1122    control_spec: Any, target_spec: Any
1123) -> tuple[tuple[CheckResult, ...], tuple[DiffBlock, ...], dict[str, Any]]:
1124    """Run the spec comparison and shape it for every surface."""
1125    result = compare_specs(control_spec, target_spec)
1126
1127    check = CheckResult(
1128        name=_SPEC_CHECK,
1129        passed=result.passed,
1130        summary=_change_summary(
1131            result.message,
1132            # A passing spec check still says what was added: an additive change
1133            # is the thing a reviewer most wants to see confirmed.
1134            result.errors if result.errors else result.warnings,
1135        ),
1136        details=tuple(result.errors),
1137    )
1138    # The incompatible changes are the check's `details`, which the report renders
1139    # under the check row and the step summary lists -- a block would repeat them.
1140    # The compatible ones have nowhere else to go, and a reviewer asking "what
1141    # else moved?" is asking about exactly those.
1142    blocks = tuple(
1143        block
1144        for block in (
1145            _change_block("Spec — compatible changes", result.warnings, "add"),
1146        )
1147        if block is not None
1148    )
1149
1150    return (check,), blocks, {"spec": _spec_detail(result)}
1151
1152
1153def _catalog_comparison(
1154    control_catalog: Any, target_catalog: Any
1155) -> tuple[tuple[CheckResult, ...], tuple[DiffBlock, ...], dict[str, Any]]:
1156    """Run the discovered-catalog comparison and shape it for every surface."""
1157    result = compare_catalog_schemas(control_catalog, target_catalog)
1158
1159    check = CheckResult(
1160        name=_CATALOG_SCHEMA_CHECK,
1161        passed=result.passed,
1162        # A catalog that only grew -- new streams, new fields, nothing removed
1163        # or retyped -- is reported and does not gate the release. Connectors
1164        # add streams constantly, and a check that goes red on every one of them
1165        # is a check people learn to click past.
1166        severity="diagnostic" if result.additive_only else "strict",
1167        summary=_change_summary(result.message, result.errors),
1168        details=tuple(result.errors),
1169    )
1170
1171    # The findings themselves are the check's `details`, which the report renders
1172    # under the check row: a stream that is missing or new has no schema diff, so
1173    # that list is the only place it appears.
1174    blocks = _stream_diff_blocks(
1175        result,
1176        block_title="Schema diff",
1177        omission_title="Schema diffs omitted for size",
1178        artifact_hint=(
1179            "Both discovered catalogs are in the artifact, under "
1180            "airbyte_messages/catalog.jsonl."
1181        ),
1182    )
1183
1184    return (check,), blocks, {"catalog_schema": _catalog_schema_detail(result)}
1185
1186
1187def _final_state_comparison(
1188    control_states: dict[tuple[str | None, str], Any] | None,
1189    target_states: dict[tuple[str | None, str], Any] | None,
1190) -> tuple[tuple[CheckResult, ...], tuple[DiffBlock, ...], dict[str, Any]]:
1191    """Run the final-state comparison and shape it for every surface."""
1192    result = compare_final_states(control_states, target_states)
1193
1194    # Structural findings first, then the value-only ones: the check lists
1195    # everything it found, and the summary quotes the front of that list.
1196    findings = [*result.errors, *result.warnings]
1197
1198    check = CheckResult(
1199        name=_FINAL_STATE_CHECK,
1200        passed=result.passed,
1201        # Two ways a non-passing state check does not gate. A cursor that only
1202        # advanced: until HTTP replay lands, the two versions call the live API
1203        # in separate runs and a timestamp cursor moves between them on its own.
1204        # And a read with no state on either side, where there was nothing to
1205        # compare rather than something wrong.
1206        severity=(
1207            "diagnostic" if result.value_only or result.inconclusive else "strict"
1208        ),
1209        summary=_change_summary(result.message, findings),
1210        details=tuple(findings),
1211    )
1212
1213    blocks = _stream_diff_blocks(
1214        result,
1215        block_title="Final state diff",
1216        omission_title="Final state diffs omitted for size",
1217        artifact_hint=(
1218            "Both versions' state messages are in the artifact, under "
1219            "airbyte_messages/state.jsonl."
1220        ),
1221    )
1222
1223    return (check,), blocks, {"final_state": _final_state_detail(result)}
1224
1225
1226def _stream_diff_blocks(
1227    result: ComparisonResult,
1228    *,
1229    block_title: str,
1230    omission_title: str,
1231    artifact_hint: str,
1232) -> tuple[DiffBlock, ...]:
1233    """One block per changed stream, within the report's diff budget.
1234
1235    In comparison order -- only a *changed* stream has a diff, so there is no
1236    severity left to sort by here -- and within the same budget the payload uses:
1237    a report that inlines every schema of a 200-stream source is not a report
1238    anyone reads. What the budget drops is named rather than skipped, since a
1239    block that is simply absent reads as "there was nothing to see".
1240    """
1241    blocks: list[DiffBlock] = []
1242    budget = _MAX_DIFF_PAYLOAD_CHARS
1243    dropped: list[str] = []
1244
1245    for name in result.failed_streams:
1246        stream = result.stream_results[name]
1247        if not stream.schema_diff:
1248            continue
1249
1250        payload = json.dumps(stream.schema_diff, indent=2, sort_keys=True)
1251        if len(payload) > budget:
1252            dropped.append(name)
1253            continue
1254
1255        budget -= len(payload)
1256        blocks.append(build_json_block(f"{block_title} — {name}", payload))
1257
1258    if dropped:
1259        blocks.append(
1260            build_diff_block(
1261                omission_title,
1262                [
1263                    DiffLine(
1264                        kind="meta",
1265                        text=(
1266                            "These streams changed, and their diffs were too "
1267                            f"large for this report. {artifact_hint}"
1268                        ),
1269                    ),
1270                    *(DiffLine(kind="del", text=name) for name in dropped),
1271                ],
1272                language="text",
1273            )
1274        )
1275
1276    return tuple(blocks)
1277
1278
1279def _change_block(title: str, changes: list[str], kind: str) -> DiffBlock | None:
1280    """One diff block listing a comparison's findings, or `None` if there are none."""
1281    if not changes:
1282        return None
1283
1284    return build_diff_block(
1285        title,
1286        [DiffLine(kind=kind, text=change) for change in changes],
1287        language="text",
1288    )
1289
1290
1291def _change_summary(message: str, changes: list[str]) -> str:
1292    """One line: what the comparison concluded, and the first few reasons."""
1293    if not changes:
1294        return message
1295
1296    listed = "; ".join(changes[:_MAX_LISTED_CHANGES])
1297    remaining = len(changes) - _MAX_LISTED_CHANGES
1298    if remaining > 0:
1299        listed = f"{listed}; and {remaining} more"
1300
1301    return f"{message}: {listed}"
1302
1303
1304def _bounded_changes(changes: list[str]) -> list[str]:
1305    """A change list the JSON payload can carry, saying what it left out."""
1306    if len(changes) <= _MAX_PAYLOAD_CHANGES:
1307        return changes
1308
1309    dropped = len(changes) - _MAX_PAYLOAD_CHANGES
1310
1311    return [
1312        *changes[:_MAX_PAYLOAD_CHANGES],
1313        f"... and {dropped} more, in the report and the console output",
1314    ]
1315
1316
1317def _spec_detail(result: ComparisonResult) -> dict[str, Any]:
1318    """The spec comparison, for the run's JSON payload."""
1319    return {
1320        "passed": result.passed,
1321        "message": result.message,
1322        "breaking_changes": _bounded_changes(result.errors),
1323        "compatible_changes": _bounded_changes(result.warnings),
1324    }
1325
1326
1327def _catalog_schema_detail(result: ComparisonResult) -> dict[str, Any]:
1328    """The catalog comparison, for the run's JSON payload."""
1329    streams, omitted, dropped = _bounded_stream_diffs(
1330        result,
1331        omission_note="Too large for this payload; diff the catalogs",
1332        artifact_hint="Both discovered catalogs are in the saved artifacts.",
1333    )
1334
1335    detail: dict[str, Any] = {
1336        "passed": result.passed,
1337        # Why a `passed: false` did not fail the run: everything found was an
1338        # addition, which is reported as a warning rather than a regression.
1339        "additive_only": result.additive_only,
1340        "message": result.message,
1341        "changed_streams": _bounded_changes(result.failed_streams),
1342        "streams": streams,
1343    }
1344    if omitted:
1345        detail["streams_with_omitted_diffs"] = omitted
1346    if dropped:
1347        detail["streams_omitted"] = dropped
1348
1349    return detail
1350
1351
1352def _final_state_detail(result: ComparisonResult) -> dict[str, Any]:
1353    """The final-state comparison, for the run's JSON payload."""
1354    streams, omitted, dropped = _bounded_stream_diffs(
1355        result,
1356        omission_note="Too large for this payload; diff the state messages",
1357        artifact_hint="Both versions' state messages are in the saved artifacts.",
1358    )
1359
1360    detail: dict[str, Any] = {
1361        "passed": result.passed,
1362        # The two reasons a `passed: false` did not fail the run, kept apart so
1363        # the payload never says "every state that moved moved in value only"
1364        # about a read where nothing moved because nothing was compared.
1365        "value_only": result.value_only,
1366        "inconclusive": result.inconclusive,
1367        "message": result.message,
1368        "changed_streams": _bounded_changes(result.failed_streams),
1369        "streams": streams,
1370    }
1371    if omitted:
1372        detail["streams_with_omitted_diffs"] = omitted
1373    if dropped:
1374        detail["streams_omitted"] = dropped
1375
1376    return detail
1377
1378
1379def _bounded_stream_diffs(
1380    result: ComparisonResult,
1381    *,
1382    omission_note: str,
1383    artifact_hint: str,
1384) -> tuple[dict[str, Any], list[str], int]:
1385    """A comparison's per-stream verdicts and diffs, sized for the JSON payload.
1386
1387    Every stream that changed is named with its verdict -- so the payload says
1388    what moved, not just that something did -- while the diffs themselves are
1389    kept only while they fit in `_MAX_DIFF_PAYLOAD_CHARS`, and the number of
1390    entries is capped in turn: a thousand-stream source would otherwise put a
1391    thousand of them into `GITHUB_OUTPUT`. Both bounds say what they dropped,
1392    on the console and in the payload: a bound nobody mentions reads as "there
1393    was nothing more to see".
1394    """
1395    streams: dict[str, Any] = {}
1396    omitted: list[str] = []
1397    budget = _MAX_DIFF_PAYLOAD_CHARS
1398
1399    # Non-passing first, so the cap falls on the streams with nothing to say.
1400    ordered = sorted(result.stream_results.items(), key=lambda item: item[1].passed)
1401    kept, dropped_entries = (
1402        ordered[:_MAX_PAYLOAD_STREAMS],
1403        ordered[_MAX_PAYLOAD_STREAMS:],
1404    )
1405
1406    for name, stream in kept:
1407        entry: dict[str, Any] = {"passed": stream.passed, "message": stream.message}
1408        streams[name] = entry
1409
1410        if not stream.schema_diff:
1411            continue
1412
1413        rendered = len(json.dumps(stream.schema_diff, separators=(",", ":")))
1414        if rendered > budget:
1415            omitted.append(name)
1416            entry["diff_omitted"] = omission_note
1417            continue
1418
1419        budget -= rendered
1420        entry["diff"] = stream.schema_diff
1421
1422    if omitted:
1423        print_warning(
1424            f"Comparison detail for {len(omitted)} streams was too large for the "
1425            f"run's JSON output ({', '.join(omitted)}); their diffs are omitted "
1426            f"there. {artifact_hint}"
1427        )
1428
1429    if dropped_entries:
1430        print_warning(
1431            f"The run's JSON output lists {_MAX_PAYLOAD_STREAMS} of "
1432            f"{len(ordered)} streams, changed ones first. The report carries "
1433            f"them all. {artifact_hint}"
1434        )
1435
1436    return streams, omitted, len(dropped_entries)
1437
1438
1439def _build_connector_image_from_source(
1440    connector_name: str,
1441    repo_root: Path | None = None,
1442    tag: str = "dev",
1443) -> str | None:
1444    """Build a connector image from source code.
1445
1446    Args:
1447        connector_name: Name of the connector (e.g., 'source-pokeapi').
1448        repo_root: Optional path to the airbyte repo root. If not provided,
1449            will attempt to auto-detect from current directory.
1450        tag: Tag to apply to the built image (default: 'dev').
1451
1452    Returns:
1453        The full image name with tag if successful, None if build fails.
1454    """
1455    if not verify_docker_installation():
1456        print_error("Docker is not installed or not running")
1457        return None
1458
1459    try:
1460        connector_directory = find_connector_root_from_name(connector_name)
1461    except FileNotFoundError:
1462        if repo_root:
1463            connector_directory = repo_root / CONNECTORS_SUBDIR / connector_name
1464            if not connector_directory.exists():
1465                print_error(f"Connector directory not found: {connector_directory}")
1466                return None
1467        else:
1468            print_error(
1469                f"Could not find connector '{connector_name}'. "
1470                "Try providing --repo-root to specify the airbyte repo location."
1471            )
1472            return None
1473
1474    metadata_file_path = connector_directory / "metadata.yaml"
1475    if not metadata_file_path.exists():
1476        print_error(f"metadata.yaml not found at {metadata_file_path}")
1477        return None
1478
1479    metadata = MetadataFile.from_file(metadata_file_path)
1480    print_success(f"Building image for connector: {connector_name}")
1481
1482    built_image = build_connector_image(
1483        connector_name=connector_name,
1484        connector_directory=connector_directory,
1485        metadata=metadata,
1486        tag=tag,
1487        no_verify=False,
1488    )
1489    print_success(f"Successfully built image: {built_image}")
1490    return built_image
1491
1492
1493def _fetch_control_image_from_metadata(connector_name: str) -> str | None:
1494    """Fetch the current released connector image from metadata.yaml on main branch.
1495
1496    This fetches the connector's metadata.yaml from the airbyte monorepo's master branch
1497    and extracts the dockerRepository and dockerImageTag to construct the control image.
1498
1499    Args:
1500        connector_name: The connector name (e.g., 'source-github').
1501
1502    Returns:
1503        The full connector image with tag (e.g., 'airbyte/source-github:1.0.0'),
1504        or None if the metadata could not be fetched or parsed.
1505    """
1506    metadata_url = (
1507        f"https://raw.githubusercontent.com/airbytehq/airbyte/master/"
1508        f"airbyte-integrations/connectors/{connector_name}/metadata.yaml"
1509    )
1510    response = requests.get(metadata_url, timeout=30)
1511    if not response.ok:
1512        print_error(
1513            f"Failed to fetch metadata for {connector_name}: "
1514            f"HTTP {response.status_code} from {metadata_url}"
1515        )
1516        return None
1517
1518    metadata = yaml.safe_load(response.text)
1519    if not isinstance(metadata, dict):
1520        print_error(f"Invalid metadata format for {connector_name}: expected dict")
1521        return None
1522
1523    data = metadata.get("data", {})
1524    docker_repository = data.get("dockerRepository")
1525    docker_image_tag = data.get("dockerImageTag")
1526
1527    if not docker_repository or not docker_image_tag:
1528        print_error(
1529            f"Could not find dockerRepository/dockerImageTag in metadata for {connector_name}"
1530        )
1531        return None
1532
1533    return f"{docker_repository}:{docker_image_tag}"
1534
1535
1536def _run_with_optional_http_metrics(
1537    connector_image: str,
1538    command: Command,
1539    output_dir: Path,
1540    target_or_control: TargetOrControl,
1541    enable_http_metrics: bool,
1542    config_path: Path | None,
1543    catalog_path: Path | None,
1544    state_path: Path | None,
1545    enable_debug_logs: bool = False,
1546    replay_expected: bool = False,
1547    replay_from: Path | None = None,
1548    replay_options: ReplayOptions | None = None,
1549    stream_large_bodies: str = DEFAULT_STREAM_LARGE_BODIES,
1550) -> tuple[dict, ComparableOutputs]:
1551    """Run a connector command with optional HTTP metrics capture.
1552
1553    When enable_http_metrics is True, starts mitmproxy to capture HTTP traffic.
1554    If mitmproxy fails to start, falls back to running without metrics.
1555
1556    Args:
1557        connector_image: Full connector image name with tag.
1558        command: The Airbyte command to run.
1559        output_dir: Directory to store output files.
1560        target_or_control: Whether this is target or control version.
1561        enable_http_metrics: Whether to capture HTTP metrics via mitmproxy.
1562        config_path: Path to connector config JSON file.
1563        catalog_path: Path to configured catalog JSON file.
1564        state_path: Path to state JSON file.
1565        enable_debug_logs: If `True`, pass `LOG_LEVEL=DEBUG` to the connector container.
1566        replay_expected: Whether this run was *meant* to replay. Separate from
1567            `replay_from`, which is only non-None once a recording actually
1568            exists -- when the proxy is unavailable there is no recording to
1569            point at, and that is exactly the case that has to be reported.
1570        replay_from: A previous run's HTTP dump to serve responses from, so this
1571            run sees the same upstream data the recorded one did. Requires
1572            `enable_http_metrics`; ignored without it.
1573        replay_options: Replay tuning; defaults apply when omitted.
1574        stream_large_bodies: Body-size threshold above which mitmproxy streams
1575            rather than records a body.
1576
1577    Returns:
1578        What `_run_connector_command` returns, with `http_metrics` added to the
1579        results dict when they were captured -- plus `http_dump_file`, the dump
1580        this run recorded, which is what a subsequent run replays from.
1581    """
1582    if not enable_http_metrics:
1583        return _run_connector_command(
1584            connector_image=connector_image,
1585            command=command,
1586            output_dir=output_dir,
1587            target_or_control=target_or_control,
1588            config_path=config_path,
1589            catalog_path=catalog_path,
1590            state_path=state_path,
1591            enable_debug_logs=enable_debug_logs,
1592        )
1593
1594    # The manager rather than `MitmproxyManager.start`: the reason a startup
1595    # failed and the fact that a shutdown had to be killed are only reachable
1596    # through it -- there is no session to carry the first, and the second is
1597    # not known until the session is gone.
1598    manager = MitmproxyManager(
1599        output_dir,
1600        replay_from=replay_from,
1601        replay_options=replay_options,
1602        stream_large_bodies=stream_large_bodies,
1603    )
1604    with manager.running() as session:
1605        if session is None:
1606            startup_failure = (
1607                manager.startup_failure_reason or "mitmproxy is unavailable"
1608            )
1609            print_error(
1610                f"Mitmproxy unavailable ({startup_failure}), "
1611                "running without HTTP metrics"
1612            )
1613            result, comparable = _run_connector_command(
1614                connector_image=connector_image,
1615                command=command,
1616                output_dir=output_dir,
1617                target_or_control=target_or_control,
1618                config_path=config_path,
1619                catalog_path=catalog_path,
1620                state_path=state_path,
1621                enable_debug_logs=enable_debug_logs,
1622            )
1623            if replay_expected:
1624                # The one skip path that produces no session, and so reported
1625                # nothing at all: no metrics, no reason, an ordinary green
1626                # comparison against live data. Keyed on `replay_expected` and
1627                # not on `replay_from`: when mitmdump is missing the *control*
1628                # run records no dump either, so there is nothing to point at
1629                # and the condition that looks natural here is never true in
1630                # the failure it defends against.
1631                #
1632                # The precise reason, not a fixed "mitmproxy is unavailable":
1633                # `wait_for_proxy_listener` distinguishes a proxy that died on
1634                # startup and a corpus that took longer than its ceiling to
1635                # load, and those point at a code fix where the hardcoded
1636                # string tells a reader to install mitmproxy.
1637                result["http_metrics"] = {
1638                    "replay_skipped_reason": (
1639                        f"{startup_failure}, so nothing could be replayed"
1640                    )
1641                }
1642            return result, comparable
1643
1644        if session.replay_source:
1645            print_success(
1646                f"Started mitmproxy on {session.proxy_url} "
1647                f"(replaying {session.replay_source}, "
1648                f"{session.unreplayable_flow_count} flow(s) not replayable)"
1649            )
1650        else:
1651            print_success(f"Started mitmproxy on {session.proxy_url}")
1652            if session.replay_skipped_reason:
1653                print_warning(
1654                    f"HTTP replay unavailable ({session.replay_skipped_reason}); "
1655                    "this run hits the live API, so it cannot be compared against "
1656                    "identical upstream data"
1657                )
1658
1659        result, comparable = _run_connector_command(
1660            connector_image=connector_image,
1661            command=command,
1662            output_dir=output_dir,
1663            target_or_control=target_or_control,
1664            config_path=config_path,
1665            catalog_path=catalog_path,
1666            state_path=state_path,
1667            proxy_url=session.proxy_url,
1668            ca_cert_path=session.ca_cert_path,
1669            enable_debug_logs=enable_debug_logs,
1670        )
1671
1672        # The path, not the metrics: this is what the next run replays from.
1673        result["http_dump_file"] = str(session.dump_file_path)
1674
1675        # Everything the metrics are built from, carried out of the block. The
1676        # dump itself cannot be read here: mitmproxy's `save` addon writes each
1677        # flow through a block-buffered file and only flushes and closes it in
1678        # `done()`, holding in-flight flows until shutdown -- so parsing inside
1679        # the `with` misses the buffered tail, and `live_flow_count` is the
1680        # feature's acceptance criterion, not a rough count.
1681        dump_file_path = session.dump_file_path
1682        replay_source = session.replay_source
1683        replay_skipped_reason = session.replay_skipped_reason
1684        unreplayable_flow_count = session.unreplayable_flow_count
1685
1686    http_metrics = parse_http_dump(dump_file_path, replay_source)
1687    metrics_payload: dict[str, object] = {
1688        "flow_count": http_metrics.flow_count,
1689        "duplicate_flow_count": http_metrics.duplicate_flow_count,
1690        # Carried into the report because it is the number that tells a
1691        # reviewer whether a record difference is a real regression or
1692        # upstream drift. Today it is a duplicate-URL heuristic; it only
1693        # becomes exact once HTTP replay lands.
1694        "cache_hits_count": http_metrics.cache_hits_count,
1695        "cache_hit_ratio": http_metrics.cache_hit_ratio,
1696        # On a replaying target run, `live_flow_count` is the acceptance
1697        # criterion for replay -- so it belongs in the report rather than
1698        # only in the raw dump.
1699        "replayed_flow_count": http_metrics.replayed_flow_count,
1700        "live_flow_count": http_metrics.live_flow_count,
1701        "replay_hit_ratio": http_metrics.replay_hit_ratio,
1702    }
1703    if replay_from is not None:
1704        # Keyed on replay having been *asked* for, not on a corpus having
1705        # survived: when every recorded flow is unreplayable there is no
1706        # corpus, and that is exactly the run whose dropped count matters.
1707        # Poor replay coverage is usually this, and the fix is a bigger
1708        # `--http-replay-max-body`.
1709        metrics_payload["unreplayable_flow_count"] = unreplayable_flow_count
1710    if replay_source:
1711        metrics_payload["replay_source"] = str(replay_source)
1712        # Only on a replaying run, where a live request is a request the corpus
1713        # could not answer and so the one thing that explains a coverage
1714        # shortfall. On a recording run every request is live, so the same list
1715        # would be a request log, which is not a surface to widen for no
1716        # diagnostic gain. Query-parameter values are masked either way -- these
1717        # URLs reach the step summary, the workflow log and `report.html`, and
1718        # plenty of APIs put credentials in the query string; see `redact_url`.
1719        live_urls = http_metrics.top_live_urls()
1720        if live_urls:
1721            metrics_payload["live_urls"] = live_urls
1722            metrics_payload["live_unique_url_count"] = len(http_metrics.live_url_counts)
1723    if manager.metrics_incomplete_reason:
1724        # Read after the block, not inside it: a killed shutdown is only known
1725        # once `_stop` has run. Recorded because every count below is then a
1726        # lower bound -- and short in the direction that makes a run look closer
1727        # to the replay acceptance criterion than it was.
1728        metrics_payload["metrics_incomplete"] = manager.metrics_incomplete_reason
1729    skipped_reason = replay_skipped_reason
1730    if replay_expected and replay_source is None and not skipped_reason:
1731        # The proxy came up but there was no recording to replay -- the
1732        # control run captured nothing. Same outcome as any other skip, so
1733        # it has to read the same way rather than as a clean replayed run.
1734        skipped_reason = "the control run recorded no HTTP dump to replay from"
1735    if skipped_reason:
1736        metrics_payload["replay_skipped_reason"] = skipped_reason
1737    result["http_metrics"] = metrics_payload
1738
1739    print_success(
1740        f"Captured {http_metrics.flow_count} HTTP flows "
1741        f"({http_metrics.duplicate_flow_count} duplicates)"
1742    )
1743    if replay_source:
1744        print_success(
1745            f"Replayed {http_metrics.replayed_flow_count} of "
1746            f"{http_metrics.flow_count} request(s) "
1747            f"({http_metrics.replay_hit_ratio}); "
1748            f"{http_metrics.live_flow_count} went live"
1749        )
1750    return result, comparable
1751
1752
1753def _discard_oversized_http_dumps(*results: dict, max_bytes: int) -> None:
1754    """Drop recorded dumps too large to keep, and say so in the report.
1755
1756    "Too large to keep" is about the runner's disk, not the artifact upload:
1757    dumps are excluded from every upload step regardless of size.
1758
1759    Deleting is not silent: the size that was discarded lands in the run's
1760    `http_metrics`, so a reader looking for the dump learns why it is not there
1761    instead of finding a dangling path.
1762    """
1763    for result in results:
1764        dump_file = result.get("http_dump_file")
1765        if not dump_file:
1766            continue
1767
1768        discarded = discard_oversized_dump(Path(dump_file), max_bytes)
1769        if discarded is None:
1770            continue
1771
1772        result.pop("http_dump_file", None)
1773        metrics = result.get("http_metrics")
1774        if isinstance(metrics, dict):
1775            metrics["discarded_dump_mb"] = round(discarded / 1024 / 1024)
1776        print_warning(
1777            f"Discarded {dump_file}: {discarded / 1024 / 1024:.0f} MB exceeds the "
1778            "retention cap, and a dump this large is what fills the runner's disk. "
1779            "Dumps are never uploaded as artifacts; raise "
1780            "--http-replay-max-dump-size-mb to keep it for a local run"
1781        )
1782
1783
1784def _compare_read_records(
1785    target_result: dict,
1786    control_result: dict,
1787    catalog_file: Path | None,
1788) -> tuple[RecordComparisonSummary | None, str | None]:
1789    """Run record-level comparisons between control and target read outputs.
1790
1791    Uses primary keys from the configured catalog (source-defined primary
1792    keys come from the discovered catalog) to compare record counts, PK
1793    presence, PK integrity, and field values. Streams without a primary key
1794    fall back to a directional record count comparison.
1795
1796    Each run's stdout is first split into per-stream record files (a single
1797    streaming pass), and the comparison then loads one stream at a time, so
1798    memory stays bounded by the largest stream instead of the full output
1799    of both runs. The per-stream files land next to each run's stdout and
1800    double as debugging artifacts. Within one stream, both sides are
1801    materialized as parsed messages plus dict copies plus PK indexes --
1802    measured at roughly 13 bytes resident per byte of record JSON, so a
1803    single stream of ~500 MB of records is the practical ceiling on a
1804    standard 7 GB runner. Start here when debugging an OOM.
1805
1806    Returns:
1807        `(summary, skip_reason)`. The summary is None when the comparison
1808        cannot run at all -- then `skip_reason` says why, so the caller can
1809        surface a diagnostic check instead of letting the run read as a pass
1810        that compared records. An unexpected error during the comparison
1811        itself returns a FAILED summary instead of raising: by this point
1812        both (expensive, live-API) runs have completed, and a comparison bug
1813        must fail the verdict, not destroy the run's outputs and reports.
1814    """
1815    if catalog_file is None or not catalog_file.exists():
1816        reason = "no configured catalog available"
1817        print_warning(f"Skipping record-level comparison: {reason}")
1818        return None, reason
1819
1820    try:
1821        configured_catalog = ConfiguredAirbyteCatalog.model_validate_json(
1822            catalog_file.read_text()
1823        )
1824    except Exception as exc:
1825        reason = f"configured catalog could not be parsed ({exc})"
1826        print_warning(f"Skipping record-level comparison: {reason}")
1827        return None, reason
1828
1829    try:
1830        primary_keys = get_primary_keys_per_stream(configured_catalog)
1831        # A stream name that repeats across namespaces merges two streams in
1832        # this name-keyed pipeline; those names degraded to count-only above,
1833        # and the reader must hear why from the summary, not just stdout.
1834        collision_warnings = [
1835            f"Stream name {name!r} appears in more than one namespace; the "
1836            "record comparison keys streams by name, so its records are "
1837            "compared by count only"
1838            for name in duplicate_stream_names(configured_catalog)
1839        ]
1840        control_record_files = split_records_per_stream(
1841            Path(control_result["stdout_file"])
1842        )
1843        target_record_files = split_records_per_stream(
1844            Path(target_result["stdout_file"])
1845        )
1846
1847        summary = run_record_comparisons_from_files(
1848            control_record_files=control_record_files,
1849            target_record_files=target_record_files,
1850            primary_keys_per_stream=primary_keys,
1851        )
1852        summary.warnings[:0] = collision_warnings
1853    except Exception as exc:
1854        summary = RecordComparisonSummary(
1855            passed=False,
1856            errored=True,
1857            count_comparison=ComparisonResult(
1858                passed=False, message="Record comparison errored"
1859            ),
1860            errors=[f"Record comparison errored: {exc!r}"],
1861        )
1862
1863    if summary.passed:
1864        print_success("Record comparison passed (counts, PKs, field values)")
1865    else:
1866        print_error("Record comparison failed:")
1867        for error in summary.errors[:20]:
1868            print_error(f"  - {error}")
1869    for warning in summary.warnings[:20]:
1870        print_warning(warning)
1871
1872    return summary, None
1873
1874
1875@connector_app.command(name="regression-test")
1876def regression_test(
1877    skip_compare: Annotated[
1878        bool,
1879        Parameter(
1880            help="If True, skip comparison and run single-version tests only. "
1881            "If False (default), run comparison tests (target vs control)."
1882        ),
1883    ] = False,
1884    skip_record_comparison: Annotated[
1885        bool,
1886        Parameter(
1887            help="If True, skip the record-level comparison for read commands "
1888            "(counts, PK presence, PK integrity, field values) and gate the "
1889            "verdict on command outcomes only. Escape hatch for sources whose "
1890            "data legitimately changes between the two live runs (rolling "
1891            "windows, feeds that age out records) until HTTP request caching "
1892            "makes strict comparison reliable."
1893        ),
1894    ] = False,
1895    test_image: Annotated[
1896        str | None,
1897        Parameter(
1898            help="Test connector image with tag (e.g., airbyte/source-github:1.0.0). "
1899            "This is the image under test - in comparison mode, it's compared against control_image."
1900        ),
1901    ] = None,
1902    control_image: Annotated[
1903        str | None,
1904        Parameter(
1905            help="Control connector image (baseline version) with tag (e.g., airbyte/source-github:1.0.0). "
1906            "Ignored if `skip_compare=True`."
1907        ),
1908    ] = None,
1909    connector_name: Annotated[
1910        str | None,
1911        Parameter(
1912            help="Connector name to build image from source (e.g., 'source-pokeapi'). "
1913            "If provided, builds the image locally with tag 'dev'. "
1914            "For comparison tests (default), this builds the target image. "
1915            "For single-version tests (skip_compare=True), this builds the test image."
1916        ),
1917    ] = None,
1918    repo_root: Annotated[
1919        str | None,
1920        Parameter(
1921            help="Path to the airbyte repo root. Required if connector_name is provided "
1922            "and the repo cannot be auto-detected."
1923        ),
1924    ] = None,
1925    command: Annotated[
1926        Literal["spec", "check", "discover", "read"],
1927        Parameter(help="The Airbyte command to run."),
1928    ] = "check",
1929    connection_id: Annotated[
1930        str | None,
1931        Parameter(
1932            help="Airbyte Cloud connection ID to fetch config/catalog from. "
1933            "Mutually exclusive with config-path/catalog-path. "
1934            "If provided, test_image/control_image can be auto-detected."
1935        ),
1936    ] = None,
1937    config_path: Annotated[
1938        str | None,
1939        Parameter(help="Path to the connector config JSON file."),
1940    ] = None,
1941    catalog_path: Annotated[
1942        str | None,
1943        Parameter(help="Path to the configured catalog JSON file (required for read)."),
1944    ] = None,
1945    state_path: Annotated[
1946        str | None,
1947        Parameter(help="Path to the state JSON file (optional for read)."),
1948    ] = None,
1949    output_dir: Annotated[
1950        str,
1951        Parameter(help="Directory to store test artifacts."),
1952    ] = "/tmp/regression_test_artifacts",
1953    enable_http_metrics: Annotated[
1954        bool,
1955        Parameter(
1956            help="Capture HTTP traffic metrics via mitmproxy. Requires mitmdump to "
1957            "be installed; without it the run proceeds without metrics. Only used "
1958            "in comparison mode, where HTTP replay turns it on anyway -- set this "
1959            "to record traffic while running both versions live "
1960            "(--no-http-replay)."
1961        ),
1962    ] = False,
1963    enable_http_replay: Annotated[
1964        bool | None,
1965        Parameter(
1966            negative="--no-http-replay",
1967            help="Replay the control run's recorded HTTP responses for the target "
1968            "run, so both versions see identical upstream data. Implies "
1969            "--enable-http-metrics. Defaults to `True` in comparison mode; has no "
1970            "effect with `--skip-compare`, which has no control run to replay "
1971            "from. Pass --no-http-replay to run both versions against the live "
1972            "API, which is the fallback whenever replay cannot be delivered "
1973            "anyway.",
1974        ),
1975    ] = None,
1976    http_replay_max_body: Annotated[
1977        str,
1978        Parameter(
1979            help="Response bodies larger than this are streamed rather than "
1980            "recorded, and therefore re-fetched live instead of replayed. Raise it "
1981            "for a connector with large pages, at the cost of replay memory. "
1982            "Accepts a mitmproxy size such as '1m' or '10m' -- never a bare number, "
1983            "which means bytes. Unlike the other --http-replay-* options this one "
1984            "governs recording, so it applies with --enable-http-metrics alone."
1985        ),
1986    ] = DEFAULT_STREAM_LARGE_BODIES,
1987    http_replay_max_dump_size_mb: Annotated[
1988        int,
1989        Parameter(
1990            help="Skip replay and run the target live when the filtered replay "
1991            "corpus is larger than this, and discard a recorded dump this large "
1992            "once the run is done with it. mitmproxy holds the whole corpus in "
1993            "memory for the entire run, and a dump grows with everything the "
1994            "connector downloads, so this is what keeps a big connector from "
1995            "exhausting the CI runner's memory and disk. Dumps are never "
1996            "uploaded as artifacts, so on a local run -- where the dump is the "
1997            "only thing that can diagnose a coverage shortfall offline -- raise "
1998            "it to keep the dump."
1999        ),
2000    ] = DEFAULT_MAX_REPLAY_DUMP_SIZE_MB,
2001    http_replay_reuse: Annotated[
2002        bool,
2003        Parameter(
2004            help="Serve a recorded response more than once. By default each "
2005            "recorded flow is consumed on first match, which preserves recorded "
2006            "ordering and lets memory fall as the run proceeds; reuse keeps the "
2007            "whole corpus resident for the entire run."
2008        ),
2009    ] = False,
2010    http_replay_ignore_params: Annotated[
2011        str | None,
2012        Parameter(
2013            help="Comma-separated query parameters to exclude from the replay match "
2014            "key, for connectors that put a timestamp or nonce in the URL."
2015        ),
2016    ] = None,
2017    strict_replay: Annotated[
2018        bool,
2019        Parameter(
2020            help="Kill requests that match nothing in the replay corpus instead of "
2021            "forwarding them upstream. Diagnostic only: it turns unmatched requests "
2022            "into connector errors, which is how you count them. mitmproxy stops "
2023            "intervening once its flowmap is empty, so combine this with "
2024            "--http-replay-reuse to cover the whole run rather than only the part "
2025            "before the corpus is consumed."
2026        ),
2027    ] = False,
2028    selected_streams: Annotated[
2029        str | None,
2030        Parameter(
2031            help="Comma-separated list of stream names to include in the read. "
2032            "Only these streams will be included in the configured catalog. "
2033            "This is useful to limit data volume by testing only specific streams."
2034        ),
2035    ] = None,
2036    enable_debug_logs: Annotated[
2037        bool,
2038        Parameter(
2039            help="Enable debug-level logging for regression test output. "
2040            "Also passed as `LOG_LEVEL=DEBUG` to the connector Docker container."
2041        ),
2042    ] = False,
2043    with_state: Annotated[
2044        bool | None,
2045        Parameter(
2046            negative="--no-state",
2047            help="Fetch and pass the connection's current state to the read command, "
2048            "producing a warm read instead of a cold read. Defaults to `True` when "
2049            "`--connection-id` is provided, `False` otherwise. Has no effect unless "
2050            "the command is `read`. Ignored when `--state-path` is explicitly provided.",
2051        ),
2052    ] = None,
2053) -> None:
2054    """Run regression tests on connectors.
2055
2056    This command supports two modes:
2057
2058    Comparison mode (skip_compare=False, default):
2059        Runs the specified Airbyte protocol command against both the target (new)
2060        and control (baseline) connector versions, then compares the results.
2061        This helps identify regressions between versions.
2062
2063    Single-version mode (skip_compare=True):
2064        Runs the specified Airbyte protocol command against a single connector
2065        and validates the output. No comparison is performed.
2066
2067    Results are written to the output directory and to GitHub Actions outputs
2068    if running in CI.
2069
2070    You can provide the test image in three ways:
2071    1. --test-image: Use a pre-built image from Docker registry
2072    2. --connector-name: Build the image locally from source code
2073    3. --connection-id: Auto-detect from an Airbyte Cloud connection
2074
2075    You can provide config/catalog either via file paths OR via a connection_id
2076    that fetches them from Airbyte Cloud.
2077    """
2078    # Configure debug logging for the regression test harness when requested
2079    if enable_debug_logs:
2080        logging.basicConfig(
2081            level=logging.DEBUG,
2082            format="%(asctime)s - %(name)s - %(levelname)s - %(message)s",
2083            force=True,
2084        )
2085
2086    output_path = Path(output_dir)
2087    output_path.mkdir(parents=True, exist_ok=True)
2088
2089    cmd = Command(command)
2090
2091    # On by default, but only where it can work: single-version mode has no
2092    # control run to replay from. Asking for it there explicitly is a mistake
2093    # worth naming; the default landing there is not, so it is not a warning.
2094    if skip_compare:
2095        if enable_http_replay:
2096            print_warning("--enable-http-replay needs a control run; ignoring it")
2097        replay_enabled = False
2098    else:
2099        replay_enabled = True if enable_http_replay is None else enable_http_replay
2100
2101    # Replay is a use of the capture, not a separate mechanism, so it turns the
2102    # capture on rather than erroring.
2103    if replay_enabled:
2104        enable_http_metrics = True
2105    elif (
2106        http_replay_reuse
2107        or strict_replay
2108        or http_replay_ignore_params
2109        or http_replay_max_dump_size_mb != DEFAULT_MAX_REPLAY_DUMP_SIZE_MB
2110    ):
2111        # Otherwise someone debugging poor coverage with --strict-replay alone
2112        # concludes strict mode found nothing wrong. `--http-replay-max-body` is
2113        # deliberately not in this list: it governs recording and applies anyway.
2114        print_warning(
2115            "Replay tuning flags do nothing with --no-http-replay; ignoring them"
2116        )
2117
2118    replay_options = ReplayOptions(
2119        reuse=http_replay_reuse,
2120        extra="kill" if strict_replay else DEFAULT_REPLAY_EXTRA,
2121        ignore_params=tuple(
2122            param.strip()
2123            for param in (http_replay_ignore_params or "").split(",")
2124            if param.strip()
2125        ),
2126        max_dump_bytes=http_replay_max_dump_size_mb * 1024 * 1024,
2127    )
2128
2129    config_file: Path | None = None
2130    catalog_file: Path | None = None
2131    state_file = Path(state_path) if state_path else None
2132
2133    selected_stream_names: set[str] = (
2134        {s.strip() for s in selected_streams.split(",") if s.strip()}
2135        if selected_streams
2136        else set()
2137    )
2138
2139    # Resolve the test image (used in both single-version and comparison modes)
2140    resolved_test_image: str | None = test_image
2141    resolved_control_image: str | None = control_image
2142
2143    # Validate conflicting parameters
2144    # Single-version mode: reject comparison-specific parameters
2145    if skip_compare and control_image:
2146        write_github_output("success", False)
2147        write_github_output(
2148            "error", "Cannot specify control_image with skip_compare=True"
2149        )
2150        exit_with_error(
2151            "Cannot specify --control-image with --skip-compare. "
2152            "Control image is only used in comparison mode."
2153        )
2154
2155    # If connector_name is provided, build the image from source
2156    if connector_name:
2157        if resolved_test_image:
2158            write_github_output("success", False)
2159            write_github_output(
2160                "error", "Cannot specify both test_image and connector_name"
2161            )
2162            exit_with_error("Cannot specify both --test-image and --connector-name")
2163
2164        repo_root_path = Path(repo_root) if repo_root else None
2165        built_image = _build_connector_image_from_source(
2166            connector_name=connector_name,
2167            repo_root=repo_root_path,
2168            tag="dev",
2169        )
2170        if not built_image:
2171            write_github_output("success", False)
2172            write_github_output("error", f"Failed to build image for {connector_name}")
2173            exit_with_error(f"Failed to build image for {connector_name}")
2174        resolved_test_image = built_image
2175
2176    if connection_id:
2177        if config_path or catalog_path:
2178            write_github_output("success", False)
2179            write_github_output(
2180                "error", "Cannot specify both connection_id and file paths"
2181            )
2182            exit_with_error(
2183                "Cannot specify both connection_id and config_path/catalog_path"
2184            )
2185
2186        print_success(f"Fetching config/catalog from connection: {connection_id}")
2187        connection_data = fetch_connection_data(connection_id)
2188
2189        # Check if we should retrieve unmasked secrets
2190        if should_use_secret_retriever():
2191            print_success(
2192                "USE_CONNECTION_SECRET_RETRIEVER enabled - enriching config with unmasked secrets..."
2193            )
2194            try:
2195                connection_data = enrich_config_with_secrets(
2196                    connection_data,
2197                    retrieval_reason="Regression test with USE_CONNECTION_SECRET_RETRIEVER=true",
2198                )
2199                print_success("Successfully retrieved unmasked secrets from database")
2200            except SecretRetrievalError as e:
2201                write_github_output("success", False)
2202                write_github_output("error", str(e))
2203                exit_with_error(
2204                    f"{e}\n\n"
2205                    f"This connection cannot be used for regression testing. "
2206                    f"Please use a connection from a non-EU workspace, or use GSM-based "
2207                    f"integration test credentials instead (by omitting --connection-id)."
2208                )
2209
2210        # Fill in stream schemas: the minimal catalog from the public API has
2211        # empty json_schema values, which silently yields zero records for
2212        # connectors that derive request fields from the catalog schema.
2213        # Only the read command consumes the catalog, so spec/check/discover
2214        # skip the extra Cloud API round trip.
2215        enrichment_state: list[dict[str, Any]] | None = None
2216        enrichment_state_fetched = False
2217        if cmd.needs_catalog():
2218            merged_count = 0
2219            try:
2220                enrichment = enrich_catalog_schemas_from_cloud(connection_data)
2221                connection_data = enrichment.connection_data
2222                merged_count = enrichment.merged_count
2223                enrichment_state = enrichment.state
2224                enrichment_state_fetched = True
2225            except Exception as exc:
2226                print_error(f"Failed to fetch stream schemas from Cloud API: {exc}")
2227            missing_schemas = streams_missing_schemas(connection_data.catalog)
2228            if selected_stream_names:
2229                missing_schemas = [
2230                    s for s in missing_schemas if s in selected_stream_names
2231                ]
2232            if missing_schemas and should_use_secret_retriever():
2233                try:
2234                    connection_data, db_merged_count = (
2235                        enrich_catalog_schemas_from_platform_db(
2236                            connection_data,
2237                            retrieval_reason="Regression test catalog schema enrichment",
2238                        )
2239                    )
2240                    merged_count += db_merged_count
2241                except Exception as exc:
2242                    print_error(f"Failed to merge stored stream schemas: {exc}")
2243                missing_schemas = streams_missing_schemas(connection_data.catalog)
2244                if selected_stream_names:
2245                    missing_schemas = [
2246                        s for s in missing_schemas if s in selected_stream_names
2247                    ]
2248            total_streams = len(connection_data.catalog.get("streams", []))
2249            checked_streams = (
2250                f"{len(missing_schemas)} of {len(selected_stream_names)} "
2251                "selected stream(s)"
2252                if selected_stream_names
2253                else f"{len(missing_schemas)} of {total_streams} stream(s)"
2254            )
2255            if not total_streams:
2256                error_msg = (
2257                    "No configured streams found in the connection catalog. "
2258                    "A read with zero streams would produce zero records on "
2259                    "both versions, so failing instead of producing a false pass."
2260                )
2261                write_github_output("success", False)
2262                write_github_output("error", error_msg)
2263                exit_with_error(error_msg)
2264            if missing_schemas:
2265                error_msg = (
2266                    f"Stream schemas missing for {checked_streams}: "
2267                    f"{', '.join(missing_schemas)}. "
2268                    "Connectors that derive request fields from the catalog "
2269                    "schema would silently return zero records, so failing "
2270                    "instead of producing a false pass."
2271                )
2272                write_github_output("success", False)
2273                write_github_output("error", error_msg)
2274                exit_with_error(error_msg)
2275            print_success(
2276                f"Stream schemas ready: merged {merged_count} of {total_streams} "
2277                "stream(s)"
2278                + (
2279                    f"; checked {len(selected_stream_names)} selected stream(s)"
2280                    if selected_stream_names
2281                    else ""
2282                )
2283            )
2284
2285        # Resolve with_state default: True when connection_id is provided
2286        resolved_with_state = with_state if with_state is not None else True
2287
2288        if not (resolved_with_state and not state_file and command == "read"):
2289            connection_data.state = None
2290        elif enrichment_state_fetched:
2291            connection_data.state = enrichment_state
2292            if connection_data.state:
2293                print_success(
2294                    f"Using state fetched during schema enrichment "
2295                    f"({len(connection_data.state)} stream(s))"
2296                )
2297            else:
2298                print_success("No state available for this connection (initial sync)")
2299        else:
2300            print_success("Fetching connection state for warm read...")
2301            try:
2302                connection_data.state = fetch_connection_state(connection_data)
2303                if connection_data.state:
2304                    print_success(
2305                        f"Fetched state with {len(connection_data.state)} stream(s)"
2306                    )
2307                else:
2308                    connection_data.state = None
2309                    print_success(
2310                        "No state available for this connection (initial sync)"
2311                    )
2312            except Exception as exc:
2313                print_error(f"Failed to fetch connection state: {exc}")
2314                print_success("Falling back to cold read (no state)")
2315
2316        config_file, catalog_file, state_file_from_conn = save_connection_data_to_files(
2317            connection_data, output_path / "connection_data"
2318        )
2319
2320        if state_file_from_conn and not state_file:
2321            state_file = state_file_from_conn
2322
2323        print_success(
2324            f"Fetched config for source: {connection_data.source_name} "
2325            f"with {len(connection_data.stream_names)} streams"
2326        )
2327
2328        # Auto-detect test/control image from connection if not provided
2329        if not resolved_test_image and connection_data.connector_image:
2330            resolved_test_image = connection_data.connector_image
2331            print_success(f"Auto-detected test image: {resolved_test_image}")
2332
2333        if (
2334            not skip_compare
2335            and not resolved_control_image
2336            and connection_data.connector_image
2337        ):
2338            resolved_control_image = connection_data.connector_image
2339            print_success(f"Auto-detected control image: {resolved_control_image}")
2340    elif config_path:
2341        config_file = Path(config_path)
2342        catalog_file = Path(catalog_path) if catalog_path else None
2343    elif connector_name:
2344        # Fallback: fetch integration test secrets from GSM using PyAirbyte API
2345        print_success(
2346            f"No connection_id or config_path provided. "
2347            f"Attempting to fetch integration test config from GSM for {connector_name}..."
2348        )
2349        gsm_config = get_first_config_from_secrets(connector_name)
2350        if gsm_config:
2351            # Write config to a temp file (not in output_path to avoid artifact upload)
2352            gsm_config_dir = Path(
2353                tempfile.mkdtemp(prefix=f"gsm-config-{connector_name}-")
2354            )
2355            gsm_config_dir.chmod(0o700)
2356            gsm_config_file = gsm_config_dir / "config.json"
2357            gsm_config_file.write_text(json.dumps(gsm_config, indent=2))
2358            gsm_config_file.chmod(0o600)
2359            config_file = gsm_config_file
2360            # Use catalog_path if provided (e.g., generated from discover output)
2361            catalog_file = Path(catalog_path) if catalog_path else None
2362            print_success(
2363                f"Fetched integration test config from GSM for {connector_name}"
2364            )
2365        else:
2366            print_error(
2367                f"Failed to fetch integration test config from GSM for {connector_name}."
2368            )
2369            config_file = None
2370            # Use catalog_path if provided (e.g., generated from discover output)
2371            catalog_file = Path(catalog_path) if catalog_path else None
2372    else:
2373        config_file = None
2374        catalog_file = Path(catalog_path) if catalog_path else None
2375
2376    # Auto-detect control_image from metadata.yaml if connector_name is provided (comparison mode only)
2377    if not skip_compare and not resolved_control_image and connector_name:
2378        resolved_control_image = _fetch_control_image_from_metadata(connector_name)
2379        if resolved_control_image:
2380            print_success(
2381                f"Auto-detected control image from metadata.yaml: {resolved_control_image}"
2382            )
2383
2384    # Validate that we have the required images
2385    if not resolved_test_image:
2386        write_github_output("success", False)
2387        write_github_output("error", "No test image specified")
2388        exit_with_error(
2389            "You must provide one of the following: a test_image, a connector_name "
2390            "to build the image from source, or a connection_id to auto-detect the image."
2391        )
2392
2393    if not skip_compare and not resolved_control_image:
2394        write_github_output("success", False)
2395        write_github_output("error", "No control image specified")
2396        exit_with_error(
2397            "You must provide one of the following: a control_image, a connection_id "
2398            "for a connection that has an associated connector image, or a connector_name "
2399            "to auto-detect the control image from the airbyte repo's metadata.yaml."
2400        )
2401
2402    # Pull images if they weren't just built locally
2403    # If connector_name was provided, we just built the test image locally
2404    if not connector_name and not ensure_image_available(resolved_test_image):
2405        write_github_output("success", False)
2406        write_github_output("error", f"Failed to pull image: {resolved_test_image}")
2407        exit_with_error(f"Failed to pull test image: {resolved_test_image}")
2408
2409    if (
2410        not skip_compare
2411        and resolved_control_image
2412        and not ensure_image_available(resolved_control_image)
2413    ):
2414        write_github_output("success", False)
2415        write_github_output("error", f"Failed to pull image: {resolved_control_image}")
2416        exit_with_error(
2417            f"Failed to pull control connector image: {resolved_control_image}"
2418        )
2419
2420    # Apply selected_streams filter to catalog if requested
2421    if selected_stream_names and catalog_file:
2422        print_success(
2423            f"Filtering catalog to {len(selected_stream_names)} selected streams: "
2424            f"{', '.join(sorted(selected_stream_names))}"
2425        )
2426        filter_configured_catalog_file(catalog_file, selected_stream_names)
2427
2428    # Upgrade READ → READ_WITH_STATE when a state file is available
2429    if cmd == Command.READ and state_file:
2430        cmd = Command.READ_WITH_STATE
2431        print_success("State available — using read-with-state (warm read)")
2432
2433    # Track telemetry for the regression test
2434    # Extract version from image tag (e.g., "airbyte/source-github:1.0.0" -> "1.0.0")
2435    target_version = (
2436        resolved_test_image.rsplit(":", 1)[-1]
2437        if ":" in resolved_test_image
2438        else "unknown"
2439    )
2440    control_version = None
2441    if resolved_control_image and ":" in resolved_control_image:
2442        control_version = resolved_control_image.rsplit(":", 1)[-1]
2443
2444    # Get tester identity from environment (GitHub Actions sets GITHUB_ACTOR)
2445    tester = os.getenv("GITHUB_ACTOR") or os.getenv("USER")
2446
2447    track_regression_test(
2448        user_id=tester,
2449        connector_image=resolved_test_image,
2450        command=command,
2451        target_version=target_version,
2452        control_version=control_version,
2453        additional_properties={
2454            "connection_id": connection_id,
2455            "skip_compare": skip_compare,
2456            "with_state": cmd == Command.READ_WITH_STATE,
2457        },
2458    )
2459
2460    # Execute the appropriate mode
2461    if skip_compare:
2462        # Single-version mode: run only the connector image. There is no control
2463        # to compare its spec or catalog against, so the declared objects go
2464        # unread here.
2465        result, _ = _run_connector_command(
2466            connector_image=resolved_test_image,
2467            command=cmd,
2468            output_dir=output_path,
2469            target_or_control=TargetOrControl.TARGET,
2470            config_path=config_file,
2471            catalog_path=catalog_file,
2472            state_path=state_file,
2473            enable_debug_logs=enable_debug_logs,
2474        )
2475
2476        # Same per-version rule as comparison mode: for CHECK the exit code alone
2477        # is not the verdict, connectionStatus is.
2478        annotated_result = _annotate_run_verdict(command, result)
2479        single_success = bool(annotated_result["success"])
2480
2481        print_json(annotated_result)
2482
2483        single_outputs: dict[str, object] = {
2484            "success": single_success,
2485            "connector": resolved_test_image,
2486            "command": command,
2487            "exit_code": result["exit_code"],
2488        }
2489        if command == "check":
2490            # Named to match the comparison-mode output the workflow summary reads.
2491            single_outputs["target_connection_status"] = result.get("connection_status")
2492
2493        write_github_outputs(single_outputs)
2494
2495        # Generate report.md with detailed metrics
2496        report_path = generate_single_version_report(
2497            connector_image=resolved_test_image,
2498            command=command,
2499            result=result,
2500            output_dir=output_path,
2501        )
2502        print_success(f"Generated report: {report_path}")
2503
2504        # The detailed surface is report.html in the uploaded artifact; the step
2505        # summary gets the concise block instead of the whole report. Both are
2506        # built from the unannotated result: the model derives the verdict itself.
2507        report_model = build_single_version_report_model(
2508            connector_image=resolved_test_image,
2509            command=command,
2510            result=result,
2511            configured_streams=_configured_stream_names(cmd, catalog_file),
2512            connection_objects=_collect_connection_objects(
2513                catalog_file, state_file, output_path / "airbyte_messages"
2514            ),
2515        )
2516        html_report_path = write_html_report(report_model, output_path)
2517        print_success(f"Generated HTML report: {html_report_path}")
2518
2519        # Write the concise summary to GITHUB_STEP_SUMMARY (if env var exists)
2520        write_github_summary(render_inline_summary(report_model))
2521
2522        if single_success:
2523            print_success(
2524                f"Single-version regression test passed for {resolved_test_image}"
2525            )
2526        elif command == "check" and result["success"]:
2527            print_error(
2528                f"Single-version regression test failed for {resolved_test_image}: "
2529                f"connectionStatus {result.get('connection_status') or 'missing'} "
2530                "(the connector itself exited 0)."
2531            )
2532        else:
2533            print_error(
2534                f"Single-version regression test failed for {resolved_test_image}"
2535            )
2536    else:
2537        # Comparison mode: run both control and target images. The control
2538        # version runs FIRST on purpose: the record comparison is
2539        # directional (extra records in the target are allowed, missing
2540        # ones fail), so data created upstream between the two live runs
2541        # must land in the target output — where it is tolerated — rather
2542        # than in the control output, where it would read as records
2543        # missing from the target and fail the comparison. Accepted side
2544        # effect (applies to every command, not just read): a hanging or
2545        # timing-out control run means the target never runs, so such a
2546        # run says nothing about the version under test.
2547        #
2548        # HTTP replay needs the same order for its own reason: the target
2549        # replays the dump the control run recorded, so that dump has to
2550        # exist before the target starts.
2551        target_output = output_path / "target"
2552        control_output = output_path / "control"
2553
2554        control_result, control_outputs = _run_with_optional_http_metrics(
2555            connector_image=resolved_control_image,  # type: ignore[arg-type]
2556            command=cmd,
2557            output_dir=control_output,
2558            target_or_control=TargetOrControl.CONTROL,
2559            enable_http_metrics=enable_http_metrics,
2560            config_path=config_file,
2561            catalog_path=catalog_file,
2562            state_path=state_file,
2563            enable_debug_logs=enable_debug_logs,
2564            stream_large_bodies=http_replay_max_body,
2565        )
2566
2567        replay_from: Path | None = None
2568        if replay_enabled:
2569            control_dump = control_result.get("http_dump_file")
2570            if control_dump:
2571                replay_from = Path(control_dump)
2572            else:
2573                print_warning(
2574                    "Control run recorded no HTTP dump; the target will run live"
2575                )
2576
2577        target_result, target_outputs = _run_with_optional_http_metrics(
2578            connector_image=resolved_test_image,
2579            command=cmd,
2580            output_dir=target_output,
2581            target_or_control=TargetOrControl.TARGET,
2582            enable_http_metrics=enable_http_metrics,
2583            config_path=config_file,
2584            catalog_path=catalog_file,
2585            state_path=state_file,
2586            enable_debug_logs=enable_debug_logs,
2587            replay_expected=replay_enabled,
2588            replay_from=replay_from,
2589            replay_options=replay_options,
2590            stream_large_bodies=http_replay_max_body,
2591        )
2592
2593        # Both dumps have served their purpose by now -- metrics parsed, replay
2594        # corpus built -- so anything too large to be a usable artifact goes,
2595        # rather than being uploaded from a runner with ~14 GB of disk.
2596        _discard_oversized_http_dumps(
2597            control_result,
2598            target_result,
2599            max_bytes=http_replay_max_dump_size_mb * 1024 * 1024,
2600        )
2601
2602        outcome = evaluate_comparison_outcome(command, target_result, control_result)
2603
2604        # Record-level comparison for READ commands (counts, PK presence,
2605        # PK integrity, field values). Only meaningful when both runs
2606        # succeeded; a comparison failure is a regression and fails the
2607        # verdict even though both containers exited cleanly. `outcome`
2608        # itself stays the pure exit-status/connectionStatus layer — the
2609        # folded verdict lives in `success`/`regression_detected` below.
2610        record_comparison: RecordComparisonSummary | None = None
2611        # A skipped comparison must be visible on the report surfaces, not
2612        # only on stdout: without this row the run reads as a pass that
2613        # compared records. Diagnostic, because the skip is not a regression.
2614        record_comparison_skipped: CheckResult | None = None
2615        if (
2616            cmd in (Command.READ, Command.READ_WITH_STATE)
2617            and outcome.both_succeeded
2618            and not skip_record_comparison
2619        ):
2620            record_comparison, skip_reason = _compare_read_records(
2621                target_result=target_result,
2622                control_result=control_result,
2623                catalog_file=catalog_file,
2624            )
2625            if record_comparison is None and skip_reason is not None:
2626                record_comparison_skipped = CheckResult(
2627                    name="Record comparison",
2628                    passed=False,
2629                    severity="diagnostic",
2630                    summary=f"Skipped — {skip_reason}; records were not compared",
2631                )
2632        elif skip_record_comparison and cmd in (Command.READ, Command.READ_WITH_STATE):
2633            print_warning(
2634                "Record-level comparison skipped (--skip-record-comparison); "
2635                "the verdict gates on command outcomes only."
2636            )
2637            record_comparison_skipped = CheckResult(
2638                name="Record comparison",
2639                passed=False,
2640                severity="diagnostic",
2641                summary="Skipped (--skip-record-comparison); records were not compared",
2642            )
2643        record_comparison_failed = (
2644            record_comparison is not None and not record_comparison.passed
2645        )
2646        # An errored comparison fails the run (fail closed) but as
2647        # INCONCLUSIVE, like both_failed: a comparator crash must not tell a
2648        # connector developer their change caused a regression.
2649        comparison_errored = record_comparison is not None and record_comparison.errored
2650
2651        # Spec/catalog-schema comparison of what the two versions declared
2652        # about themselves; a failed check here is a regression as much as
2653        # a dropped record.
2654        (
2655            comparison_checks,
2656            comparison_diff_blocks,
2657            comparison_detail,
2658        ) = _compare_outputs(
2659            command=command,
2660            outcome=outcome,
2661            control_outputs=control_outputs,
2662            target_outputs=target_outputs,
2663        )
2664        failed_comparisons = failed_checks(comparison_checks)
2665
2666        success = (
2667            outcome.success and not record_comparison_failed and not failed_comparisons
2668        )
2669        regression_detected = (
2670            outcome.regression_detected
2671            or (record_comparison_failed and not comparison_errored)
2672            or bool(failed_comparisons)
2673        )
2674
2675        # Annotate copies of each side, never the originals: `generate_regression_report`
2676        # below reads `result["success"]` as the exit status, so handing it a dict whose
2677        # `success` is the verdict would corrupt the report's per-version wording.
2678        combined_result = {
2679            "target": _annotate_run_verdict(command, target_result),
2680            "control": _annotate_run_verdict(command, control_result),
2681            "both_succeeded": outcome.both_succeeded,
2682            "regression_detected": regression_detected,
2683            # The authoritative verdict, so consumers of this payload (stdout and
2684            # the `regression_report` output) never have to re-derive it from
2685            # `regression_detected` -- which is false for an inconclusive run.
2686            "success": success,
2687            "both_failed": outcome.both_failed,
2688        }
2689        if record_comparison is not None:
2690            combined_result["record_comparison"] = record_comparison.to_dict()
2691        if comparison_detail:
2692            combined_result["comparison"] = comparison_detail
2693
2694        print_json(combined_result)
2695
2696        gh_outputs: dict[str, object] = {
2697            "success": success,
2698            "target_image": resolved_test_image,
2699            "control_image": resolved_control_image,
2700            "command": command,
2701            "target_exit_code": target_result["exit_code"],
2702            "control_exit_code": control_result["exit_code"],
2703            "regression_detected": regression_detected,
2704            "both_failed": outcome.both_failed,
2705        }
2706
2707        if command == "check":
2708            gh_outputs["target_connection_status"] = target_result.get(
2709                "connection_status"
2710            )
2711            gh_outputs["control_connection_status"] = control_result.get(
2712                "connection_status"
2713            )
2714
2715        if record_comparison is not None:
2716            gh_outputs["record_comparison_passed"] = record_comparison.passed
2717
2718        write_github_outputs(gh_outputs)
2719
2720        write_json_output("regression_report", combined_result)
2721
2722        report_path = generate_regression_report(
2723            target_image=resolved_test_image,
2724            control_image=resolved_control_image,  # type: ignore[arg-type]
2725            command=command,
2726            target_result=target_result,
2727            control_result=control_result,
2728            output_dir=output_path,
2729            record_comparison=(
2730                record_comparison.to_dict() if record_comparison is not None else None
2731            ),
2732            # Without these the markdown report would announce "no regression"
2733            # over the same run the HTML report and the CLI call a failure --
2734            # and, for a warning, would show no finding at all where the other
2735            # two surfaces show one.
2736            comparison_findings=[
2737                (check.status, f"{check.name}: {check.summary}")
2738                for check in comparison_checks
2739                if check.status in {"fail", "warn"}
2740            ],
2741        )
2742        print_success(f"Generated regression report: {report_path}")
2743
2744        # Same rule as `generate_regression_report` above: the model reads
2745        # `result["success"]` as the exit status and derives the per-version
2746        # verdict itself, so it gets the unannotated originals. `outcome` is
2747        # passed through so the run is evaluated exactly once.
2748        report_model = build_comparison_report_model(
2749            target_image=resolved_test_image,
2750            control_image=resolved_control_image,  # type: ignore[arg-type]
2751            command=command,
2752            target_result=target_result,
2753            control_result=control_result,
2754            outcome=outcome,
2755            record_comparison=record_comparison,
2756            extra_checks=(
2757                *comparison_checks,
2758                *([record_comparison_skipped] if record_comparison_skipped else []),
2759            ),
2760            extra_diff_blocks=comparison_diff_blocks,
2761            configured_streams=_configured_stream_names(cmd, catalog_file),
2762            connection_objects=_collect_connection_objects(
2763                catalog_file, state_file, target_output / "airbyte_messages"
2764            ),
2765        )
2766        html_report_path = write_html_report(report_model, output_path)
2767        print_success(f"Generated regression HTML report: {html_report_path}")
2768
2769        # Write the concise summary to GITHUB_STEP_SUMMARY (if env var exists)
2770        write_github_summary(render_inline_summary(report_model))
2771
2772        if failed_comparisons:
2773            # Only reachable when both commands ran cleanly -- a comparison is
2774            # skipped otherwise -- so this says what the exit codes could not.
2775            print_error(
2776                f"Regression detected between {resolved_test_image} and "
2777                f"{resolved_control_image}: "
2778                + "; ".join(
2779                    f"{check.name.lower()} — {check.summary}"
2780                    for check in failed_comparisons
2781                )
2782            )
2783        elif outcome.check_improvement:
2784            print_success(
2785                f"Regression test passed for {resolved_test_image} vs "
2786                f"{resolved_control_image}: target connectionStatus SUCCEEDED "
2787                "while control FAILED (improvement)."
2788            )
2789        elif outcome.both_succeeded and comparison_errored:
2790            print_error(
2791                f"Record comparison errored for {resolved_test_image} vs "
2792                f"{resolved_control_image}. Test failed/inconclusive; this "
2793                "is a tooling failure, not evidence of a regression."
2794            )
2795        elif outcome.both_succeeded and record_comparison_failed:
2796            print_error(
2797                f"Regression detected between {resolved_test_image} and "
2798                f"{resolved_control_image}: both versions ran, but the "
2799                "record comparison failed."
2800            )
2801        elif outcome.both_succeeded:
2802            print_success(
2803                f"Regression test passed for {resolved_test_image} vs {resolved_control_image}"
2804            )
2805        elif outcome.regression_detected:
2806            print_error(
2807                f"Regression detected between {resolved_test_image} and {resolved_control_image}"
2808            )
2809        elif outcome.both_failed:
2810            print_error(
2811                f"Both versions failed for {resolved_test_image} vs {resolved_control_image}. "
2812                "Test failed/inconclusive; cannot rule out a regression."
2813            )
2814        else:
2815            print_error(
2816                f"Control {resolved_control_image} failed while target "
2817                f"{resolved_test_image} succeeded. Test failed/inconclusive; "
2818                "nothing to compare the target against."
2819            )
2820
2821
2822def _configured_stream_names(
2823    command: Command, catalog_file: Path | None
2824) -> tuple[str, ...]:
2825    """Every stream the run was configured to read.
2826
2827    The report's per-stream tables are keyed on emitted records, so a stream
2828    that returned nothing has no entry and would be missing from the report
2829    entirely. Reading the configured catalog is what lets it be shown as zero.
2830
2831    Only for commands that read records. `--connection-id` fetches a catalog for
2832    every command, and handing it to `spec`, `check` or `discover` would give
2833    them an all-zero stream table and a coverage warning saying no stream
2834    returned anything -- true, meaningless, and enough noise to train a reviewer
2835    to ignore the warning on the one command where it means something.
2836    """
2837    if command not in (Command.READ, Command.READ_WITH_STATE):
2838        return ()
2839
2840    if catalog_file is None or not catalog_file.exists():
2841        return ()
2842
2843    try:
2844        catalog = ConfiguredAirbyteCatalog.parse_raw(catalog_file.read_text())
2845    except Exception as exc:
2846        print_warning(f"Could not read the configured catalog for the report: {exc}")
2847        return ()
2848
2849    return tuple(
2850        sorted(
2851            configured.stream.name
2852            for configured in catalog.streams
2853            if configured.stream
2854        )
2855    )
2856
2857
2858def _collect_connection_objects(
2859    catalog_file: Path | None,
2860    state_file: Path | None,
2861    messages_dir: Path | None,
2862) -> tuple[DiffBlock, ...]:
2863    """Gather the objects the run was given or discovered, for the report.
2864
2865    The connector config is deliberately absent: it is the one object that
2866    carries secrets, and the report is an artifact anyone with repo access can
2867    download.
2868
2869    Args:
2870        catalog_file: The configured catalog handed to the run.
2871        state_file: The state handed to the run, for a warm read.
2872        messages_dir: The target side's `airbyte_messages` directory, which
2873            holds the catalog the connector discovered.
2874
2875    Returns:
2876        One JSON block per object that exists, in the order a reader wants them.
2877    """
2878    sources: list[tuple[str, Path | None]] = [
2879        ("Configured catalog (input)", catalog_file),
2880        ("State (input)", state_file),
2881        (
2882            "Discovered catalog (from the target)",
2883            (messages_dir / "catalog.jsonl") if messages_dir else None,
2884        ),
2885    ]
2886
2887    blocks: list[DiffBlock] = []
2888    for title, path in sources:
2889        if path is None or not path.is_file():
2890            continue
2891
2892        try:
2893            payload = path.read_text()
2894        except OSError:
2895            continue
2896
2897        if payload.strip():
2898            blocks.append(build_json_block(title, payload))
2899
2900    return tuple(blocks)
2901
2902
2903# The order commands are presented in, matching how a sync runs them. Anything
2904# else found under the artifacts directory follows, sorted.
2905_REPORT_COMMAND_ORDER = ("spec", "check", "discover", "read")
2906
2907
2908@connector_app.command(name="consolidate-regression-reports")
2909def consolidate_regression_reports(
2910    artifacts_dir: Annotated[
2911        str,
2912        Parameter(
2913            help="Directory holding one subdirectory per command, each with the "
2914            "`report.html` that `regression-test --output-dir` wrote."
2915        ),
2916    ],
2917    output_path: Annotated[
2918        str | None,
2919        Parameter(
2920            name=["--output"],
2921            help="Where to write the consolidated report. "
2922            "Defaults to `report.html` in `--artifacts-dir`.",
2923        ),
2924    ] = None,
2925) -> None:
2926    """Fold one run's per-command regression reports into a single page.
2927
2928    The regression test runs the CLI once per Airbyte command, so each command
2929    writes its own report into its own artifact. This produces the run-level
2930    view -- a verdict table across commands plus every report inlined -- so a
2931    reviewer downloads one file instead of four.
2932
2933    Never fails the run: a command whose report is missing or unreadable is
2934    skipped, because this is a reporting convenience and the per-command
2935    artifacts remain the source of truth.
2936    """
2937    artifacts_path = Path(artifacts_dir)
2938    ordered = [artifacts_path / command for command in _REPORT_COMMAND_ORDER]
2939    extra = sorted(
2940        path
2941        for path in artifacts_path.glob("*")
2942        if path.is_dir() and path not in ordered
2943    )
2944    report_paths = [
2945        directory / HTML_REPORT_FILENAME for directory in [*ordered, *extra]
2946    ]
2947
2948    consolidated = consolidate_reports(report_paths)
2949    if consolidated is None:
2950        print_warning(
2951            f"No per-command reports found under {artifacts_path}; "
2952            "nothing to consolidate."
2953        )
2954        return
2955
2956    destination = Path(output_path) if output_path else artifacts_path / "report.html"
2957    destination.parent.mkdir(parents=True, exist_ok=True)
2958    destination.write_text(consolidated)
2959    print_success(f"Consolidated regression report: {destination}")
2960
2961
2962@connector_app.command(name="fetch-connection-config")
2963def fetch_connection_config_cmd(
2964    connection_id: Annotated[
2965        str,
2966        Parameter(help="The UUID of the Airbyte Cloud connection."),
2967    ],
2968    output_path: Annotated[
2969        str | None,
2970        Parameter(
2971            help="Path to output file or directory. "
2972            "If directory, writes connection-<id>-config.json inside it. "
2973            "Default: platform temp directory (e.g. /tmp/connection-<id>-config.json)"
2974        ),
2975    ] = None,
2976    with_secrets: Annotated[
2977        bool,
2978        Parameter(
2979            name="--with-secrets",
2980            negative="--no-secrets",
2981            help="If set, fetches unmasked secrets from the internal database. "
2982            "Requires GCP_PROD_DB_ACCESS_CREDENTIALS env var or `gcloud auth application-default login`. "
2983            "Must be used with --oc-issue-url.",
2984        ),
2985    ] = False,
2986    oc_issue_url: Annotated[
2987        str | None,
2988        Parameter(
2989            help="OC issue URL for audit logging. Required when using --with-secrets."
2990        ),
2991    ] = None,
2992) -> None:
2993    """Fetch connection configuration from Airbyte Cloud to a local file.
2994
2995    This command retrieves the source configuration for a given connection ID
2996    and writes it to a JSON file. When `--output-path` is omitted the file is
2997    written to the platform temp directory to avoid accidentally committing
2998    secrets to a git repository.
2999
3000    Requires authentication via AIRBYTE_CLOUD_CLIENT_ID and
3001    AIRBYTE_CLOUD_CLIENT_SECRET environment variables.
3002
3003    When --with-secrets is specified, the command fetches unmasked secrets from
3004    the internal database using the connection-retriever. This additionally requires:
3005    - An OC issue URL for audit logging (--oc-issue-url)
3006    - GCP credentials via `GCP_PROD_DB_ACCESS_CREDENTIALS` env var or `gcloud auth application-default login`
3007    - Cloud SQL Python Connector access to the Prod DB replica.
3008    """
3009    path = Path(output_path) if output_path else None
3010    result = fetch_connection_config(
3011        connection_id=connection_id,
3012        output_path=path,
3013        with_secrets=with_secrets,
3014        oc_issue_url=oc_issue_url,
3015    )
3016    if result.success:
3017        print_success(result.message)
3018    else:
3019        print_error(result.message)
3020    print_json(result.model_dump())
3021
3022
3023def _read_json_input(
3024    json_str: str | None,
3025    input_file: str | None,
3026) -> dict:
3027    """Read JSON from a positional arg, --input file, or STDIN (in that priority)."""
3028    if json_str is not None:
3029        return json.loads(json_str)
3030    if input_file is not None:
3031        return json.loads(Path(input_file).read_text())
3032    if not sys.stdin.isatty():
3033        return json.loads(sys.stdin.read())
3034    exit_with_error(
3035        "No JSON input provided. Pass it as a positional argument, "
3036        "via --input <file>, or pipe to STDIN."
3037    )
3038    raise SystemExit(1)  # unreachable, but satisfies type checker
3039
3040
3041@state_app.command(name="get")
3042def get_connection_state(
3043    connection_id: Annotated[
3044        str,
3045        Parameter(help="The connection ID (UUID) to fetch state for."),
3046    ],
3047    stream_name: Annotated[
3048        str | None,
3049        Parameter(
3050            help="Optional stream name to filter state for a single stream.",
3051            name="--stream",
3052        ),
3053    ] = None,
3054    stream_namespace: Annotated[
3055        str | None,
3056        Parameter(
3057            help="Optional stream namespace to narrow the stream filter.",
3058            name="--namespace",
3059        ),
3060    ] = None,
3061    output: Annotated[
3062        str | None,
3063        Parameter(
3064            help="Path to write the state JSON output to a file.",
3065            name="--output",
3066        ),
3067    ] = None,
3068) -> None:
3069    """Get the current state for an Airbyte Cloud connection."""
3070    workspace = CloudWorkspace.from_env()
3071    conn = workspace.get_connection(connection_id)
3072
3073    if stream_name is not None:
3074        stream_state = conn.get_stream_state(
3075            stream_name=stream_name,
3076            stream_namespace=stream_namespace,
3077        )
3078        result = {
3079            "connection_id": connection_id,
3080            "stream_name": stream_name,
3081            "stream_namespace": stream_namespace,
3082            "stream_state": stream_state,
3083        }
3084    else:
3085        result = conn.dump_raw_state()
3086
3087    json_output = json.dumps(result, indent=2, default=str)
3088    if output is not None:
3089        Path(output).write_text(json_output + "\n")
3090        print(f"State written to {output}", file=sys.stderr)
3091    else:
3092        print(json_output)
3093
3094
3095@state_app.command(name="set")
3096def set_connection_state(
3097    connection_id: Annotated[
3098        str,
3099        Parameter(help="The connection ID (UUID) to update state for."),
3100    ],
3101    state_json: Annotated[
3102        str | None,
3103        Parameter(
3104            help="The connection state as a JSON string. When --stream is used, "
3105            'this is just the stream\'s state blob (e.g., \'{"cursor": "2024-01-01"}\').'
3106            " Otherwise, must include 'stateType', 'connectionId', and the appropriate "
3107            "state field. Can also be provided via --input or STDIN."
3108        ),
3109    ] = None,
3110    stream_name: Annotated[
3111        str | None,
3112        Parameter(
3113            help="Optional stream name to update state for a single stream only.",
3114            name="--stream",
3115        ),
3116    ] = None,
3117    stream_namespace: Annotated[
3118        str | None,
3119        Parameter(
3120            help="Optional stream namespace to identify the stream.",
3121            name="--namespace",
3122        ),
3123    ] = None,
3124    input_file: Annotated[
3125        str | None,
3126        Parameter(
3127            help="Path to a JSON file containing the state to set.",
3128            name="--input",
3129        ),
3130    ] = None,
3131) -> None:
3132    """Set the state for an Airbyte Cloud connection.
3133
3134    State JSON can be provided as a positional argument, via --input <file>,
3135    or piped through STDIN.
3136
3137    Uses the safe variant that prevents updates while a sync is running.
3138    When --stream is provided, only that stream's state is updated within
3139    the existing connection state.
3140    """
3141    parsed_state = _read_json_input(state_json, input_file)
3142    workspace = CloudWorkspace.from_env()
3143    conn = workspace.get_connection(connection_id)
3144
3145    if stream_name is not None:
3146        conn.set_stream_state(
3147            stream_name=stream_name,
3148            state_blob_dict=parsed_state,
3149            stream_namespace=stream_namespace,
3150        )
3151    else:
3152        conn.import_raw_state(parsed_state)
3153
3154    result = conn.dump_raw_state()
3155    json_output = json.dumps(result, indent=2, default=str)
3156    print(json_output)
3157
3158
3159@state_app.command(name="reset")
3160def reset_stream_state(
3161    connection_id: Annotated[
3162        str,
3163        Parameter(help="The connection ID (UUID) to update state for."),
3164    ],
3165    stream_name: Annotated[
3166        str,
3167        Parameter(
3168            help="The configured stream name whose state should be reset.",
3169            name="--stream",
3170        ),
3171    ],
3172    stream_namespace: Annotated[
3173        str | None,
3174        Parameter(
3175            help="Optional stream namespace to identify the stream.",
3176            name="--namespace",
3177        ),
3178    ] = None,
3179) -> None:
3180    """Reset a configured stream's state so the next sync full-refreshes it.
3181
3182    Uses the safe variant that prevents updates while a sync is running.
3183    Returns a restorable `previous_state_backup` in raw Config API format.
3184    """
3185    workspace = CloudWorkspace.from_env()
3186    conn = workspace.get_connection(connection_id)
3187
3188    try:
3189        result = reset_cloud_stream_state(
3190            conn,
3191            stream_name=stream_name,
3192            stream_namespace=stream_namespace,
3193        )
3194    except PyAirbyteInputError as e:
3195        exit_with_error(str(e))
3196
3197    print(json.dumps(result.model_dump(), indent=2, default=str))
3198
3199
3200@catalog_app.command(name="get")
3201def get_connection_catalog(
3202    connection_id: Annotated[
3203        str,
3204        Parameter(help="The connection ID (UUID) to fetch catalog for."),
3205    ],
3206    output: Annotated[
3207        str | None,
3208        Parameter(
3209            help="Path to write the catalog JSON output to a file.",
3210            name="--output",
3211        ),
3212    ] = None,
3213) -> None:
3214    """Get the configured catalog for an Airbyte Cloud connection."""
3215    workspace = CloudWorkspace.from_env()
3216    conn = workspace.get_connection(connection_id)
3217    result = conn.dump_raw_catalog()
3218    if result is None:
3219        exit_with_error("No configured catalog found for this connection.")
3220        raise SystemExit(1)  # unreachable, but satisfies type checker
3221
3222    json_output = json.dumps(result, indent=2, default=str)
3223    if output is not None:
3224        Path(output).write_text(json_output + "\n")
3225        print(f"Catalog written to {output}", file=sys.stderr)
3226    else:
3227        print(json_output)
3228
3229
3230@catalog_app.command(name="set")
3231def set_connection_catalog(
3232    connection_id: Annotated[
3233        str,
3234        Parameter(help="The connection ID (UUID) to update catalog for."),
3235    ],
3236    catalog_json: Annotated[
3237        str | None,
3238        Parameter(
3239            help="The configured catalog as a JSON string. "
3240            "Can also be provided via --input or STDIN."
3241        ),
3242    ] = None,
3243    input_file: Annotated[
3244        str | None,
3245        Parameter(
3246            help="Path to a JSON file containing the catalog to set.",
3247            name="--input",
3248        ),
3249    ] = None,
3250) -> None:
3251    """Set the configured catalog for an Airbyte Cloud connection.
3252
3253    Catalog JSON can be provided as a positional argument, via --input <file>,
3254    or piped through STDIN.
3255
3256    WARNING: This replaces the entire configured catalog.
3257    """
3258    parsed_catalog = _read_json_input(catalog_json, input_file)
3259    workspace = CloudWorkspace.from_env()
3260    conn = workspace.get_connection(connection_id)
3261    conn.import_raw_catalog(parsed_catalog)
3262
3263    result = conn.dump_raw_catalog()
3264    if result is None:
3265        exit_with_error("Failed to retrieve catalog after import.")
3266        raise SystemExit(1)  # unreachable, but satisfies type checker
3267
3268    json_output = json.dumps(result, indent=2, default=str)
3269    print(json_output)
3270
3271
3272@logs_app.command(name="lookup-cloud-backend-error")
3273def lookup_cloud_backend_error(
3274    error_id: Annotated[
3275        str,
3276        Parameter(
3277            help=(
3278                "The error ID (UUID) to search for. This is typically returned "
3279                "in API error responses as {'errorId': '...'}"
3280            )
3281        ),
3282    ],
3283    lookback_days: Annotated[
3284        int,
3285        Parameter(help="Number of days to look back in logs."),
3286    ] = 7,
3287    min_severity_filter: Annotated[
3288        GCPSeverity | None,
3289        Parameter(
3290            help="Optional minimum severity level to filter logs.",
3291        ),
3292    ] = None,
3293    raw: Annotated[
3294        bool,
3295        Parameter(help="Output raw JSON instead of formatted text."),
3296    ] = False,
3297) -> None:
3298    """Look up error details from GCP Cloud Logging by error ID.
3299
3300    When an Airbyte Cloud API returns an error response with only an error ID
3301    (e.g., {"errorId": "3173452e-8f22-4286-a1ec-b0f16c1e078a"}), this command
3302    fetches the full stack trace and error details from GCP Cloud Logging.
3303
3304    Requires GCP credentials with Logs Viewer role on the target project.
3305    Set up credentials with: gcloud auth application-default login
3306    """
3307    print(f"Searching for error ID: {error_id}", file=sys.stderr)
3308    print(f"Lookback days: {lookback_days}", file=sys.stderr)
3309    if min_severity_filter:
3310        print(f"Severity filter: {min_severity_filter}", file=sys.stderr)
3311    print(file=sys.stderr)
3312
3313    result = fetch_error_logs(
3314        error_id=error_id,
3315        lookback_days=lookback_days,
3316        min_severity_filter=min_severity_filter,
3317    )
3318
3319    if raw:
3320        print_json(result.model_dump())
3321        return
3322
3323    print(f"Found {result.total_entries_found} log entries", file=sys.stderr)
3324    print(file=sys.stderr)
3325
3326    if result.payloads:
3327        for i, payload in enumerate(result.payloads):
3328            print(f"=== Log Group {i + 1} ===")
3329            print(f"Timestamp: {payload.timestamp}")
3330            print(f"Severity: {payload.severity}")
3331            if payload.resource.labels.pod_name:
3332                print(f"Pod: {payload.resource.labels.pod_name}")
3333            print(f"Lines: {payload.num_log_lines}")
3334            print()
3335            print(payload.message)
3336            print()
3337    elif result.entries:
3338        print("No grouped payloads, showing raw entries:", file=sys.stderr)
3339        for entry in result.entries:
3340            print(f"[{entry.timestamp}] {entry.severity}: {entry.payload}")
3341    else:
3342        print_error("No log entries found for this error ID.")
3343
3344
3345# =============================================================================
3346# Organization commands
3347# =============================================================================
3348
3349
3350@organization_app.command(name="search")
3351def organization_search(
3352    name_contains: Annotated[
3353        str,
3354        Parameter(
3355            name="--name-contains",
3356            help="Case-insensitive substring to match against organization name or email.",
3357        ),
3358    ],
3359    limit: Annotated[
3360        int,
3361        Parameter(
3362            name="--limit",
3363            help="Maximum number of results to return.",
3364        ),
3365    ] = 20,
3366) -> None:
3367    """Search organizations by name or email substring."""
3368    rows = _search_organizations(name_contains=name_contains, limit=limit)
3369    if not rows:
3370        print_error(f"No organizations found matching '{name_contains}'.")
3371        return
3372
3373    print_json([dict(row) for row in rows])
3374
3375
3376@organization_app.command(name="info")
3377def organization_info(
3378    organization_id: Annotated[
3379        str,
3380        Parameter(help="The organization UUID."),
3381    ],
3382) -> None:
3383    """Get information about an organization (stub — not yet implemented)."""
3384    exit_with_error(
3385        "organization info is not yet implemented. "
3386        "Use 'airbyte-ops cloud organization search' to find organizations by name."
3387    )
3388
3389
3390# =============================================================================
3391# Workspace commands
3392# =============================================================================
3393
3394
3395@workspace_app.command(name="search")
3396def workspace_search(
3397    name_contains: Annotated[
3398        str | None,
3399        Parameter(
3400            name="--name-contains",
3401            help="Case-insensitive substring to match against workspace name or slug.",
3402        ),
3403    ] = None,
3404    email_domain: Annotated[
3405        str | None,
3406        Parameter(
3407            name="--email-domain",
3408            help="Email domain to search for (e.g. 'motherduck.com').",
3409        ),
3410    ] = None,
3411    limit: Annotated[
3412        int,
3413        Parameter(
3414            name="--limit",
3415            help="Maximum number of results to return.",
3416        ),
3417    ] = 100,
3418) -> None:
3419    """Search workspaces by name substring or email domain."""
3420    if not name_contains and not email_domain:
3421        exit_with_error(
3422            "At least one of --name-contains or --email-domain must be provided."
3423        )
3424
3425    if name_contains:
3426        rows = _search_workspaces(name_contains=name_contains, limit=limit)
3427    else:
3428        rows = query_workspaces_by_email_domain(
3429            email_domain=email_domain.lstrip("@"),  # type: ignore[union-attr]
3430            limit=limit,
3431        )
3432
3433    if not rows:
3434        search_term = name_contains or email_domain
3435        print_error(f"No workspaces found matching '{search_term}'.")
3436        return
3437
3438    print_json([dict(row) for row in rows])
3439
3440
3441@workspace_app.command(name="info")
3442def workspace_info_cmd(
3443    workspace_id: Annotated[
3444        str,
3445        Parameter(help="The workspace UUID."),
3446    ],
3447) -> None:
3448    """Get information about a workspace."""
3449    result = query_workspace_info(workspace_id=workspace_id)
3450    if not result:
3451        print_error(f"No workspace found with ID '{workspace_id}'.")
3452        return
3453
3454    print_json(dict(result))