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-idand--actor-definition-id.organization: pins ALL instances of a connector type across an organization. Requires--organization-idand--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-idand--actor-definition-id.organization: clears the organization-level pin for a connector type. Requires--organization-idand--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:
- --test-image: Use a pre-built image from Docker registry
- --connector-name: Build the image locally from source code
- --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 ifskip_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 toTruein 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 asLOG_LEVEL=DEBUGto 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 toTruewhen--connection-idis provided,Falseotherwise. Has no effect unless the command isread. Ignored when--state-pathis 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 thereport.htmlthatregression-test --output-dirwrote. [required]OUTPUT, --output: Where to write the consolidated report. Defaults toreport.htmlin--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_CREDENTIALSenv var orgcloud 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 orgcloud 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
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
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
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
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))