airbyte_ops_mcp.registry
Registry operations for Airbyte connectors.
This package provides functionality for reading, listing, publishing, compiling, and validating connector registry artifacts.
1# Copyright (c) 2025 Airbyte, Inc., all rights reserved. 2"""Registry operations for Airbyte connectors. 3 4This package provides functionality for reading, listing, publishing, compiling, 5and validating connector registry artifacts. 6""" 7 8from __future__ import annotations 9 10from airbyte_ops_mcp.registry._constants import ( 11 DEFAULT_METADATA_SERVICE_BUCKET_NAME, 12 DEV_METADATA_SERVICE_BUCKET_NAME, 13 LATEST_GCS_FOLDER_NAME, 14 METADATA_FILE_NAME, 15 METADATA_FOLDER, 16 PROD_METADATA_SERVICE_BUCKET_NAME, 17 RELEASE_CANDIDATE_GCS_FOLDER_NAME, 18 SONAR_DEV_BUCKET_NAME, 19 SONAR_PROD_BUCKET_NAME, 20) 21from airbyte_ops_mcp.registry._enums import ( 22 ConnectorLanguage, 23 ConnectorType, 24 SupportLevel, 25) 26from airbyte_ops_mcp.registry.audit import ( 27 AuditResult, 28 UnpublishedConnector, 29 find_unpublished_connectors, 30) 31from airbyte_ops_mcp.registry.compile import ( 32 CompileResult, 33 PurgeLatestResult, 34 compile_registry, 35 purge_latest_dirs, 36) 37from airbyte_ops_mcp.registry.generate import ( 38 GenerateResult, 39 generate_version_artifacts, 40) 41from airbyte_ops_mcp.registry.models import ( 42 ConnectorListResult, 43 ConnectorMetadata, 44 MetadataPublishResult, 45 RegistryEntryResult, 46 VersionListResult, 47) 48from airbyte_ops_mcp.registry.operations import ( 49 get_registry_entry, 50 get_registry_spec, 51 list_connector_versions, 52 list_registry_connectors, 53 list_registry_connectors_filtered, 54) 55from airbyte_ops_mcp.registry.publish import ( 56 CONNECTOR_PATH_PREFIX, 57 get_connector_metadata, 58 get_gcs_publish_path, 59 publish_connector_metadata, 60) 61from airbyte_ops_mcp.registry.publish_artifacts import ( 62 PublishArtifactsResult, 63 publish_version_artifacts, 64) 65from airbyte_ops_mcp.registry.rebuild import ( 66 OutputMode, 67 RebuildResult, 68 rebuild_registry, 69) 70from airbyte_ops_mcp.registry.registry_store_base import ( 71 Registry, 72 get_registry, 73) 74from airbyte_ops_mcp.registry.store import ( 75 REGISTRY_STORE_ENV_VAR, 76 RegistryStore, 77 StoreType, 78 resolve_registry_store, 79) 80from airbyte_ops_mcp.registry.validate import ( 81 ValidateOptions, 82 ValidationResult, 83 validate_metadata, 84) 85from airbyte_ops_mcp.registry.yank import ( 86 YANK_FILE_NAME, 87 YankResult, 88 unyank_connector_version, 89 yank_connector_version, 90) 91 92__all__ = [ 93 "CONNECTOR_PATH_PREFIX", 94 "DEFAULT_METADATA_SERVICE_BUCKET_NAME", 95 "DEV_METADATA_SERVICE_BUCKET_NAME", 96 "LATEST_GCS_FOLDER_NAME", 97 "METADATA_FILE_NAME", 98 "METADATA_FOLDER", 99 "PROD_METADATA_SERVICE_BUCKET_NAME", 100 "REGISTRY_STORE_ENV_VAR", 101 "RELEASE_CANDIDATE_GCS_FOLDER_NAME", 102 "SONAR_DEV_BUCKET_NAME", 103 "SONAR_PROD_BUCKET_NAME", 104 "YANK_FILE_NAME", 105 "AuditResult", 106 "CompileResult", 107 "ConnectorLanguage", 108 "ConnectorListResult", 109 "ConnectorMetadata", 110 "ConnectorType", 111 "GenerateResult", 112 "MetadataPublishResult", 113 "OutputMode", 114 "PublishArtifactsResult", 115 "PurgeLatestResult", 116 "RebuildResult", 117 "Registry", 118 "RegistryEntryResult", 119 "RegistryStore", 120 "StoreType", 121 "SupportLevel", 122 "UnpublishedConnector", 123 "ValidateOptions", 124 "ValidationResult", 125 "VersionListResult", 126 "YankResult", 127 "compile_registry", 128 "find_unpublished_connectors", 129 "generate_version_artifacts", 130 "get_connector_metadata", 131 "get_gcs_publish_path", 132 "get_registry", 133 "get_registry_entry", 134 "get_registry_spec", 135 "list_connector_versions", 136 "list_registry_connectors", 137 "list_registry_connectors_filtered", 138 "publish_connector_metadata", 139 "publish_version_artifacts", 140 "purge_latest_dirs", 141 "rebuild_registry", 142 "resolve_registry_store", 143 "unyank_connector_version", 144 "validate_metadata", 145 "yank_connector_version", 146]
37@dataclass 38class AuditResult: 39 """Result of auditing which connectors have unpublished versions.""" 40 41 unpublished: list[UnpublishedConnector] = field(default_factory=list) 42 checked_count: int = 0 43 skipped_archived: list[str] = field(default_factory=list) 44 skipped_rc: list[str] = field(default_factory=list) 45 skipped_disabled: list[str] = field(default_factory=list) 46 errors: list[str] = field(default_factory=list)
Result of auditing which connectors have unpublished versions.
230@dataclass 231class CompileResult: 232 """Result of a registry compile operation.""" 233 234 target: str 235 connectors_scanned: int = 0 236 versions_found: int = 0 237 yanked_versions: int = 0 238 latest_updated: int = 0 239 latest_already_current: int = 0 240 cloud_registry_entries: int = 0 241 oss_registry_entries: int = 0 242 composite_registry_entries: int = 0 243 metrics_connector_count: int = 0 244 metrics_registry_entries: int = 0 245 metrics_source: str | None = None 246 metrics_error: str | None = None 247 version_indexes_written: int = 0 248 version_indexes_skipped: int = 0 249 specs_secrets_mask_properties: int = 0 250 errors: list[str] = field(default_factory=list) 251 dry_run: bool = False 252 253 @property 254 def status(self) -> str: 255 if self.dry_run: 256 return "dry-run" 257 if self.errors: 258 return "completed-with-errors" 259 return "success" 260 261 def summary(self) -> str: 262 return ( 263 f"[{self.status}] Scanned {self.connectors_scanned} connectors, " 264 f"{self.versions_found} versions ({self.yanked_versions} yanked). " 265 f"Latest updated: {self.latest_updated}, " 266 f"already current: {self.latest_already_current}. " 267 f"Registry entries: cloud={self.cloud_registry_entries}, " 268 f"oss={self.oss_registry_entries}, " 269 f"composite={self.composite_registry_entries}. " 270 f"Metrics loaded for {self.metrics_connector_count} connectors, " 271 f"injected into {self.metrics_registry_entries} registry entries. " 272 f"Version indexes: {self.version_indexes_written} written, " 273 f"{self.version_indexes_skipped} unchanged. " 274 f"Specs secrets mask: {self.specs_secrets_mask_properties} properties. " 275 f"Errors: {len(self.errors)}." 276 )
Result of a registry compile operation.
261 def summary(self) -> str: 262 return ( 263 f"[{self.status}] Scanned {self.connectors_scanned} connectors, " 264 f"{self.versions_found} versions ({self.yanked_versions} yanked). " 265 f"Latest updated: {self.latest_updated}, " 266 f"already current: {self.latest_already_current}. " 267 f"Registry entries: cloud={self.cloud_registry_entries}, " 268 f"oss={self.oss_registry_entries}, " 269 f"composite={self.composite_registry_entries}. " 270 f"Metrics loaded for {self.metrics_connector_count} connectors, " 271 f"injected into {self.metrics_registry_entries} registry entries. " 272 f"Version indexes: {self.version_indexes_written} written, " 273 f"{self.version_indexes_skipped} unchanged. " 274 f"Specs secrets mask: {self.specs_secrets_mask_properties} properties. " 275 f"Errors: {len(self.errors)}." 276 )
88class ConnectorLanguage(StrEnum): 89 """Connector implementation languages.""" 90 91 PYTHON = "python" 92 JAVA = "java" 93 LOW_CODE = "low-code" 94 MANIFEST_ONLY = "manifest-only" 95 96 @classmethod 97 def parse(cls, value: str) -> ConnectorLanguage: 98 """Parse a string into a `ConnectorLanguage`, raising `ValueError` on mismatch.""" 99 try: 100 return cls(value) 101 except ValueError: 102 valid = ", ".join(f"`{m.value}`" for m in cls) 103 raise ValueError( 104 f"Unrecognized language: {value!r}. Expected one of: {valid}." 105 ) from None
Connector implementation languages.
96 @classmethod 97 def parse(cls, value: str) -> ConnectorLanguage: 98 """Parse a string into a `ConnectorLanguage`, raising `ValueError` on mismatch.""" 99 try: 100 return cls(value) 101 except ValueError: 102 valid = ", ".join(f"`{m.value}`" for m in cls) 103 raise ValueError( 104 f"Unrecognized language: {value!r}. Expected one of: {valid}." 105 ) from None
Parse a string into a ConnectorLanguage, raising ValueError on mismatch.
73class ConnectorListResult(BaseModel): 74 """Result of listing connectors in the registry.""" 75 76 bucket_name: str = Field(description="The GCS bucket name") 77 connector_count: int = Field(description="Number of connectors found") 78 connectors: list[str] = Field(description="List of connector names")
Result of listing connectors in the registry.
12class ConnectorMetadata(BaseModel): 13 """Connector metadata from metadata.yaml. 14 15 This model represents the essential metadata about a connector 16 read from its metadata.yaml file in the Airbyte monorepo. 17 """ 18 19 name: str = Field(description="The connector technical name") 20 docker_repository: str = Field(description="The Docker repository") 21 docker_image_tag: str = Field(description="The Docker image tag/version") 22 support_level: str | None = Field( 23 default=None, description="The support level (certified, community, etc.)" 24 ) 25 definition_id: str | None = Field( 26 default=None, description="The connector definition ID" 27 )
Connector metadata from metadata.yaml.
This model represents the essential metadata about a connector read from its metadata.yaml file in the Airbyte monorepo.
70class ConnectorType(StrEnum): 71 """Connector type: source or destination.""" 72 73 SOURCE = "source" 74 DESTINATION = "destination" 75 76 @classmethod 77 def parse(cls, value: str) -> ConnectorType: 78 """Parse a string into a `ConnectorType`, raising `ValueError` on mismatch.""" 79 try: 80 return cls(value) 81 except ValueError: 82 valid = ", ".join(f"`{m.value}`" for m in cls) 83 raise ValueError( 84 f"Unrecognized connector type: {value!r}. Expected one of: {valid}." 85 ) from None
Connector type: source or destination.
76 @classmethod 77 def parse(cls, value: str) -> ConnectorType: 78 """Parse a string into a `ConnectorType`, raising `ValueError` on mismatch.""" 79 try: 80 return cls(value) 81 except ValueError: 82 valid = ", ".join(f"`{m.value}`" for m in cls) 83 raise ValueError( 84 f"Unrecognized connector type: {value!r}. Expected one of: {valid}." 85 ) from None
Parse a string into a ConnectorType, raising ValueError on mismatch.
129@dataclass 130class GenerateResult: 131 """Result of a local artifact generation run.""" 132 133 connector_name: str 134 version: str 135 docker_image: str 136 output_dir: str 137 artifacts_written: list[str] = field(default_factory=list) 138 errors: list[str] = field(default_factory=list) 139 validation_errors: list[str] = field(default_factory=list) 140 dry_run: bool = False 141 142 @property 143 def success(self) -> bool: 144 return len(self.errors) == 0 and len(self.validation_errors) == 0 145 146 def to_dict(self) -> dict[str, Any]: 147 return { 148 "connector_name": self.connector_name, 149 "version": self.version, 150 "docker_image": self.docker_image, 151 "output_dir": self.output_dir, 152 "artifacts_written": self.artifacts_written, 153 "errors": self.errors, 154 "validation_errors": self.validation_errors, 155 "dry_run": self.dry_run, 156 "success": self.success, 157 }
Result of a local artifact generation run.
146 def to_dict(self) -> dict[str, Any]: 147 return { 148 "connector_name": self.connector_name, 149 "version": self.version, 150 "docker_image": self.docker_image, 151 "output_dir": self.output_dir, 152 "artifacts_written": self.artifacts_written, 153 "errors": self.errors, 154 "validation_errors": self.validation_errors, 155 "dry_run": self.dry_run, 156 "success": self.success, 157 }
30class MetadataPublishResult(BaseModel): 31 """Result of a metadata publish operation to GCS. 32 33 This model provides detailed information about the outcome of 34 publishing connector metadata to the registry. 35 """ 36 37 connector_name: str = Field(description="The connector technical name") 38 version: str = Field(description="The version that was published") 39 bucket_name: str = Field(description="The GCS bucket name") 40 versioned_path: str = Field(description="The versioned GCS path") 41 latest_path: str | None = Field( 42 default=None, description="The latest GCS path if updated" 43 ) 44 versioned_uploaded: bool = Field( 45 default=False, description="Whether the versioned metadata was uploaded" 46 ) 47 latest_uploaded: bool = Field( 48 default=False, description="Whether the latest metadata was uploaded" 49 ) 50 status: Literal["success", "dry-run", "already-up-to-date"] = Field( 51 description="The status of the operation" 52 ) 53 message: str = Field(description="Status message describing the outcome") 54 55 def __str__(self) -> str: 56 """Return a string representation of the publish result.""" 57 return f"[{self.status}] {self.connector_name}:{self.version} -> {self.versioned_path}"
Result of a metadata publish operation to GCS.
This model provides detailed information about the outcome of publishing connector metadata to the registry.
55@dataclass 56class PublishArtifactsResult: 57 """Result of a version-artifacts publish operation.""" 58 59 connector_name: str 60 version: str 61 target: str 62 gcs_destination: str 63 files_uploaded: list[str] = field(default_factory=list) 64 errors: list[str] = field(default_factory=list) 65 validation_errors: list[str] = field(default_factory=list) 66 progressive_rollout_overridden_by_breaking_change: bool = False 67 progressive_rollout_overridden_by_published_ga: bool = False 68 dry_run: bool = False 69 70 @property 71 def success(self) -> bool: 72 return len(self.errors) == 0 and len(self.validation_errors) == 0 73 74 @property 75 def status(self) -> str: 76 if self.dry_run: 77 return "dry-run" 78 if self.errors or self.validation_errors: 79 return "completed-with-errors" 80 return "success"
Result of a version-artifacts publish operation.
279@dataclass 280class PurgeLatestResult: 281 """Result of a purge-latest operation.""" 282 283 target: str 284 connectors_found: int = 0 285 latest_dirs_deleted: int = 0 286 errors: list[str] = field(default_factory=list) 287 dry_run: bool = False 288 289 @property 290 def status(self) -> str: 291 if self.dry_run: 292 return "dry-run" 293 if self.errors: 294 return "completed-with-errors" 295 return "success" 296 297 def summary(self) -> str: 298 return ( 299 f"[{self.status}] Found {self.connectors_found} connectors, " 300 f"deleted {self.latest_dirs_deleted} latest/ directories. " 301 f"Errors: {len(self.errors)}." 302 )
Result of a purge-latest operation.
43@dataclass 44class RebuildResult: 45 """Result of a registry rebuild operation.""" 46 47 source_bucket: str 48 output_mode: OutputMode 49 output_root: str 50 connectors_processed: int = 0 51 blobs_copied: int = 0 52 blobs_skipped: int = 0 53 errors: list[str] = field(default_factory=list) 54 dry_run: bool = False 55 56 @property 57 def status(self) -> str: 58 """Return the status of the rebuild operation.""" 59 if self.dry_run: 60 return "dry-run" 61 if self.errors: 62 return "completed-with-errors" 63 return "success" 64 65 def summary(self) -> str: 66 """Return a human-readable summary.""" 67 return ( 68 f"[{self.status}] Rebuilt {self.connectors_processed} connectors, " 69 f"{self.blobs_copied} blobs copied, {self.blobs_skipped} skipped, " 70 f"{len(self.errors)} errors. Output: {self.output_root}" 71 )
Result of a registry rebuild operation.
56 @property 57 def status(self) -> str: 58 """Return the status of the rebuild operation.""" 59 if self.dry_run: 60 return "dry-run" 61 if self.errors: 62 return "completed-with-errors" 63 return "success"
Return the status of the rebuild operation.
65 def summary(self) -> str: 66 """Return a human-readable summary.""" 67 return ( 68 f"[{self.status}] Rebuilt {self.connectors_processed} connectors, " 69 f"{self.blobs_copied} blobs copied, {self.blobs_skipped} skipped, " 70 f"{len(self.errors)} errors. Output: {self.output_root}" 71 )
Return a human-readable summary.
44class Registry(ABC): 45 """A configured connector registry store. 46 47 A Registry is bound to a specific `airbyte_ops_mcp.registry.store.RegistryStore` 48 (store type + env + optional prefix), and provides methods used by the CLI to 49 read/write registry contents. 50 """ 51 52 def __init__(self, store: RegistryStore) -> None: 53 self.store = store 54 55 @property 56 def store_type(self) -> StoreType: 57 return self.store.store_type 58 59 @property 60 def bucket_name(self) -> str: 61 return self.store.bucket 62 63 @property 64 def prefix(self) -> str: 65 return self.store.prefix 66 67 def _require_no_prefix(self, op_name: str) -> None: 68 """Raise if this op doesn't support prefixed targets.""" 69 70 if self.prefix: 71 raise NotImplementedError( 72 f"Operation '{op_name}' does not yet support store prefixes (got prefix='{self.prefix}')." 73 ) 74 75 # --------------------------------------------------------------------- 76 # Read operations 77 # --------------------------------------------------------------------- 78 79 @abstractmethod 80 def list_connectors( 81 self, 82 *, 83 support_level: SupportLevel | None = None, 84 min_support_level: SupportLevel | None = None, 85 connector_type: ConnectorType | None = None, 86 language: ConnectorLanguage | None = None, 87 ) -> list[str]: 88 raise NotImplementedError( 89 _op_not_implemented_message(self.store_type, "list_connectors") 90 ) 91 92 def list_connector_versions(self, connector_name: str) -> list[str]: 93 raise NotImplementedError( 94 _op_not_implemented_message(self.store_type, "list_connector_versions") 95 ) 96 97 def get_connector_metadata( 98 self, connector_name: str, version: str = "latest" 99 ) -> dict[str, Any]: 100 raise NotImplementedError( 101 _op_not_implemented_message(self.store_type, "get_connector_metadata") 102 ) 103 104 def list_yanked_versions( 105 self, 106 *, 107 with_details: bool = True, 108 ) -> list[YankedVersion]: 109 raise NotImplementedError( 110 _op_not_implemented_message(self.store_type, "list_yanked_versions") 111 ) 112 113 def get_yank_marker( 114 self, 115 connector_name: str, 116 version: str, 117 ) -> YankMarkerDetail | None: 118 raise NotImplementedError( 119 _op_not_implemented_message(self.store_type, "get_yank_marker") 120 ) 121 122 # --------------------------------------------------------------------- 123 # Write / mutate operations 124 # --------------------------------------------------------------------- 125 126 def yank( 127 self, 128 connector_name: str, 129 version: str, 130 reason: str = "", 131 approval_url: str = "", 132 dry_run: bool = False, 133 ) -> YankResult: 134 raise NotImplementedError(_op_not_implemented_message(self.store_type, "yank")) 135 136 def unyank( 137 self, 138 connector_name: str, 139 version: str, 140 dry_run: bool = False, 141 ) -> YankResult: 142 raise NotImplementedError( 143 _op_not_implemented_message(self.store_type, "unyank") 144 ) 145 146 def finalize_progressive_rollout_marker( 147 self, 148 connector_name: str, 149 outcome: Literal["promoted", "aborted"], 150 version: str | None = None, 151 dry_run: bool = False, 152 ) -> ProgressiveRolloutMarkerResult: 153 raise NotImplementedError( 154 _op_not_implemented_message( 155 self.store_type, 156 "finalize_progressive_rollout_marker", 157 ) 158 ) 159 160 def publish_version_artifacts( 161 self, 162 connector_name: str, 163 version: str, 164 artifacts_dir: Path, 165 dry_run: bool = False, 166 with_validate: bool = True, 167 ) -> PublishArtifactsResult: 168 raise NotImplementedError( 169 _op_not_implemented_message(self.store_type, "publish_version_artifacts") 170 ) 171 172 def delete_dev_latest( 173 self, 174 connector_name: list[str] | None = None, 175 dry_run: bool = False, 176 ) -> PurgeLatestResult: 177 raise NotImplementedError( 178 _op_not_implemented_message(self.store_type, "delete_dev_latest") 179 ) 180 181 def compile( 182 self, 183 output_store: RegistryStore | None = None, 184 connector_name: list[str] | None = None, 185 dry_run: bool = False, 186 with_secrets_mask: bool = False, 187 with_legacy_migration: str | None = None, 188 with_metrics: bool = True, 189 force: bool = False, 190 with_full_restate: bool = False, 191 release_attribution_index: Path | None = None, 192 ) -> CompileResult: 193 raise NotImplementedError( 194 _op_not_implemented_message(self.store_type, "compile") 195 ) 196 197 def marketing_stubs_check(self, repo_root: Path) -> dict[str, Any]: 198 raise NotImplementedError( 199 _op_not_implemented_message(self.store_type, "marketing_stubs_check") 200 ) 201 202 def marketing_stubs_sync( 203 self, 204 repo_root: Path, 205 dry_run: bool = False, 206 ) -> dict[str, Any]: 207 raise NotImplementedError( 208 _op_not_implemented_message(self.store_type, "marketing_stubs_sync") 209 ) 210 211 def mirror( 212 self, 213 output_mode: OutputMode, 214 output_path_root: str | None = None, 215 gcs_bucket: str | None = None, 216 s3_bucket: str | None = None, 217 dry_run: bool = False, 218 connector_name: list[str] | None = None, 219 ) -> RebuildResult: 220 raise NotImplementedError( 221 _op_not_implemented_message(self.store_type, "mirror") 222 )
A configured connector registry store.
A Registry is bound to a specific airbyte_ops_mcp.registry.store.RegistryStore
(store type + env + optional prefix), and provides methods used by the CLI to
read/write registry contents.
79 @abstractmethod 80 def list_connectors( 81 self, 82 *, 83 support_level: SupportLevel | None = None, 84 min_support_level: SupportLevel | None = None, 85 connector_type: ConnectorType | None = None, 86 language: ConnectorLanguage | None = None, 87 ) -> list[str]: 88 raise NotImplementedError( 89 _op_not_implemented_message(self.store_type, "list_connectors") 90 )
146 def finalize_progressive_rollout_marker( 147 self, 148 connector_name: str, 149 outcome: Literal["promoted", "aborted"], 150 version: str | None = None, 151 dry_run: bool = False, 152 ) -> ProgressiveRolloutMarkerResult: 153 raise NotImplementedError( 154 _op_not_implemented_message( 155 self.store_type, 156 "finalize_progressive_rollout_marker", 157 ) 158 )
160 def publish_version_artifacts( 161 self, 162 connector_name: str, 163 version: str, 164 artifacts_dir: Path, 165 dry_run: bool = False, 166 with_validate: bool = True, 167 ) -> PublishArtifactsResult: 168 raise NotImplementedError( 169 _op_not_implemented_message(self.store_type, "publish_version_artifacts") 170 )
181 def compile( 182 self, 183 output_store: RegistryStore | None = None, 184 connector_name: list[str] | None = None, 185 dry_run: bool = False, 186 with_secrets_mask: bool = False, 187 with_legacy_migration: str | None = None, 188 with_metrics: bool = True, 189 force: bool = False, 190 with_full_restate: bool = False, 191 release_attribution_index: Path | None = None, 192 ) -> CompileResult: 193 raise NotImplementedError( 194 _op_not_implemented_message(self.store_type, "compile") 195 )
211 def mirror( 212 self, 213 output_mode: OutputMode, 214 output_path_root: str | None = None, 215 gcs_bucket: str | None = None, 216 s3_bucket: str | None = None, 217 dry_run: bool = False, 218 connector_name: list[str] | None = None, 219 ) -> RebuildResult: 220 raise NotImplementedError( 221 _op_not_implemented_message(self.store_type, "mirror") 222 )
60class RegistryEntryResult(BaseModel): 61 """Result of reading a registry entry from GCS. 62 63 This model wraps the raw metadata dictionary with additional context. 64 """ 65 66 connector_name: str = Field(description="The connector technical name") 67 version: str = Field(description="The version that was read") 68 bucket_name: str = Field(description="The GCS bucket name") 69 gcs_path: str = Field(description="The GCS path that was read") 70 metadata: dict = Field(description="The raw metadata dictionary")
Result of reading a registry entry from GCS.
This model wraps the raw metadata dictionary with additional context.
134@dataclass(frozen=True, kw_only=True) 135class RegistryStore: 136 """Parsed store target (type + environment + optional prefix). 137 138 `read_only` is an internal invariant used to ensure a source store 139 cannot yield a writable filesystem client. 140 141 Examples:: 142 143 RegistryStore.parse("coral:dev") 144 # -> RegistryStore(store_type=StoreType.CORAL, env="dev", prefix="") 145 RegistryStore.parse("coral:dev/aj-test") 146 # -> RegistryStore(store_type=StoreType.CORAL, env="dev", prefix="aj-test") 147 RegistryStore.parse("sonar:prod") 148 # -> RegistryStore(store_type=StoreType.SONAR, env="prod", prefix="") 149 """ 150 151 store_type: StoreType 152 env: str 153 prefix: str = "" 154 read_only: bool = False 155 local_path: Path | None = None 156 157 def __post_init__(self) -> None: 158 """Validate local target fields and allocate omitted paths.""" 159 if self.env == "local": 160 if self.prefix: 161 raise ValueError("Local store targets cannot have a prefix.") 162 if self.local_path is None: 163 object.__setattr__( 164 self, 165 "local_path", 166 Path(tempfile.mkdtemp(prefix="airbyte-registry-")), 167 ) 168 elif self.local_path is not None: 169 raise ValueError("Only local store targets may have a local path.") 170 171 # -- Derived helpers ----------------------------------------------------- 172 173 @property 174 def bucket(self) -> str: 175 """Resolve the concrete bucket name for this target.""" 176 if self.env == "local": 177 raise ValueError("Local stores do not have a bucket.") 178 env_map = BUCKET_MAP.get(self.store_type) 179 if env_map is None: 180 raise ValueError(f"Unknown store type: {self.store_type!r}") 181 bucket_name = env_map.get(self.env) 182 if bucket_name is None: 183 raise ValueError( 184 f"Unknown environment '{self.env}' for store type '{self.store_type.value}'. " 185 f"Expected one of: {', '.join(sorted(env_map))}." 186 ) 187 return bucket_name 188 189 @property 190 def bucket_root(self) -> str: 191 """Bucket or local root with optional prefix appended.""" 192 if self.env == "local": 193 assert self.local_path is not None 194 return str(self.local_path) 195 if self.prefix: 196 return f"{self.bucket}/{self.prefix}" 197 return self.bucket 198 199 # -- Parsing ------------------------------------------------------------- 200 201 @classmethod 202 def parse(cls, target: str) -> RegistryStore: 203 """Parse a store target string. 204 205 Accepted formats: 206 207 "coral:dev" 208 209 "coral:prod" 210 "coral:dev/aj-test100" 211 "sonar:prod" 212 "coral:local:/tmp/registry-output" 213 214 Raises: 215 ValueError: If the string cannot be parsed or references an 216 unknown store type / environment. 217 """ 218 if ":" not in target: 219 raise ValueError( 220 f"Invalid store target '{target}'. " 221 "Expected format: '<store_type>:<env>[/<prefix>]' " 222 "(e.g. 'coral:dev', 'sonar:prod', 'coral:dev/my-test')." 223 ) 224 225 store_part, env_part = target.split(":", 1) 226 227 # Validate store type 228 store_part_lower = store_part.lower() 229 valid_types = {t.value: t for t in StoreType} 230 if store_part_lower not in valid_types: 231 raise ValueError( 232 f"Unknown store type '{store_part}'. " 233 f"Expected one of: {', '.join(sorted(valid_types))}." 234 ) 235 store_type = valid_types[store_part_lower] 236 237 # Split env and prefix 238 if env_part.startswith("local:"): 239 local_path_text = env_part.removeprefix("local:") 240 return cls( 241 store_type=store_type, 242 env="local", 243 local_path=Path(local_path_text) if local_path_text else None, 244 ) 245 env_key, _, prefix = env_part.partition("/") 246 prefix = prefix.strip("/") 247 248 # Validate env 249 env_map = BUCKET_MAP.get(store_type, {}) 250 if env_key not in env_map: 251 raise ValueError( 252 f"Unknown environment '{env_key}' for store type '{store_type.value}'. " 253 f"Expected one of: {', '.join(sorted(env_map))}." 254 ) 255 256 return cls(store_type=store_type, env=env_key, prefix=prefix)
Parsed store target (type + environment + optional prefix).
read_only is an internal invariant used to ensure a source store
cannot yield a writable filesystem client.
Examples::
RegistryStore.parse("coral:dev")
# -> RegistryStore(store_type=StoreType.CORAL, env="dev", prefix="")
RegistryStore.parse("coral:dev/aj-test")
# -> RegistryStore(store_type=StoreType.CORAL, env="dev", prefix="aj-test")
RegistryStore.parse("sonar:prod")
# -> RegistryStore(store_type=StoreType.SONAR, env="prod", prefix="")
173 @property 174 def bucket(self) -> str: 175 """Resolve the concrete bucket name for this target.""" 176 if self.env == "local": 177 raise ValueError("Local stores do not have a bucket.") 178 env_map = BUCKET_MAP.get(self.store_type) 179 if env_map is None: 180 raise ValueError(f"Unknown store type: {self.store_type!r}") 181 bucket_name = env_map.get(self.env) 182 if bucket_name is None: 183 raise ValueError( 184 f"Unknown environment '{self.env}' for store type '{self.store_type.value}'. " 185 f"Expected one of: {', '.join(sorted(env_map))}." 186 ) 187 return bucket_name
Resolve the concrete bucket name for this target.
189 @property 190 def bucket_root(self) -> str: 191 """Bucket or local root with optional prefix appended.""" 192 if self.env == "local": 193 assert self.local_path is not None 194 return str(self.local_path) 195 if self.prefix: 196 return f"{self.bucket}/{self.prefix}" 197 return self.bucket
Bucket or local root with optional prefix appended.
201 @classmethod 202 def parse(cls, target: str) -> RegistryStore: 203 """Parse a store target string. 204 205 Accepted formats: 206 207 "coral:dev" 208 209 "coral:prod" 210 "coral:dev/aj-test100" 211 "sonar:prod" 212 "coral:local:/tmp/registry-output" 213 214 Raises: 215 ValueError: If the string cannot be parsed or references an 216 unknown store type / environment. 217 """ 218 if ":" not in target: 219 raise ValueError( 220 f"Invalid store target '{target}'. " 221 "Expected format: '<store_type>:<env>[/<prefix>]' " 222 "(e.g. 'coral:dev', 'sonar:prod', 'coral:dev/my-test')." 223 ) 224 225 store_part, env_part = target.split(":", 1) 226 227 # Validate store type 228 store_part_lower = store_part.lower() 229 valid_types = {t.value: t for t in StoreType} 230 if store_part_lower not in valid_types: 231 raise ValueError( 232 f"Unknown store type '{store_part}'. " 233 f"Expected one of: {', '.join(sorted(valid_types))}." 234 ) 235 store_type = valid_types[store_part_lower] 236 237 # Split env and prefix 238 if env_part.startswith("local:"): 239 local_path_text = env_part.removeprefix("local:") 240 return cls( 241 store_type=store_type, 242 env="local", 243 local_path=Path(local_path_text) if local_path_text else None, 244 ) 245 env_key, _, prefix = env_part.partition("/") 246 prefix = prefix.strip("/") 247 248 # Validate env 249 env_map = BUCKET_MAP.get(store_type, {}) 250 if env_key not in env_map: 251 raise ValueError( 252 f"Unknown environment '{env_key}' for store type '{store_type.value}'. " 253 f"Expected one of: {', '.join(sorted(env_map))}." 254 ) 255 256 return cls(store_type=store_type, env=env_key, prefix=prefix)
Parse a store target string.
Accepted formats:
"coral:dev"
"coral:prod" "coral:dev/aj-test100" "sonar:prod" "coral:local:/tmp/registry-output"
Raises:
- ValueError: If the string cannot be parsed or references an unknown store type / environment.
61class StoreType(str, Enum): 62 """Registry store type identifier.""" 63 64 SONAR = "sonar" 65 CORAL = "coral" 66 67 # -- Auto-detection class methods ---------------------------------------- 68 69 @classmethod 70 def get_from_connector_name(cls, name: str) -> StoreType: 71 """Infer the store type from a connector's technical name. 72 73 Connectors whose name starts with `source-` or `destination-` belong 74 to the **coral** registry. All other names belong to **sonar**. 75 76 Args: 77 name: Connector technical name (e.g. `"source-github"` or `"stripe"`). 78 79 Returns: 80 The inferred `StoreType`. 81 """ 82 if name.startswith("source-") or name.startswith("destination-"): 83 return cls.CORAL 84 return cls.SONAR 85 86 @classmethod 87 def detect_from_repo_dir(cls, path: Path | None = None) -> StoreType | None: 88 """Infer the store type from a repository working directory. 89 90 Checks for well-known directory markers: 91 92 * **sonar** -- `integrations/` alongside `connector-sdk/` 93 * **coral** -- `airbyte-integrations/connectors/` 94 95 Args: 96 path: Directory to inspect. Defaults to `Path.cwd`. 97 98 Returns: 99 The inferred `StoreType`, or `None` if the directory does 100 not match any known registry repository layout. 101 """ 102 if path is None: 103 path = Path.cwd() 104 105 # Sonar markers 106 if (path / "integrations").is_dir() and (path / "connector-sdk").is_dir(): 107 return cls.SONAR 108 109 # Coral markers 110 if (path / "airbyte-integrations" / "connectors").is_dir(): 111 return cls.CORAL 112 113 return None
Registry store type identifier.
69 @classmethod 70 def get_from_connector_name(cls, name: str) -> StoreType: 71 """Infer the store type from a connector's technical name. 72 73 Connectors whose name starts with `source-` or `destination-` belong 74 to the **coral** registry. All other names belong to **sonar**. 75 76 Args: 77 name: Connector technical name (e.g. `"source-github"` or `"stripe"`). 78 79 Returns: 80 The inferred `StoreType`. 81 """ 82 if name.startswith("source-") or name.startswith("destination-"): 83 return cls.CORAL 84 return cls.SONAR
Infer the store type from a connector's technical name.
Connectors whose name starts with source- or destination- belong
to the coral registry. All other names belong to sonar.
Arguments:
- name: Connector technical name (e.g.
"source-github"or"stripe").
Returns:
The inferred
StoreType.
86 @classmethod 87 def detect_from_repo_dir(cls, path: Path | None = None) -> StoreType | None: 88 """Infer the store type from a repository working directory. 89 90 Checks for well-known directory markers: 91 92 * **sonar** -- `integrations/` alongside `connector-sdk/` 93 * **coral** -- `airbyte-integrations/connectors/` 94 95 Args: 96 path: Directory to inspect. Defaults to `Path.cwd`. 97 98 Returns: 99 The inferred `StoreType`, or `None` if the directory does 100 not match any known registry repository layout. 101 """ 102 if path is None: 103 path = Path.cwd() 104 105 # Sonar markers 106 if (path / "integrations").is_dir() and (path / "connector-sdk").is_dir(): 107 return cls.SONAR 108 109 # Coral markers 110 if (path / "airbyte-integrations" / "connectors").is_dir(): 111 return cls.CORAL 112 113 return None
Infer the store type from a repository working directory.
Checks for well-known directory markers:
- sonar --
integrations/alongsideconnector-sdk/ - coral --
airbyte-integrations/connectors/
Arguments:
- path: Directory to inspect. Defaults to
Path.cwd.
Returns:
The inferred
StoreType, orNoneif the directory does not match any known registry repository layout.
14class SupportLevel(StrEnum): 15 """Connector support levels ordered by precedence.""" 16 17 ARCHIVED = "archived" 18 COMMUNITY = "community" 19 CERTIFIED = "certified" 20 21 @property 22 def precedence(self) -> int: 23 """Numeric precedence for ordering comparisons. 24 25 Higher values indicate higher support commitment. 26 """ 27 return _SUPPORT_LEVEL_PRECEDENCE[self] 28 29 @classmethod 30 def from_precedence(cls, precedence: int) -> SupportLevel: 31 """Look up a `SupportLevel` by its numeric precedence value. 32 33 Raises `ValueError` when the precedence is not recognised. 34 """ 35 for member in cls: 36 if _SUPPORT_LEVEL_PRECEDENCE[member] == precedence: 37 return member 38 valid = ", ".join(f"`{_SUPPORT_LEVEL_PRECEDENCE[m]}`" for m in cls) 39 raise ValueError( 40 f"Unrecognized support-level precedence: {precedence!r}. " 41 f"Expected one of: {valid}." 42 ) from None 43 44 @classmethod 45 def parse(cls, value: str) -> SupportLevel: 46 """Parse a string into a `SupportLevel`. 47 48 Accepts a keyword (`archived`, `community`, `certified`) 49 or a legacy integer string (`100`, `200`, `300`). 50 51 Raises `ValueError` when the value is not recognised. 52 """ 53 try: 54 return cls(value) 55 except ValueError: 56 pass 57 # Fallback: try interpreting as an integer precedence value. 58 try: 59 return cls.from_precedence(int(value)) 60 except (ValueError, KeyError): 61 pass 62 valid_kw = ", ".join(f"`{m.value}`" for m in cls) 63 valid_int = ", ".join(f"`{_SUPPORT_LEVEL_PRECEDENCE[m]}`" for m in cls) 64 raise ValueError( 65 f"Unrecognized support level: {value!r}. " 66 f"Expected keyword ({valid_kw}) or integer ({valid_int})." 67 ) from None
Connector support levels ordered by precedence.
21 @property 22 def precedence(self) -> int: 23 """Numeric precedence for ordering comparisons. 24 25 Higher values indicate higher support commitment. 26 """ 27 return _SUPPORT_LEVEL_PRECEDENCE[self]
Numeric precedence for ordering comparisons.
Higher values indicate higher support commitment.
29 @classmethod 30 def from_precedence(cls, precedence: int) -> SupportLevel: 31 """Look up a `SupportLevel` by its numeric precedence value. 32 33 Raises `ValueError` when the precedence is not recognised. 34 """ 35 for member in cls: 36 if _SUPPORT_LEVEL_PRECEDENCE[member] == precedence: 37 return member 38 valid = ", ".join(f"`{_SUPPORT_LEVEL_PRECEDENCE[m]}`" for m in cls) 39 raise ValueError( 40 f"Unrecognized support-level precedence: {precedence!r}. " 41 f"Expected one of: {valid}." 42 ) from None
Look up a SupportLevel by its numeric precedence value.
Raises ValueError when the precedence is not recognised.
44 @classmethod 45 def parse(cls, value: str) -> SupportLevel: 46 """Parse a string into a `SupportLevel`. 47 48 Accepts a keyword (`archived`, `community`, `certified`) 49 or a legacy integer string (`100`, `200`, `300`). 50 51 Raises `ValueError` when the value is not recognised. 52 """ 53 try: 54 return cls(value) 55 except ValueError: 56 pass 57 # Fallback: try interpreting as an integer precedence value. 58 try: 59 return cls.from_precedence(int(value)) 60 except (ValueError, KeyError): 61 pass 62 valid_kw = ", ".join(f"`{m.value}`" for m in cls) 63 valid_int = ", ".join(f"`{_SUPPORT_LEVEL_PRECEDENCE[m]}`" for m in cls) 64 raise ValueError( 65 f"Unrecognized support level: {value!r}. " 66 f"Expected keyword ({valid_kw}) or integer ({valid_int})." 67 ) from None
Parse a string into a SupportLevel.
Accepts a keyword (archived, community, certified)
or a legacy integer string (100, 200, 300).
Raises ValueError when the value is not recognised.
29@dataclass 30class UnpublishedConnector: 31 """A connector whose current local version is not published on GCS.""" 32 33 connector_name: str 34 local_version: str
A connector whose current local version is not published on GCS.
37@dataclass(frozen=True) 38class ValidateOptions: 39 """Options that influence which validators run and how.""" 40 41 docs_path: str | None = None 42 """Path to the connector's documentation file (for `validate_docs_path_exists`).""" 43 44 is_prerelease: bool = False 45 """Whether this is a pre-release build (skips version-decrement checks)."""
Options that influence which validators run and how.
48@dataclass 49class ValidationResult: 50 """Aggregate result from running all validators.""" 51 52 passed: bool = True 53 errors: list[str] = field(default_factory=list) 54 validators_run: int = 0 55 56 def add_error(self, message: str) -> None: 57 self.passed = False 58 self.errors.append(message)
Aggregate result from running all validators.
81class VersionListResult(BaseModel): 82 """Result of listing versions for a connector.""" 83 84 connector_name: str = Field(description="The connector technical name") 85 bucket_name: str = Field(description="The GCS bucket name") 86 version_count: int = Field(description="Number of versions found") 87 versions: list[str] = Field(description="List of version strings")
Result of listing versions for a connector.
34@dataclass 35class YankResult: 36 """Result of a yank or unyank operation.""" 37 38 connector_name: str 39 version: str 40 bucket_name: str 41 action: str # "yank" or "unyank" 42 success: bool 43 message: str 44 dry_run: bool = False 45 46 def to_dict(self) -> dict[str, Any]: 47 """Convert to a dictionary for JSON serialization.""" 48 return { 49 "connector_name": self.connector_name, 50 "version": self.version, 51 "bucket_name": self.bucket_name, 52 "action": self.action, 53 "success": self.success, 54 "message": self.message, 55 "dry_run": self.dry_run, 56 }
Result of a yank or unyank operation.
46 def to_dict(self) -> dict[str, Any]: 47 """Convert to a dictionary for JSON serialization.""" 48 return { 49 "connector_name": self.connector_name, 50 "version": self.version, 51 "bucket_name": self.bucket_name, 52 "action": self.action, 53 "success": self.success, 54 "message": self.message, 55 "dry_run": self.dry_run, 56 }
Convert to a dictionary for JSON serialization.
1862def compile_registry( 1863 *, 1864 store: RegistryStore, 1865 output_store: RegistryStore | None = None, 1866 connector_name: list[str] | None = None, 1867 dry_run: bool = False, 1868 with_secrets_mask: bool = False, 1869 with_legacy_migration: str | None = None, 1870 with_metrics: bool = True, 1871 force: bool = False, 1872 with_full_restate: bool = False, 1873 release_attribution_index: ReleaseAttributionIndex | None = None, 1874) -> CompileResult: 1875 """Compile the registry: sync latest/ dirs and write index files. 1876 1877 Steps: 1878 1. Glob for all `metadata.yaml` to discover (connector, version) pairs. 1879 2. Glob for active marker files. 1880 3. Compute the latest GA semver per connector. 1881 4. Compute active release candidates from versioned markers. 1882 5. Glob for `version=*` markers in `latest/` dirs for a fast check. 1883 6. Delete stale `latest/` dirs and recursively copy the versioned dir. 1884 7. Synthesize missing registry entries from pinned latest overrides. 1885 8. (Optional) Legacy migration: delete disabled registry entries. 1886 9. (Optional) Read latest connector metrics. 1887 10. Write global registry JSONs. When an output store is supplied, 1888 read each connector entry from its computed version directory 1889 instead of `latest/`. 1890 11. Write composite registry JSON. 1891 12. Write per-connector `versions.json`. 1892 13. (Optional) Regenerate `specs_secrets_mask.yaml`. 1893 1894 Args: 1895 store: Registry store (bucket + optional prefix). 1896 connector_name: If provided, only resync `latest/` directories for 1897 these connectors (steps 5-6). Global registry indexes always 1898 operate on the full set so they remain complete; per-connector 1899 `versions.json` writes are scoped to the requested connectors. 1900 dry_run: If True, report what would be done without writing. 1901 with_secrets_mask: If True, regenerate `specs_secrets_mask.yaml`. 1902 with_legacy_migration: If set, run the named migration step. 1903 Currently supported: `"v1"` — delete `{registry_type}.json` 1904 files for connectors whose `registryOverrides.{registry}.enabled` 1905 is `false`. 1906 with_metrics: If True, inject latest connector metrics from the 1907 analytics JSONL export into `generated.metrics`. 1908 force: If True, resync all connectors' latest/ directories even if the 1909 existing version marker matches the computed latest version. This 1910 is useful when metadata content changes without a version bump. 1911 with_full_restate: If `True`, re-read every version's `metadata.yaml` 1912 instead of carrying forward existing release attribution. This is 1913 expensive for all connectors, at approximately 33.5k reads. When 1914 no `release_attribution_index` is supplied, historical release 1915 blocks that exist only in the attribution index cannot be 1916 re-derived and will be discarded. 1917 release_attribution_index: Optional Phase 1 index used to seed missing 1918 per-version release blocks without rewriting historical metadata. 1919 1920 The compile performs one additional GCS read per connector for its prior 1921 `versions.json`, trading approximately 716 reads for avoiding historical 1922 metadata reads on the default path. 1923 1924 Returns: 1925 A `CompileResult` describing what was done. 1926 """ 1927 if with_legacy_migration and with_legacy_migration not in LEGACY_MIGRATION_VERSIONS: 1928 raise ValueError( 1929 f"Unknown legacy migration version: {with_legacy_migration!r}. " 1930 f"Supported: {', '.join(LEGACY_MIGRATION_VERSIONS)}" 1931 ) 1932 if output_store is not None and with_legacy_migration is not None: 1933 raise ValueError( 1934 "Legacy migration is incompatible with --output-store because " 1935 "the source store must remain read-only." 1936 ) 1937 if output_store is not None and output_store.env == "prod": 1938 raise ValueError("Production stores cannot be compile output targets.") 1939 if output_store is not None and output_store.env != "local": 1940 raise ValueError("--output-store must use a local environment.") 1941 if output_store is not None and output_store.store_type is not store.store_type: 1942 raise ValueError( 1943 "Source and output stores must use the same registry type: " 1944 f"{store.store_type.value!r} versus {output_store.store_type.value!r}." 1945 ) 1946 1947 output_target = output_store or store 1948 result = CompileResult(target=output_target.bucket_root, dry_run=dry_run) 1949 1950 source_store = replace(store, read_only=True) if output_store is not None else store 1951 fs = _filesystem_for_store(source_store) 1952 output_fs = _filesystem_for_store(output_target) 1953 1954 # --- Steps 1 and 2: Scan versions and active markers --- 1955 # Always scan ALL connectors so that index rebuilds are complete. 1956 _log_progress("Step 1-2: Scanning versions and active markers...") 1957 connector_versions, yanked, progressive_rollouts = _scan_versions_and_markers( 1958 fs, 1959 store=store, 1960 connector_name=None, 1961 ) 1962 result.connectors_scanned = len(connector_versions) 1963 result.versions_found = sum(len(v) for v in connector_versions.values()) 1964 result.yanked_versions = len(yanked) 1965 _log_progress( 1966 " Found %d connectors, %d versions, %d yanked", 1967 result.connectors_scanned, 1968 result.versions_found, 1969 result.yanked_versions, 1970 ) 1971 _log_progress(" Found %d progressive rollout markers", len(progressive_rollouts)) 1972 1973 # --- Step 3: Compute latest --- 1974 _log_progress("Step 3: Computing latest GA version per connector...") 1975 latest_versions = _compute_latest_versions( 1976 connector_versions=connector_versions, 1977 yanked=yanked, 1978 progressive_rollouts=progressive_rollouts, 1979 ) 1980 _log_progress(" Computed latest for %d connectors", len(latest_versions)) 1981 1982 # --- Step 4: Compute release candidates --- 1983 _log_progress("Step 4: Computing active release candidates...") 1984 rc_versions = _compute_release_candidates( 1985 connector_versions=connector_versions, 1986 yanked=yanked, 1987 progressive_rollouts=progressive_rollouts, 1988 ) 1989 _log_progress(" Computed %d active release candidates", len(rc_versions)) 1990 1991 # --- Step 5: Check existing latest markers --- 1992 # When --connector-name is set, only check/sync those connectors (steps 5-6). 1993 # Index rebuilds always use the full unfiltered data. 1994 if output_store is not None: 1995 _log_progress("Output store supplied; source store will remain read-only.") 1996 sync_scope = {} 1997 elif connector_name: 1998 connector_name_set = set(connector_name) 1999 sync_scope = { 2000 c: v for c, v in latest_versions.items() if c in connector_name_set 2001 } 2002 _log_progress( 2003 " --connector-name filter: syncing %d of %d connectors", 2004 len(sync_scope), 2005 len(latest_versions), 2006 ) 2007 else: 2008 sync_scope = latest_versions 2009 2010 _log_progress("Step 5: Checking existing latest/ markers...") 2011 sync_scope_names = list(sync_scope) if connector_name else None 2012 existing_markers = ( 2013 {} 2014 if output_store is not None 2015 else _scan_latest_markers( 2016 fs, 2017 store=store, 2018 connector_name=sync_scope_names, 2019 ) 2020 ) 2021 _log_progress(" Found %d existing markers", len(existing_markers)) 2022 2023 stale_connectors: list[str] = [] 2024 pinned_override_synthesis_connectors: list[str] = [] 2025 for connector, expected_version in sync_scope.items(): 2026 current_marker = existing_markers.get(connector) 2027 if force or current_marker != expected_version: 2028 stale_connectors.append(connector) 2029 continue 2030 2031 if _requires_pinned_override_synthesis( 2032 fs, 2033 store=store, 2034 connector=connector, 2035 version=expected_version, 2036 ): 2037 pinned_override_synthesis_connectors.append(connector) 2038 continue 2039 2040 result.latest_already_current += 1 2041 2042 _log_progress( 2043 " %d connectors need latest/ update, %d already current", 2044 len(stale_connectors), 2045 result.latest_already_current, 2046 ) 2047 2048 # --- Step 6: Resync stale latest/ dirs (parallel) --- 2049 if stale_connectors: 2050 _log_progress( 2051 "Step 6: Syncing %d stale latest/ directories (max_workers=%d)...", 2052 len(stale_connectors), 2053 _COMPILE_SYNC_MAX_WORKERS, 2054 ) 2055 2056 def _sync_one_connector(connector: str) -> None: 2057 """Sync a single connector's latest/ dir.""" 2058 version = latest_versions[connector] 2059 _sync_latest_dir( 2060 fs, 2061 store=store, 2062 connector=connector, 2063 version=version, 2064 dry_run=dry_run, 2065 ) 2066 if not dry_run: 2067 _apply_overrides_to_latest_entry( 2068 fs, 2069 store=store, 2070 connector=connector, 2071 version=version, 2072 ) 2073 2074 sorted_stale = sorted(stale_connectors) 2075 with ThreadPoolExecutor(max_workers=_COMPILE_SYNC_MAX_WORKERS) as pool: 2076 futures = {pool.submit(_sync_one_connector, c): c for c in sorted_stale} 2077 for i, future in enumerate(as_completed(futures), 1): 2078 connector = futures[future] 2079 try: 2080 future.result() 2081 result.latest_updated += 1 2082 except Exception as exc: 2083 error_msg = f"Failed to sync latest/ for {connector}: {exc}" 2084 logger.error(error_msg) 2085 result.errors.append(error_msg) 2086 # Delete the (possibly partial) latest/ dir so the next 2087 # compile retries this connector from scratch. 2088 try: 2089 _delete_latest_dir( 2090 fs, 2091 store=store, 2092 connector=connector, 2093 ) 2094 logger.info( 2095 "Cleaned up partial latest/ for %s after failure", 2096 connector, 2097 ) 2098 except Exception as cleanup_exc: 2099 logger.warning( 2100 "Could not clean up latest/ for %s: %s", 2101 connector, 2102 cleanup_exc, 2103 ) 2104 if i % 100 == 0: 2105 _log_progress(" Synced %d / %d...", i, len(sorted_stale)) 2106 else: 2107 _log_progress("Step 6: All latest/ directories are current, nothing to sync.") 2108 2109 # --- Step 7: Synthesize missing latest entries from pinned overrides --- 2110 if pinned_override_synthesis_connectors: 2111 _log_progress( 2112 "Step 7: Synthesizing %d latest/ registry entries from pinned overrides...", 2113 len(pinned_override_synthesis_connectors), 2114 ) 2115 for connector in sorted(pinned_override_synthesis_connectors): 2116 if dry_run: 2117 _log_progress( 2118 " [DRY RUN] Would synthesize pinned latest entries for %s", 2119 connector, 2120 ) 2121 result.latest_updated += 1 2122 continue 2123 try: 2124 _apply_overrides_to_latest_entry( 2125 fs, 2126 store=store, 2127 connector=connector, 2128 version=sync_scope[connector], 2129 ) 2130 result.latest_updated += 1 2131 except Exception as exc: 2132 error_msg = ( 2133 f"Failed to synthesize pinned latest entries for {connector}: {exc}" 2134 ) 2135 logger.error(error_msg) 2136 result.errors.append(error_msg) 2137 2138 # --- Step 8: Legacy migration (optional) --- 2139 if with_legacy_migration == "v1": 2140 _log_progress( 2141 "Step 8: Legacy migration v1 — deleting disabled registry entries..." 2142 ) 2143 migration_deleted = _cleanup_disabled_registry_entries( 2144 fs, 2145 store=store, 2146 connector_versions=connector_versions, 2147 dry_run=dry_run, 2148 ) 2149 total_deleted = sum(len(v) for v in migration_deleted.values()) 2150 if migration_deleted: 2151 for conn, paths in sorted(migration_deleted.items()): 2152 _log_progress( 2153 " %s: %s %d files", 2154 conn, 2155 "would delete" if dry_run else "deleted", 2156 len(paths), 2157 ) 2158 _log_progress( 2159 " Migration v1: %s %d files across %d connectors", 2160 "would delete" if dry_run else "deleted", 2161 total_deleted, 2162 len(migration_deleted), 2163 ) 2164 2165 # --- Step 9: Read latest connector metrics (optional) --- 2166 metrics_bundle = None 2167 if with_metrics and store.store_type == StoreType.CORAL: 2168 _log_progress("Step 9: Reading latest connector metrics JSONL...") 2169 try: 2170 metrics_bundle = read_latest_connector_metrics() 2171 result.metrics_source = metrics_bundle.blob_path 2172 result.metrics_connector_count = metrics_bundle.connector_count 2173 if metrics_bundle.blob_path: 2174 _log_progress( 2175 " Loaded metrics for %d connectors from gs://%s", 2176 metrics_bundle.connector_count, 2177 metrics_bundle.blob_path, 2178 ) 2179 else: 2180 _log_progress(" No connector metrics JSONL file found.") 2181 except Exception as exc: 2182 error_msg = f"Failed to read connector metrics JSONL: {exc}" 2183 logger.warning(error_msg) 2184 result.metrics_error = error_msg 2185 _log_progress(" %s", error_msg) 2186 elif with_metrics: 2187 _log_progress("Step 9: Skipping connector metrics for non-coral registry.") 2188 else: 2189 _log_progress("Step 9: Connector metrics injection disabled.") 2190 2191 # --- Step 10: Compile global registry JSONs --- 2192 _log_progress("Step 10: Compiling global registry JSON files...") 2193 all_registry_entries: list[dict[str, Any]] = [] # collected for Step 13 2194 entries_by_registry_type: dict[str, list[dict[str, Any]]] = {} 2195 for registry_type in VALID_REGISTRIES: 2196 entries = _compile_global_registry( 2197 fs, 2198 store=store, 2199 latest_versions=latest_versions, 2200 registry_type=registry_type, 2201 source_from_version_dirs=output_store is not None, 2202 ) 2203 2204 # Inject release candidate info into entries that have active RCs. 2205 if rc_versions: 2206 rc_entries: dict[str, list[dict[str, Any]]] = {} 2207 for connector, rc_ver_list in rc_versions.items(): 2208 for rc_ver in rc_ver_list: 2209 rc_entry = _read_rc_registry_entry( 2210 fs, 2211 store=store, 2212 connector=connector, 2213 rc_version=rc_ver, 2214 registry_type=registry_type, 2215 ) 2216 if rc_entry: 2217 docker_repo = rc_entry.get( 2218 "dockerRepository", 2219 f"airbyte/{connector}", 2220 ) 2221 rc_entries.setdefault(docker_repo, []).append( 2222 { 2223 "version": rc_ver, 2224 "entry": rc_entry, 2225 } 2226 ) 2227 if rc_entries: 2228 entries = _apply_release_candidates_to_entries(entries, rc_entries) 2229 total_rcs = sum(len(v) for v in rc_entries.values()) 2230 _log_progress( 2231 " Injected %d release candidate(s) for %d connector(s) into %s registry", 2232 total_rcs, 2233 len(rc_entries), 2234 registry_type, 2235 ) 2236 2237 if metrics_bundle is not None: 2238 injected = apply_metrics_to_registry_entries(entries, metrics_bundle) 2239 result.metrics_registry_entries += injected 2240 _log_progress( 2241 " Injected metrics into %d %s registry entries", 2242 injected, 2243 registry_type, 2244 ) 2245 2246 all_registry_entries.extend(entries) 2247 entries_by_registry_type[registry_type] = entries 2248 registry_json = _build_global_registry_json(entries) 2249 entry_count = len(registry_json["sources"]) + len(registry_json["destinations"]) 2250 2251 if registry_type == "cloud": 2252 result.cloud_registry_entries = entry_count 2253 else: 2254 result.oss_registry_entries = entry_count 2255 2256 if dry_run: 2257 _log_progress( 2258 " [DRY RUN] Would write %s_registry.json (%d entries)", 2259 registry_type, 2260 entry_count, 2261 ) 2262 else: 2263 content = json.dumps(registry_json, indent=2, sort_keys=True) + "\n" 2264 path = ( 2265 f"{output_target.bucket_root}/{_REGISTRIES_PREFIX}/" 2266 f"{registry_type}_registry.json" 2267 ) 2268 _write_registry_blob( 2269 output_fs, 2270 store=output_target, 2271 path=path, 2272 content=content, 2273 cache_control=_REGISTRY_INDEX_CACHE_CONTROL, 2274 ) 2275 _log_progress( 2276 " Wrote %s_registry.json (%d entries)", 2277 registry_type, 2278 entry_count, 2279 ) 2280 2281 # --- Step 11: Compile composite registry JSON (superset) --- 2282 _log_progress("Step 11: Compiling composite_registry.json (superset)...") 2283 composite_json = _build_composite_registry_json( 2284 cloud_entries=entries_by_registry_type.get("cloud", []), 2285 oss_entries=entries_by_registry_type.get("oss", []), 2286 ) 2287 composite_entry_count = len(composite_json["sources"]) + len( 2288 composite_json["destinations"] 2289 ) 2290 result.composite_registry_entries = composite_entry_count 2291 if dry_run: 2292 _log_progress( 2293 " [DRY RUN] Would write composite_registry.json (%d entries)", 2294 composite_entry_count, 2295 ) 2296 else: 2297 composite_content = json.dumps(composite_json, indent=2, sort_keys=True) + "\n" 2298 path = ( 2299 f"{output_target.bucket_root}/{_REGISTRIES_PREFIX}/composite_registry.json" 2300 ) 2301 _write_registry_blob( 2302 output_fs, 2303 store=output_target, 2304 path=path, 2305 content=composite_content, 2306 cache_control=_REGISTRY_INDEX_CACHE_CONTROL, 2307 ) 2308 _log_progress( 2309 " Wrote composite_registry.json (%d entries)", 2310 composite_entry_count, 2311 ) 2312 2313 # --- Step 12: Per-connector version indexes (parallel) --- 2314 _log_progress( 2315 "Step 12: Writing per-connector version indexes (max_workers=%d)...", 2316 _COMPILE_WRITE_MAX_WORKERS, 2317 ) 2318 base = f"{output_target.bucket_root}/{METADATA_FOLDER}/airbyte" 2319 sorted_connectors = sorted( 2320 set(connector_versions) 2321 if connector_name is None 2322 else set(connector_name) & set(connector_versions) 2323 ) 2324 2325 def _write_one_version_index(connector: str) -> bool: 2326 """Build and write a single connector's versions.json.""" 2327 versions = connector_versions[connector] 2328 latest_v = latest_versions.get(connector) 2329 rc_v_list = rc_versions.get(connector) 2330 index_path = f"{base}/{connector}/versions.json" 2331 previous_index = _read_previous_version_index( 2332 fs, 2333 f"{store.bucket_root}/{METADATA_FOLDER}/airbyte/{connector}/versions.json", 2334 connector=connector, 2335 ) 2336 output_previous_index = ( 2337 previous_index 2338 if output_store is None 2339 else _read_previous_version_index( 2340 output_fs, 2341 index_path, 2342 connector=connector, 2343 ) 2344 ) 2345 index = _build_version_index( 2346 fs, 2347 store=store, 2348 connector=connector, 2349 versions=versions, 2350 yanked=yanked, 2351 latest_version=latest_v, 2352 rc_version=rc_v_list[0] if rc_v_list else None, 2353 rc_versions_all=rc_v_list, 2354 previous_index=previous_index, 2355 full_restate=with_full_restate, 2356 release_attribution=( 2357 release_attribution_index.connectors.get(connector) 2358 if release_attribution_index 2359 else None 2360 ), 2361 ) 2362 content = _version_index_content(index) 2363 if _version_index_is_unchanged(output_previous_index, index): 2364 action = "Would skip unchanged" if dry_run else "Skipped unchanged" 2365 _log_progress(" %s %s/versions.json", action, connector) 2366 return False 2367 if dry_run: 2368 _log_progress( 2369 " [DRY RUN] Would write %s/versions.json (%d versions)", 2370 connector, 2371 len(versions), 2372 ) 2373 else: 2374 parent = index_path.rsplit("/", 1)[0] 2375 if output_target.env == "local": 2376 output_fs.makedirs(parent, exist_ok=True) 2377 with output_fs.open(index_path, "w") as f: 2378 f.write(content) 2379 return True 2380 2381 with ThreadPoolExecutor(max_workers=_COMPILE_WRITE_MAX_WORKERS) as pool: 2382 futures = { 2383 pool.submit(_write_one_version_index, c): c for c in sorted_connectors 2384 } 2385 for i, future in enumerate(as_completed(futures), 1): 2386 connector = futures[future] 2387 try: 2388 wrote = future.result() 2389 if wrote: 2390 result.version_indexes_written += 1 2391 else: 2392 result.version_indexes_skipped += 1 2393 except Exception as exc: 2394 error_msg = f"Failed to write versions.json for {connector}: {exc}" 2395 logger.error(error_msg) 2396 result.errors.append(error_msg) 2397 if i % 100 == 0: 2398 _log_progress( 2399 " Wrote %d / %d version indexes...", i, len(sorted_connectors) 2400 ) 2401 2402 # --- Step 13: Specs secrets mask (optional) --- 2403 if with_secrets_mask: 2404 _log_progress("Step 13: Generating specs secrets mask...") 2405 # Reuse entries collected during Step 10 to avoid redundant GCS reads. 2406 secret_names = _extract_secret_property_names(all_registry_entries) 2407 sorted_names = sorted(secret_names) 2408 result.specs_secrets_mask_properties = len(sorted_names) 2409 mask_content = yaml.dump({"properties": sorted_names}, default_flow_style=False) 2410 mask_path = ( 2411 f"{output_target.bucket_root}/{_REGISTRIES_PREFIX}/" 2412 f"{_SPECS_SECRETS_MASK_FILENAME}" 2413 ) 2414 2415 _log_progress( 2416 " Found %d secret properties: %s", 2417 len(sorted_names), 2418 ", ".join(sorted_names), 2419 ) 2420 2421 if dry_run: 2422 _log_progress( 2423 " [DRY RUN] Would write %s", 2424 _SPECS_SECRETS_MASK_FILENAME, 2425 ) 2426 else: 2427 parent = mask_path.rsplit("/", 1)[0] 2428 if output_target.env == "local": 2429 output_fs.makedirs(parent, exist_ok=True) 2430 with output_fs.open(mask_path, "w") as f: 2431 f.write(mask_content) 2432 _log_progress( 2433 " Wrote %s", 2434 _SPECS_SECRETS_MASK_FILENAME, 2435 ) 2436 2437 if output_store is not None and not dry_run: 2438 compile_metadata = ( 2439 json.dumps( 2440 { 2441 "source_store": store.bucket_root, 2442 "output_store": output_store.bucket_root, 2443 "compiled_at": datetime.now(timezone.utc).isoformat(), 2444 }, 2445 indent=2, 2446 ) 2447 + "\n" 2448 ) 2449 _write_registry_blob( 2450 output_fs, 2451 store=output_target, 2452 path=f"{output_target.bucket_root}/{_COMPILE_METADATA_FILENAME}", 2453 content=compile_metadata, 2454 ) 2455 2456 _log_progress(result.summary()) 2457 return result
Compile the registry: sync latest/ dirs and write index files.
Steps:
- Glob for all
metadata.yamlto discover (connector, version) pairs.- Glob for active marker files.
- Compute the latest GA semver per connector.
- Compute active release candidates from versioned markers.
- Glob for
version=*markers inlatest/dirs for a fast check.- Delete stale
latest/dirs and recursively copy the versioned dir.- Synthesize missing registry entries from pinned latest overrides.
- (Optional) Legacy migration: delete disabled registry entries.
- (Optional) Read latest connector metrics.
- Write global registry JSONs. When an output store is supplied, read each connector entry from its computed version directory instead of
latest/.- Write composite registry JSON.
- Write per-connector
versions.json.- (Optional) Regenerate
specs_secrets_mask.yaml.
Arguments:
- store: Registry store (bucket + optional prefix).
- connector_name: If provided, only resync
latest/directories for these connectors (steps 5-6). Global registry indexes always operate on the full set so they remain complete; per-connectorversions.jsonwrites are scoped to the requested connectors. - dry_run: If True, report what would be done without writing.
- with_secrets_mask: If True, regenerate
specs_secrets_mask.yaml. - with_legacy_migration: If set, run the named migration step.
Currently supported:
"v1"— delete{registry_type}.jsonfiles for connectors whoseregistryOverrides.{registry}.enabledisfalse. - with_metrics: If True, inject latest connector metrics from the
analytics JSONL export into
generated.metrics. - force: If True, resync all connectors' latest/ directories even if the existing version marker matches the computed latest version. This is useful when metadata content changes without a version bump.
- with_full_restate: If
True, re-read every version'smetadata.yamlinstead of carrying forward existing release attribution. This is expensive for all connectors, at approximately 33.5k reads. When norelease_attribution_indexis supplied, historical release blocks that exist only in the attribution index cannot be re-derived and will be discarded. - release_attribution_index: Optional Phase 1 index used to seed missing per-version release blocks without rewriting historical metadata.
The compile performs one additional GCS read per connector for its prior
versions.json, trading approximately 716 reads for avoiding historical
metadata reads on the default path.
Returns:
A
CompileResultdescribing what was done.
90def find_unpublished_connectors( 91 repo_path: str | Path, 92 bucket_name: str, 93 connector_names: list[str] | None = None, 94) -> AuditResult: 95 """Find connectors whose local version is not published on GCS. 96 97 For each connector in the local checkout, reads `dockerImageTag` from 98 `metadata.yaml` and checks whether `metadata/<docker-repo>/<version>/metadata.yaml` 99 exists in the GCS bucket. Connectors that are archived, disabled on all 100 registries, or have RC versions are skipped. 101 102 Args: 103 repo_path: Path to the Airbyte monorepo checkout. 104 bucket_name: GCS bucket name to check against. 105 connector_names: Optional list of connector names to check. 106 If `None`, discovers all connectors in the repo. 107 108 Returns: 109 An `AuditResult` containing unpublished connectors and metadata. 110 """ 111 repo_path = Path(repo_path) 112 connectors_dir = repo_path / CONNECTOR_PATH_PREFIX 113 114 if not connectors_dir.exists(): 115 raise ValueError(f"Connectors directory not found: {connectors_dir}") 116 117 # Discover connector names if not provided 118 if connector_names is None: 119 connector_names = sorted( 120 d.name 121 for d in connectors_dir.iterdir() 122 if d.is_dir() and (d / METADATA_FILE_NAME).exists() 123 ) 124 125 result = AuditResult() 126 127 # Collect connectors and their versions first, then batch-check GCS 128 to_check: list[tuple[str, str]] = [] # (connector_name, version) 129 130 for name in connector_names: 131 metadata_path = connectors_dir / name / METADATA_FILE_NAME 132 metadata = _read_local_metadata(metadata_path) 133 if metadata is None: 134 result.errors.append(f"{name}: metadata.yaml not found or unreadable") 135 continue 136 137 if _is_archived(metadata): 138 result.skipped_archived.append(name) 139 continue 140 141 if _is_disabled_on_all_registries(metadata): 142 result.skipped_disabled.append(name) 143 continue 144 145 data = metadata.get("data", {}) 146 version = data.get("dockerImageTag") 147 if not version: 148 result.errors.append(f"{name}: no dockerImageTag in metadata") 149 continue 150 151 if _is_rc_version(version): 152 result.skipped_rc.append(name) 153 continue 154 155 to_check.append((name, version)) 156 157 if not to_check: 158 result.checked_count = 0 159 return result 160 161 # Check GCS for each connector version 162 storage_client = get_gcs_storage_client() 163 bucket = storage_client.bucket(bucket_name) 164 165 for name, version in to_check: 166 result.checked_count += 1 167 blob_path = f"{METADATA_FOLDER}/airbyte/{name}/{version}/{METADATA_FILE_NAME}" 168 blob = bucket.blob(blob_path) 169 170 try: 171 exists = blob.exists() 172 except Exception as e: 173 result.errors.append(f"{name}: GCS check failed: {e}") 174 continue 175 176 if not exists: 177 logger.info( 178 "Unpublished: %s version %s (checked %s)", 179 name, 180 version, 181 blob_path, 182 ) 183 result.unpublished.append( 184 UnpublishedConnector(connector_name=name, local_version=version) 185 ) 186 187 logger.info( 188 "Audit complete: %d checked, %d unpublished, %d archived-skipped, %d disabled-skipped, %d rc-skipped", 189 result.checked_count, 190 len(result.unpublished), 191 len(result.skipped_archived), 192 len(result.skipped_disabled), 193 len(result.skipped_rc), 194 ) 195 196 return result
Find connectors whose local version is not published on GCS.
For each connector in the local checkout, reads dockerImageTag from
metadata.yaml and checks whether metadata/<docker-repo>/<version>/metadata.yaml
exists in the GCS bucket. Connectors that are archived, disabled on all
registries, or have RC versions are skipped.
Arguments:
- repo_path: Path to the Airbyte monorepo checkout.
- bucket_name: GCS bucket name to check against.
- connector_names: Optional list of connector names to check.
If
None, discovers all connectors in the repo.
Returns:
An
AuditResultcontaining unpublished connectors and metadata.
709def generate_version_artifacts( 710 metadata_file: Path, 711 docker_image: str, 712 output_dir: Path | None = None, 713 repo_root: Path | None = None, 714 dry_run: bool = False, 715 with_validate: bool = True, 716 with_dependency_dump: bool = True, 717 with_sbom: bool = True, 718 pr_number: int | None = None, 719 with_github_lookup: bool = True, 720) -> GenerateResult: 721 """Generate all version artifacts for a connector release. 722 723 Artifacts are enriched with git commit info, SBOM URL, and (when applicable) 724 components SHA before writing. Validation is run after generation by default. 725 726 Args: 727 metadata_file: Path to the connector's `metadata.yaml`. 728 docker_image: Docker image to run spec against (e.g. `airbyte/source-faker:6.2.38`). 729 output_dir: Directory to write artifacts to. If `None`, a temp directory is created. 730 repo_root: Root of the Airbyte repo checkout (for resolving `doc.md`). 731 If `None`, inferred by walking up from `metadata_file`. 732 dry_run: If `True`, report what would be generated without writing or running docker. 733 with_validate: If `True` (default), run metadata validators after generation. 734 Pass `False` (`--no-validate`) to skip. 735 with_dependency_dump: If `True` (default), generate `dependencies.json` 736 for Python connectors. Pass `False` (`--no-dependency-dump`) to skip. 737 with_sbom: If `True` (default), generate `spdx.json` (SBOM) for 738 connectors. Pass `False` (`--no-sbom`) to skip. 739 with_github_lookup: If `True` (default), resolve publish PR authors 740 through GitHub GraphQL when a PR number is available. 741 742 Returns: 743 A `GenerateResult` describing what was produced. 744 """ 745 # --- Load metadata --- 746 if not metadata_file.exists(): 747 raise FileNotFoundError(f"Metadata file not found: {metadata_file}") 748 749 raw_metadata: dict[str, Any] = yaml.safe_load(metadata_file.read_text()) 750 metadata_data: dict[str, Any] = raw_metadata.get("data", {}) 751 752 connector_name = metadata_data.get("dockerRepository", "unknown").replace( 753 "airbyte/", "" 754 ) 755 version = metadata_data.get("dockerImageTag", "unknown") 756 757 # --- Resolve output directory --- 758 if output_dir is None: 759 output_dir = Path( 760 tempfile.mkdtemp(prefix=f"connector-artifacts-{connector_name}-{version}-") 761 ) 762 output_dir.mkdir(parents=True, exist_ok=True) 763 764 result = GenerateResult( 765 connector_name=connector_name, 766 version=version, 767 docker_image=docker_image, 768 output_dir=str(output_dir), 769 dry_run=dry_run, 770 ) 771 772 if dry_run: 773 logger.info("[DRY RUN] Would generate artifacts to %s", output_dir) 774 result.artifacts_written = [ 775 "metadata.yaml", 776 "icon.svg", 777 "doc.md", 778 "cloud.json", 779 "oss.json", 780 "manifest.yaml (if present)", 781 "components.zip (if components.py present)", 782 "components.zip.sha256 (if components.py present)", 783 f"version={version}", 784 ] 785 if with_sbom: 786 result.artifacts_written.append(SBOM_FILE_NAME) 787 if with_dependency_dump: 788 result.artifacts_written.append("dependencies.json (if Python connector)") 789 return result 790 791 # --- Prepare metadata output --- 792 metadata_out = output_dir / "metadata.yaml" 793 result.artifacts_written.append("metadata.yaml") 794 795 # --- Enrich metadata with git info *before* building registry entries so 796 # that `generated.git` propagates into `cloud.json` / `oss.json`. --- 797 raw_metadata = _enrich_metadata_git_info(raw_metadata, metadata_file) 798 raw_metadata = _enrich_metadata_release( 799 raw_metadata, 800 metadata_file, 801 pr_number=pr_number, 802 with_github_lookup=with_github_lookup, 803 ) 804 805 # --- Generate SBOM from the connector Docker image --- 806 sbom_generated = False 807 if not with_sbom: 808 logger.info("SBOM generation disabled via --no-sbom.") 809 else: 810 try: 811 sbom_path = generate_sbom(docker_image, output_dir) 812 except RuntimeError as exc: 813 logger.warning("SBOM generation failed (non-fatal): %s", exc) 814 except (FileNotFoundError, subprocess.TimeoutExpired): 815 logger.warning("Docker not available or SBOM generation timed out.") 816 else: 817 result.artifacts_written.append(SBOM_FILE_NAME) 818 sbom_generated = True 819 logger.info("Generated SBOM: %s", sbom_path) 820 821 # --- Enrich metadata with SBOM URL --- 822 raw_metadata = _enrich_metadata_sbom_url( 823 raw_metadata, sbom_generated=sbom_generated 824 ) 825 826 # --- Run docker spec for cloud and oss --- 827 specs: dict[str, dict[str, Any]] = {} 828 for mode in VALID_REGISTRIES: 829 try: 830 specs[mode] = _run_docker_spec(docker_image, mode) 831 logger.info("Got %s spec from docker image %s", mode, docker_image) 832 except RuntimeError as exc: 833 error_msg = f"Failed to get {mode} spec: {exc}" 834 logger.error(error_msg) 835 result.errors.append(error_msg) 836 837 # --- Generate dependencies.json for Python connectors --- 838 # This must happen *before* building registry entries so that the 839 # local dependencies data can be used for packageInfo without a GCS 840 # round-trip. 841 local_dependencies: dict[str, Any] | None = None 842 if not with_dependency_dump: 843 logger.info("Dependency generation disabled via --no-dependency-dump.") 844 elif _is_python_connector(metadata_data): 845 logger.info("Python connector detected — generating dependencies.json") 846 local_dependencies = generate_python_dependencies_file( 847 metadata_data=metadata_data, 848 docker_image=docker_image, 849 output_dir=output_dir, 850 ) 851 if local_dependencies is not None: 852 result.artifacts_written.append(CONNECTOR_DEPENDENCY_FILE_NAME) 853 else: 854 logger.info( 855 "Non-Python connector (%s) — skipping dependencies.json generation.", 856 connector_name, 857 ) 858 859 # --- Generate registry entries (cloud.json, oss.json) --- 860 for registry_type in VALID_REGISTRIES: 861 if not is_registry_enabled(metadata_data, registry_type): 862 logger.info( 863 "Registry type %s is not enabled for %s, skipping %s.json generation.", 864 registry_type, 865 connector_name, 866 registry_type, 867 ) 868 continue 869 870 spec = specs.get(registry_type) 871 if spec is None: 872 error_msg = ( 873 f"Cannot generate {registry_type}.json: no spec available " 874 f"(docker spec for {registry_type} failed or was not run)." 875 ) 876 result.errors.append(error_msg) 877 continue 878 879 registry_entry = _build_registry_entry( 880 metadata_data, 881 registry_type, 882 spec, 883 local_dependencies=local_dependencies, 884 ) 885 886 out_path = output_dir / f"{registry_type}.json" 887 out_path.write_text( 888 json.dumps(registry_entry, indent=2, sort_keys=True, default=_json_serial) 889 + "\n" 890 ) 891 result.artifacts_written.append(f"{registry_type}.json") 892 logger.info("Wrote %s", out_path) 893 894 # --- Copy icon.svg (sibling of metadata.yaml in the connector directory) --- 895 icon_source = metadata_file.parent / "icon.svg" 896 if icon_source.is_file(): 897 icon_out = output_dir / "icon.svg" 898 shutil.copy2(icon_source, icon_out) 899 result.artifacts_written.append("icon.svg") 900 logger.info("Wrote %s", icon_out) 901 else: 902 logger.warning("No icon.svg found at %s.", icon_source) 903 result.errors.append("Icon file is missing.") 904 905 # --- Copy doc.md (derived from documentationUrl in metadata) --- 906 if repo_root is None: 907 # Infer repo root by walking up from metadata_file looking for .git 908 # Note: .git can be a directory (normal clone) or a file (git worktree) 909 # Resolve to absolute path first so the walk-up works with relative paths. 910 candidate = metadata_file.resolve().parent 911 while candidate != candidate.parent: 912 git_indicator = candidate / ".git" 913 if git_indicator.is_dir() or git_indicator.is_file(): 914 repo_root = candidate 915 break 916 candidate = candidate.parent 917 918 if repo_root is not None: 919 doc_source = _resolve_doc_path(metadata_data, repo_root) 920 if doc_source is not None and doc_source.is_file(): 921 doc_out = output_dir / DOC_FILE_NAME 922 shutil.copy2(doc_source, doc_out) 923 result.artifacts_written.append(DOC_FILE_NAME) 924 logger.info("Wrote %s (from %s)", doc_out, doc_source) 925 else: 926 error_msg = ( 927 f"Documentation file not found: {doc_source}. " 928 f"Derived from documentationUrl in metadata." 929 ) 930 logger.error(error_msg) 931 result.errors.append(error_msg) 932 else: 933 error_msg = "Cannot resolve doc.md: repo root not found." 934 logger.error(error_msg) 935 result.errors.append(error_msg) 936 937 # --- Copy manifest.yaml (from connector root, if present) --- 938 connector_dir = metadata_file.parent 939 manifest_source = connector_dir / MANIFEST_FILE_NAME 940 components_sha256: str | None = None 941 if manifest_source.is_file(): 942 manifest_out = output_dir / MANIFEST_FILE_NAME 943 shutil.copy2(manifest_source, manifest_out) 944 result.artifacts_written.append(MANIFEST_FILE_NAME) 945 logger.info("Wrote %s", manifest_out) 946 947 # --- Generate components.zip if components.py exists --- 948 components_source = connector_dir / COMPONENTS_PY_FILE_NAME 949 if components_source.is_file(): 950 zip_path, sha256_path = _create_components_zip( 951 manifest_path=manifest_source, 952 components_path=components_source, 953 output_dir=output_dir, 954 ) 955 result.artifacts_written.append(COMPONENTS_ZIP_FILE_NAME) 956 result.artifacts_written.append(COMPONENTS_ZIP_SHA256_FILE_NAME) 957 logger.info("Wrote %s and %s", zip_path, sha256_path) 958 # Read back the SHA256 for metadata enrichment 959 components_sha256 = sha256_path.read_text().strip() 960 else: 961 logger.info( 962 "No manifest.yaml at %s — skipping manifest artifacts.", manifest_source 963 ) 964 965 # --- Enrich metadata with components SHA (after zip creation) --- 966 raw_metadata = _enrich_metadata_components_sha(raw_metadata, components_sha256) 967 968 # --- Write final enriched metadata.yaml --- 969 # Use sort_keys=True to match the legacy pipeline's alphabetical key ordering. 970 # After Registry 2.0 launches we are free to change the key ordering. 971 metadata_out.write_text( 972 yaml.dump(raw_metadata, default_flow_style=False, sort_keys=True) 973 ) 974 logger.info("Wrote enriched %s", metadata_out) 975 976 # --- Write version marker file (version=<semver>) --- 977 # This zero-byte file is used by the compile step as a fast-check marker. 978 # Including it in the generated artifacts means `latest/` gets the marker 979 # for free via a recursive copy, removing the need for a separate write. 980 marker_file = output_dir / f"version={version}" 981 marker_file.write_bytes(b"") 982 result.artifacts_written.append(f"version={version}") 983 logger.info("Wrote version marker %s", marker_file) 984 985 # --- Validate metadata (after generation) --- 986 if with_validate: 987 logger.info("Running post-generation validation...") 988 doc_path: str | None = None 989 if repo_root is not None: 990 resolved = _resolve_doc_path(metadata_data, repo_root) 991 doc_path = str(resolved) if resolved else None 992 validation = validate_metadata( 993 metadata_data=metadata_data, 994 opts=ValidateOptions(docs_path=doc_path), 995 ) 996 if not validation.passed: 997 for err in validation.errors: 998 logger.error("Validation error: %s", err) 999 result.validation_errors = validation.errors 1000 else: 1001 logger.info("Validation passed (%d validators).", validation.validators_run) 1002 1003 return result
Generate all version artifacts for a connector release.
Artifacts are enriched with git commit info, SBOM URL, and (when applicable) components SHA before writing. Validation is run after generation by default.
Arguments:
- metadata_file: Path to the connector's
metadata.yaml. - docker_image: Docker image to run spec against (e.g.
airbyte/source-faker:6.2.38). - output_dir: Directory to write artifacts to. If
None, a temp directory is created. - repo_root: Root of the Airbyte repo checkout (for resolving
doc.md). IfNone, inferred by walking up frommetadata_file. - dry_run: If
True, report what would be generated without writing or running docker. - with_validate: If
True(default), run metadata validators after generation. PassFalse(--no-validate) to skip. - with_dependency_dump: If
True(default), generatedependencies.jsonfor Python connectors. PassFalse(--no-dependency-dump) to skip. - with_sbom: If
True(default), generatespdx.json(SBOM) for connectors. PassFalse(--no-sbom) to skip. - with_github_lookup: If
True(default), resolve publish PR authors through GitHub GraphQL when a PR number is available.
Returns:
A
GenerateResultdescribing what was produced.
37def get_connector_metadata(repo_path: Path, connector_name: str) -> ConnectorMetadata: 38 """Read connector metadata from metadata.yaml. 39 40 Args: 41 repo_path: Path to the Airbyte monorepo. 42 connector_name: The connector technical name (e.g., 'source-github'). 43 44 Returns: 45 ConnectorMetadata object with the connector's metadata. 46 47 Raises: 48 FileNotFoundError: If the connector directory or metadata file doesn't exist. 49 """ 50 connector_dir = repo_path / CONNECTOR_PATH_PREFIX / connector_name 51 if not connector_dir.exists(): 52 raise FileNotFoundError(f"Connector directory not found: {connector_dir}") 53 54 metadata_file = connector_dir / METADATA_FILE_NAME 55 if not metadata_file.exists(): 56 raise FileNotFoundError(f"Metadata file not found: {metadata_file}") 57 58 with open(metadata_file) as f: 59 metadata = yaml.safe_load(f) 60 61 data = metadata.get("data", {}) 62 return ConnectorMetadata( 63 name=connector_name, 64 docker_repository=data.get("dockerRepository", f"airbyte/{connector_name}"), 65 docker_image_tag=data.get("dockerImageTag", "unknown"), 66 support_level=data.get("supportLevel"), 67 definition_id=data.get("definitionId"), 68 )
Read connector metadata from metadata.yaml.
Arguments:
- repo_path: Path to the Airbyte monorepo.
- connector_name: The connector technical name (e.g., 'source-github').
Returns:
ConnectorMetadata object with the connector's metadata.
Raises:
- FileNotFoundError: If the connector directory or metadata file doesn't exist.
71def get_gcs_publish_path( 72 connector_name: str, 73 artifact_type: str, 74 version: str = LATEST_GCS_FOLDER_NAME, 75) -> str: 76 """Compute the GCS path for a connector artifact for publishing. 77 78 All connectors use the airbyte/{connector_name} convention. 79 """ 80 artifact_files = { 81 "metadata": METADATA_FILE_NAME, 82 "spec": "spec.json", 83 "icon": "icon.svg", 84 "doc": "doc.md", 85 } 86 87 if artifact_type not in artifact_files: 88 raise ValueError( 89 f"Unknown artifact type: {artifact_type}. " 90 f"Valid types are: {', '.join(artifact_files.keys())}" 91 ) 92 93 file_name = artifact_files[artifact_type] 94 return f"{METADATA_FOLDER}/airbyte/{connector_name}/{version}/{file_name}"
Compute the GCS path for a connector artifact for publishing.
All connectors use the airbyte/{connector_name} convention.
225def get_registry(store: RegistryStore) -> Registry: 226 """Factory for obtaining the right store implementation.""" 227 228 if store.store_type == StoreType.CORAL: 229 from airbyte_ops_mcp.registry.coral_registry_store import CoralRegistry 230 231 return CoralRegistry(store) 232 233 if store.store_type == StoreType.SONAR: 234 from airbyte_ops_mcp.registry.sonar_registry_store import SonarRegistry 235 236 return SonarRegistry(store) 237 238 # defensive: StoreType is an Enum, but keep this for readability 239 raise ValueError(f"Unknown store type: {store.store_type}")
Factory for obtaining the right store implementation.
43def get_registry_entry( 44 connector_name: str, 45 bucket_name: str, 46 version: str = LATEST_GCS_FOLDER_NAME, 47 prefix: str = "", 48) -> dict[str, Any]: 49 """Get a connector's registry entry from GCS. 50 51 Reads metadata for a connector from the registry stored in GCS. 52 53 Args: 54 connector_name: The connector name (e.g., "source-faker", "destination-postgres") 55 bucket_name: Name of the GCS bucket containing the registry 56 version: Version folder name (e.g., "latest", "1.2.3") 57 prefix: Optional path prefix within the bucket; leading and trailing 58 slashes are ignored. 59 60 Returns: 61 dict: The connector's metadata as a dictionary 62 63 Raises: 64 ValueError: If GCS credentials are not configured, or if the metadata has an invalid structure 65 FileNotFoundError: If the connector metadata is not found in the registry 66 yaml.YAMLError: If the metadata file contains invalid YAML syntax 67 """ 68 storage_client = get_gcs_storage_client() 69 bucket = storage_client.bucket(bucket_name) 70 71 # Construct the path to the metadata file 72 # Pattern: metadata/airbyte/{connector_name}/{version}/metadata.yaml 73 normalized_prefix = prefix.strip("/") 74 prefix_part = f"{normalized_prefix}/" if normalized_prefix else "" 75 blob_path = ( 76 f"{prefix_part}{METADATA_FOLDER}/airbyte/" 77 f"{connector_name}/{version}/{METADATA_FILE_NAME}" 78 ) 79 blob = bucket.blob(blob_path) 80 81 logger.info(f"Reading registry entry for {connector_name} from {blob_path}") 82 83 # Read the file 84 content = safe_read_gcs_file(blob) 85 if content is None: 86 raise FileNotFoundError( 87 f"Connector metadata not found in registry: {connector_name}. " 88 f"Checked path: {blob_path}" 89 ) 90 91 # Parse YAML 92 try: 93 metadata = yaml.safe_load(content) 94 if metadata is None or not isinstance(metadata, dict): 95 raise ValueError(f"Metadata file {blob_path} has an invalid structure") 96 return metadata 97 except yaml.YAMLError as e: 98 logger.error( 99 "Failed to parse metadata for %s from %s: %s", 100 connector_name, 101 blob_path, 102 e, 103 ) 104 raise
Get a connector's registry entry from GCS.
Reads metadata for a connector from the registry stored in GCS.
Arguments:
- connector_name: The connector name (e.g., "source-faker", "destination-postgres")
- bucket_name: Name of the GCS bucket containing the registry
- version: Version folder name (e.g., "latest", "1.2.3")
- prefix: Optional path prefix within the bucket; leading and trailing slashes are ignored.
Returns:
dict: The connector's metadata as a dictionary
Raises:
- ValueError: If GCS credentials are not configured, or if the metadata has an invalid structure
- FileNotFoundError: If the connector metadata is not found in the registry
- yaml.YAMLError: If the metadata file contains invalid YAML syntax
107def get_registry_spec( 108 connector_name: str, 109 bucket_name: str, 110 version: str = LATEST_GCS_FOLDER_NAME, 111) -> dict[str, Any]: 112 """Get a connector's spec from GCS. 113 114 Reads the connector specification from the registry stored in GCS. 115 116 Args: 117 connector_name: The connector name (e.g., "source-faker", "destination-postgres") 118 bucket_name: Name of the GCS bucket containing the registry 119 version: Version folder name (e.g., "latest", "1.2.3") 120 121 Returns: 122 dict: The connector's spec as a dictionary 123 124 Raises: 125 ValueError: If GCS credentials are not configured, or if the spec is not a JSON object 126 FileNotFoundError: If the connector spec is not found in the registry 127 json.JSONDecodeError: If the spec file contains invalid JSON syntax 128 """ 129 storage_client = get_gcs_storage_client() 130 bucket = storage_client.bucket(bucket_name) 131 132 # Construct the path to the spec file 133 # Pattern: metadata/airbyte/{connector_name}/{version}/spec.json 134 blob_path = f"{METADATA_FOLDER}/airbyte/{connector_name}/{version}/{SPEC_FILE_NAME}" 135 blob = bucket.blob(blob_path) 136 137 logger.info(f"Reading spec for {connector_name} from {blob_path}") 138 139 # Read the file 140 content = safe_read_gcs_file(blob) 141 if content is None: 142 raise FileNotFoundError( 143 f"Connector spec not found in registry: {connector_name}. " 144 f"Checked path: {blob_path}" 145 ) 146 147 # Parse JSON 148 try: 149 spec = json.loads(content) 150 if spec is None or not isinstance(spec, dict): 151 raise ValueError( 152 f"Spec file for {connector_name} at {blob_path} is not a JSON object" 153 ) 154 return spec 155 except json.JSONDecodeError as e: 156 logger.error( 157 "Failed to parse spec for %s from %s: %s", 158 connector_name, 159 blob_path, 160 e, 161 ) 162 raise
Get a connector's spec from GCS.
Reads the connector specification from the registry stored in GCS.
Arguments:
- connector_name: The connector name (e.g., "source-faker", "destination-postgres")
- bucket_name: Name of the GCS bucket containing the registry
- version: Version folder name (e.g., "latest", "1.2.3")
Returns:
dict: The connector's spec as a dictionary
Raises:
- ValueError: If GCS credentials are not configured, or if the spec is not a JSON object
- FileNotFoundError: If the connector spec is not found in the registry
- json.JSONDecodeError: If the spec file contains invalid JSON syntax
331def list_connector_versions(connector_name: str, bucket_name: str) -> list[str]: 332 """List all versions of a connector in the registry. 333 334 Scans the GCS bucket to find all versions of a specific connector. 335 336 Args: 337 connector_name: The connector name (e.g., "source-faker") 338 bucket_name: Name of the GCS bucket containing the registry 339 340 Returns: 341 list[str]: Sorted list of version strings (excluding 'latest' and 'release_candidate') 342 343 Raises: 344 ValueError: If GCS credentials are not configured 345 """ 346 storage_client = get_gcs_storage_client() 347 bucket = storage_client.bucket(bucket_name) 348 349 # List all blobs matching the pattern: metadata/airbyte/{connector_name}/*/metadata.yaml 350 glob_pattern = f"{METADATA_FOLDER}/airbyte/{connector_name}/*/{METADATA_FILE_NAME}" 351 logger.info(f"Listing versions for {connector_name} with pattern: {glob_pattern}") 352 353 try: 354 blobs = bucket.list_blobs(match_glob=glob_pattern) 355 except Exception as e: 356 logger.error(f"Error listing blobs in bucket {bucket_name}: {e}") 357 raise 358 359 # Extract versions from blob paths 360 # Path format: metadata/airbyte/{connector-name}/{version}/metadata.yaml 361 versions: set[str] = set() 362 for blob in blobs: 363 path_parts = blob.name.split("/") 364 # Path should be: metadata / airbyte / connector-name / version / metadata.yaml 365 if len(path_parts) >= 5: 366 version = path_parts[3] 367 # Exclude special folders 368 if version not in ("latest", "release_candidate"): 369 versions.add(version) 370 371 return sorted(versions)
List all versions of a connector in the registry.
Scans the GCS bucket to find all versions of a specific connector.
Arguments:
- connector_name: The connector name (e.g., "source-faker")
- bucket_name: Name of the GCS bucket containing the registry
Returns:
list[str]: Sorted list of version strings (excluding 'latest' and 'release_candidate')
Raises:
- ValueError: If GCS credentials are not configured
165def list_registry_connectors(bucket_name: str) -> list[str]: 166 """List all connectors in the registry. 167 168 Scans the GCS bucket to find all connectors that have metadata files. 169 170 Args: 171 bucket_name: Name of the GCS bucket containing the registry 172 173 Returns: 174 list[str]: Sorted list of connector names 175 176 Raises: 177 ValueError: If GCS credentials are not configured 178 """ 179 storage_client = get_gcs_storage_client() 180 bucket = storage_client.bucket(bucket_name) 181 182 # List all blobs matching the pattern: metadata/airbyte/*/latest/metadata.yaml 183 glob_pattern = ( 184 f"{METADATA_FOLDER}/airbyte/*/{LATEST_GCS_FOLDER_NAME}/{METADATA_FILE_NAME}" 185 ) 186 logger.info(f"Listing connectors with pattern: {glob_pattern}") 187 188 try: 189 blobs = bucket.list_blobs(match_glob=glob_pattern) 190 except Exception as e: 191 logger.error(f"Error listing blobs in bucket {bucket_name}: {e}") 192 raise 193 194 # Extract connector names from blob paths 195 # Path format: metadata/airbyte/{connector-name}/latest/metadata.yaml 196 connector_names: set[str] = set() 197 for blob in blobs: 198 path_parts = blob.name.split("/") 199 # Path should be: metadata / airbyte / connector-name / latest / metadata.yaml 200 if len(path_parts) >= 5: 201 connector_name = path_parts[2] 202 connector_names.add(connector_name) 203 204 return sorted(connector_names)
List all connectors in the registry.
Scans the GCS bucket to find all connectors that have metadata files.
Arguments:
- bucket_name: Name of the GCS bucket containing the registry
Returns:
list[str]: Sorted list of connector names
Raises:
- ValueError: If GCS credentials are not configured
207def list_registry_connectors_filtered( 208 bucket_name: str, 209 *, 210 support_level: SupportLevel | None = None, 211 min_support_level: SupportLevel | None = None, 212 connector_type: ConnectorType | None = None, 213 language: ConnectorLanguage | None = None, 214 prefix: str = "", 215) -> list[str]: 216 """List connectors from the compiled cloud registry index with filtering. 217 218 When any filter is applied, reads the compiled `cloud_registry.json` index 219 instead of globbing individual metadata blobs. This is significantly faster 220 because the index is a single JSON file containing all connector entries. 221 222 When no filters are applied, falls back to the existing glob-based search 223 which captures all connectors (including OSS-only connectors not in the 224 Cloud index). 225 226 Args: 227 bucket_name: Name of the GCS bucket containing the registry. 228 support_level: Exact support level to match (e.g., `SupportLevel.CERTIFIED`). 229 min_support_level: Minimum support level threshold. Returns connectors 230 at or above this level. 231 connector_type: Filter by connector type (`ConnectorType.SOURCE` or 232 `ConnectorType.DESTINATION`). 233 language: Filter by implementation language (e.g., `ConnectorLanguage.PYTHON`). 234 prefix: Optional bucket prefix (e.g., `"aj-test100"`). 235 236 Returns: 237 Sorted list of connector technical names (e.g., `"source-github"`). 238 239 Raises: 240 ValueError: If `support_level` and `min_support_level` are both provided. 241 """ 242 has_filters = any([support_level, min_support_level, connector_type, language]) 243 244 if not has_filters: 245 return list_registry_connectors(bucket_name=bucket_name) 246 247 if support_level and min_support_level: 248 raise ValueError( 249 "Cannot specify both `support_level` and `min_support_level`. " 250 "Use `support_level` for an exact match or `min_support_level` for a threshold." 251 ) 252 253 entries = _read_cloud_registry_index(bucket_name=bucket_name, prefix=prefix) 254 255 # Apply support_level exact match 256 if support_level: 257 entries = [e for e in entries if e.get("supportLevel") == support_level] 258 259 # Apply min_support_level threshold 260 if min_support_level: 261 threshold = min_support_level.precedence 262 known_levels = {m.value for m in SupportLevel} 263 entries = [ 264 e 265 for e in entries 266 if e.get("supportLevel") 267 and e["supportLevel"] in known_levels 268 and SupportLevel(e["supportLevel"]).precedence >= threshold 269 ] 270 271 # Apply connector_type filter 272 if connector_type == ConnectorType.SOURCE: 273 entries = [e for e in entries if "sourceDefinitionId" in e] 274 elif connector_type == ConnectorType.DESTINATION: 275 entries = [e for e in entries if "destinationDefinitionId" in e] 276 277 # Apply language filter 278 if language: 279 entries = [e for e in entries if e.get("language") == language] 280 281 # Extract connector names from dockerRepository (e.g., "airbyte/source-github" -> "source-github") 282 names: set[str] = set() 283 for entry in entries: 284 docker_repo = entry.get("dockerRepository", "") 285 if "/" in docker_repo: 286 names.add(docker_repo.split("/", 1)[1]) 287 elif docker_repo: 288 names.add(docker_repo) 289 290 return sorted(names)
List connectors from the compiled cloud registry index with filtering.
When any filter is applied, reads the compiled cloud_registry.json index
instead of globbing individual metadata blobs. This is significantly faster
because the index is a single JSON file containing all connector entries.
When no filters are applied, falls back to the existing glob-based search which captures all connectors (including OSS-only connectors not in the Cloud index).
Arguments:
- bucket_name: Name of the GCS bucket containing the registry.
- support_level: Exact support level to match (e.g.,
SupportLevel.CERTIFIED). - min_support_level: Minimum support level threshold. Returns connectors at or above this level.
- connector_type: Filter by connector type (
ConnectorType.SOURCEorConnectorType.DESTINATION). - language: Filter by implementation language (e.g.,
ConnectorLanguage.PYTHON). - prefix: Optional bucket prefix (e.g.,
"aj-test100").
Returns:
Sorted list of connector technical names (e.g.,
"source-github").
Raises:
- ValueError: If
support_levelandmin_support_levelare both provided.
97def publish_connector_metadata( 98 connector_name: str, 99 metadata: dict[str, Any], 100 bucket_name: str, 101 version: str, 102 update_latest: bool = True, 103 dry_run: bool = False, 104) -> MetadataPublishResult: 105 """Publish connector metadata to GCS. 106 107 Uploads the metadata to the registry bucket at a versioned path, and optionally 108 also updates the 'latest' pointer. Uses MD5 hash comparison to avoid re-uploading 109 unchanged files. 110 111 Requires GCS_CREDENTIALS environment variable to be set. 112 """ 113 if not isinstance(metadata, dict): 114 raise ValueError("Metadata must be a dictionary") 115 116 if "data" not in metadata: 117 raise ValueError("Metadata must contain 'data' field") 118 119 # Construct GCS paths using airbyte/{connector_name} convention 120 versioned_blob_path = get_gcs_publish_path(connector_name, "metadata", version) 121 latest_blob_path = get_gcs_publish_path( 122 connector_name, "metadata", LATEST_GCS_FOLDER_NAME 123 ) 124 125 if dry_run: 126 message = f"[DRY RUN] Would upload metadata to gs://{bucket_name}/{versioned_blob_path}" 127 if update_latest: 128 message += f" and gs://{bucket_name}/{latest_blob_path}" 129 logger.info(message) 130 return MetadataPublishResult( 131 connector_name=connector_name, 132 version=version, 133 bucket_name=bucket_name, 134 versioned_path=versioned_blob_path, 135 latest_path=latest_blob_path if update_latest else None, 136 versioned_uploaded=False, 137 latest_uploaded=False, 138 status="dry-run", 139 message=message, 140 ) 141 142 # Get GCS client and bucket 143 storage_client = get_gcs_storage_client() 144 bucket = storage_client.bucket(bucket_name) 145 146 # Write metadata to temp file 147 with tempfile.NamedTemporaryFile( 148 mode="w", suffix=".yaml", delete=False 149 ) as tmp_file: 150 yaml.dump(metadata, tmp_file) 151 tmp_path = Path(tmp_file.name) 152 153 try: 154 # Upload versioned file 155 versioned_uploaded, _ = upload_file_if_changed( 156 local_file_path=tmp_path, 157 bucket=bucket, 158 blob_path=versioned_blob_path, 159 disable_cache=True, 160 ) 161 162 if versioned_uploaded: 163 logger.info( 164 f"Uploaded metadata for {connector_name} v{version} to {versioned_blob_path}" 165 ) 166 else: 167 logger.info( 168 f"Versioned metadata for {connector_name} v{version} is already up to date" 169 ) 170 171 # Optionally update latest pointer 172 latest_uploaded = False 173 if update_latest: 174 latest_uploaded, _ = upload_file_if_changed( 175 local_file_path=tmp_path, 176 bucket=bucket, 177 blob_path=latest_blob_path, 178 disable_cache=True, 179 ) 180 if latest_uploaded: 181 logger.info(f"Updated latest pointer for {connector_name}") 182 else: 183 logger.info( 184 f"Latest pointer for {connector_name} is already up to date" 185 ) 186 finally: 187 # Clean up temp file even if upload fails 188 tmp_path.unlink(missing_ok=True) 189 190 # Determine status 191 if versioned_uploaded or latest_uploaded: 192 status = "success" 193 message = f"Published metadata for {connector_name} v{version}" 194 if versioned_uploaded: 195 message += f" to {versioned_blob_path}" 196 if latest_uploaded: 197 message += " and updated latest" 198 else: 199 status = "already-up-to-date" 200 message = f"Metadata for {connector_name} v{version} is already up to date" 201 202 return MetadataPublishResult( 203 connector_name=connector_name, 204 version=version, 205 bucket_name=bucket_name, 206 versioned_path=versioned_blob_path, 207 latest_path=latest_blob_path if update_latest else None, 208 versioned_uploaded=versioned_uploaded, 209 latest_uploaded=latest_uploaded, 210 status=status, 211 message=message, 212 )
Publish connector metadata to GCS.
Uploads the metadata to the registry bucket at a versioned path, and optionally also updates the 'latest' pointer. Uses MD5 hash comparison to avoid re-uploading unchanged files.
Requires GCS_CREDENTIALS environment variable to be set.
212def publish_version_artifacts( 213 connector_name: str, 214 version: str, 215 artifacts_dir: Path, 216 store: RegistryStore, 217 dry_run: bool = False, 218 with_validate: bool = True, 219) -> PublishArtifactsResult: 220 """Publish locally generated artifacts to a GCS registry bucket. 221 222 Uses `gcsfs.GCSFileSystem` to upload the local *artifacts_dir* to the 223 versioned path inside the target GCS bucket. 224 225 The target GCS path is: 226 `gs://<bucket>/[<prefix>/]metadata/airbyte/<connector>/<version>/` 227 228 Before uploading, this function validates that the `connector_name` 229 (derived from the connector directory) matches the `dockerRepository` 230 declared in `metadata.yaml`. A mismatch would cause the registry 231 compile step to see duplicate definition-ID entries and fail. 232 233 Args: 234 connector_name: Connector name (e.g. `source-faker`). 235 version: Version string (e.g. `6.2.38`). 236 artifacts_dir: Local directory containing artifacts from `generate`. 237 store: Parsed store target containing bucket, prefix, and stage info. 238 dry_run: If `True`, report what would be uploaded without writing. 239 with_validate: If `True` (default), validate metadata before uploading. 240 Pass `False` (`--no-validate`) to skip. 241 242 Returns: 243 A `PublishArtifactsResult` describing what was published. 244 245 Raises: 246 ValueError: If the connector directory name does not match 247 `dockerRepository` in the generated metadata. 248 """ 249 if not artifacts_dir.is_dir(): 250 raise FileNotFoundError(f"Artifacts directory not found: {artifacts_dir}") 251 252 # Fail fast if the connector directory name doesn't match dockerRepository. 253 # A mismatch would publish artifacts under the wrong GCS path and corrupt 254 # the registry (duplicate definition-IDs under different directory names). 255 mismatch_error = _check_connector_name_matches_docker_repo( 256 connector_name, artifacts_dir 257 ) 258 if mismatch_error: 259 raise ValueError(mismatch_error) 260 261 # Build the GCS destination path 262 bucket_name = store.bucket 263 prefix = store.prefix 264 blob_root = versioned_blob_root( 265 connector_name=connector_name, version=version, store=store 266 ) 267 versioned_dest = f"gcs://{bucket_name}/{blob_root}" 268 269 target_label = f"{bucket_name}/{prefix}" if prefix else bucket_name 270 progressive_rollout_enabled = _metadata_enables_progressive_rollout(artifacts_dir) 271 rollout_overridden_by_breaking_change = ( 272 progressive_rollout_enabled 273 and _metadata_declares_breaking_change(artifacts_dir, version) 274 ) 275 published_latest_version = ( 276 _published_latest_version(connector_name, store) 277 if progressive_rollout_enabled and not rollout_overridden_by_breaking_change 278 else None 279 ) 280 rollout_overridden_by_published_ga = ( 281 published_latest_version == version 282 if published_latest_version is not None 283 else False 284 ) 285 if rollout_overridden_by_breaking_change: 286 logger.info( 287 "Ignoring progressive rollout settings for %s@%s; version is declared " 288 "as a breaking change. The rollout marker will not be published.", 289 connector_name, 290 version, 291 ) 292 elif rollout_overridden_by_published_ga: 293 logger.info( 294 "Ignoring progressive rollout settings for %s@%s; version is already " 295 "the published Default GA. The rollout marker will not be published.", 296 connector_name, 297 version, 298 ) 299 should_publish_rollout_marker = ( 300 progressive_rollout_enabled 301 and not rollout_overridden_by_breaking_change 302 and not rollout_overridden_by_published_ga 303 ) 304 result = PublishArtifactsResult( 305 connector_name=connector_name, 306 version=version, 307 target=target_label, 308 gcs_destination=versioned_dest, 309 progressive_rollout_overridden_by_breaking_change=rollout_overridden_by_breaking_change, 310 progressive_rollout_overridden_by_published_ga=rollout_overridden_by_published_ga, 311 dry_run=dry_run, 312 ) 313 314 # --- Pre-publish validation --- 315 if with_validate: 316 metadata_file = artifacts_dir / "metadata.yaml" 317 if metadata_file.is_file(): 318 raw_metadata = yaml.safe_load(metadata_file.read_text()) 319 metadata_data = (raw_metadata or {}).get("data", {}) 320 validation = validate_metadata(metadata_data=metadata_data) 321 if not validation.passed: 322 for err in validation.errors: 323 logger.error("Pre-publish validation error: %s", err) 324 result.validation_errors = validation.errors 325 return result 326 logger.info( 327 "Pre-publish validation passed (%d validators).", 328 validation.validators_run, 329 ) 330 else: 331 logger.warning("No metadata.yaml in artifacts dir; skipping validation.") 332 333 # Enumerate local files 334 local_files = sorted(f for f in artifacts_dir.rglob("*") if f.is_file()) 335 if not local_files: 336 result.errors.append(f"No files found in {artifacts_dir}.") 337 return result 338 339 _log_progress( 340 "Publishing %d artifacts for %s@%s → %s", 341 len(local_files), 342 connector_name, 343 version, 344 versioned_dest, 345 ) 346 347 # Build references used by both dry-run and real upload paths 348 deps_file = artifacts_dir / CONNECTOR_DEPENDENCY_FILE_NAME 349 has_deps = deps_file.is_file() 350 deps_gcs_key = dependencies_blob_path( 351 connector_name=connector_name, version=version, store=store 352 ) 353 354 sbom_file = artifacts_dir / SBOM_FILE_NAME 355 has_sbom = sbom_file.is_file() 356 rollout_marker_path = f"{blob_root}/{PROGRESSIVE_ROLLOUT_MARKER_FILE}" 357 358 if dry_run: 359 for f in local_files: 360 rel = f.relative_to(artifacts_dir) 361 result.files_uploaded.append(str(rel)) 362 _log_progress(" [DRY RUN] would upload: %s", rel) 363 # Report the dual-load of dependencies.json to connector_dependencies/ 364 if has_deps: 365 result.files_uploaded.append(deps_gcs_key) 366 _log_progress( 367 " [DRY RUN] would also dual-load: %s → gs://%s/%s", 368 CONNECTOR_DEPENDENCY_FILE_NAME, 369 bucket_name, 370 deps_gcs_key, 371 ) 372 # Report the separate sbom/ upload 373 if has_sbom: 374 sbom_gcs_key = sbom_blob_path( 375 connector_name=connector_name, 376 version=version, 377 store=store, 378 ) 379 result.files_uploaded.append(sbom_gcs_key) 380 _log_progress( 381 " [DRY RUN] would also upload: %s → gs://%s/%s", 382 SBOM_FILE_NAME, 383 bucket_name, 384 sbom_gcs_key, 385 ) 386 if should_publish_rollout_marker: 387 result.files_uploaded.append(PROGRESSIVE_ROLLOUT_MARKER_FILE) 388 _log_progress( 389 " [DRY RUN] would write active marker: gs://%s/%s", 390 bucket_name, 391 rollout_marker_path, 392 ) 393 return result 394 395 # Authenticate 396 token = get_gcs_credentials_token() 397 fs = gcsfs.GCSFileSystem(token=token) 398 399 # Strip gcs:// prefix for gcsfs path 400 dest_path = versioned_dest.replace("gcs://", "") 401 402 # Upload all files to the versioned path 403 _log_progress("Uploading to: %s", versioned_dest) 404 for f in local_files: 405 rel = f.relative_to(artifacts_dir) 406 remote_path = f"{dest_path}/{rel}" 407 fs.put(str(f), remote_path) 408 result.files_uploaded.append(str(rel)) 409 _log_progress(" Uploaded: %s", rel) 410 411 # Delete remote files that don't exist locally (sync semantics) 412 try: 413 remote_files = fs.ls(dest_path, detail=False) 414 local_rel_paths = {str(f.relative_to(artifacts_dir)) for f in local_files} 415 for remote_file in remote_files: 416 # Skip the directory entry itself if it appears in the listing 417 if remote_file == dest_path: 418 continue 419 # Derive the remote relative path, matching upload semantics 420 if remote_file.startswith(dest_path + "/"): 421 remote_rel = remote_file[len(dest_path) + 1 :] 422 else: 423 remote_rel = remote_file.split("/")[-1] 424 if is_registry_state_marker_file(Path(remote_rel).name): 425 continue 426 if remote_rel not in local_rel_paths: 427 fs.rm(remote_file) 428 _log_progress(" Deleted stale remote file: %s", remote_rel) 429 except FileNotFoundError: 430 pass # Destination doesn't exist yet, nothing to clean 431 432 _log_progress("Uploaded %d files to %s", len(local_files), versioned_dest) 433 434 if should_publish_rollout_marker: 435 marker_remote = f"{bucket_name}/{rollout_marker_path}" 436 with fs.open(marker_remote, "w") as marker_file: 437 marker_file.write(_progressive_rollout_marker_content()) 438 result.files_uploaded.append(PROGRESSIVE_ROLLOUT_MARKER_FILE) 439 _log_progress("Wrote active marker: gs://%s", marker_remote) 440 441 # --- Dual-load dependencies.json to the connector_dependencies/ path --- 442 if not has_deps: 443 logger.debug( 444 "No %s in artifacts dir — skipping dual-load.", 445 CONNECTOR_DEPENDENCY_FILE_NAME, 446 ) 447 else: 448 deps_remote = f"{bucket_name}/{deps_gcs_key}" 449 _log_progress( 450 "Dual-loading %s to gs://%s", 451 CONNECTOR_DEPENDENCY_FILE_NAME, 452 deps_remote, 453 ) 454 fs.put(str(deps_file), deps_remote) 455 result.files_uploaded.append(deps_gcs_key) 456 _log_progress(" Uploaded %s (dual-load)", CONNECTOR_DEPENDENCY_FILE_NAME) 457 458 # --- Upload SBOM to the dedicated sbom/ path in GCS --- 459 if not has_sbom: 460 logger.debug( 461 "No %s in artifacts dir — skipping SBOM dual-load.", 462 SBOM_FILE_NAME, 463 ) 464 else: 465 sbom_gcs_uri = upload_sbom( 466 sbom_path=sbom_file, 467 connector_name=connector_name, 468 version=version, 469 store=store, 470 dry_run=dry_run, 471 ) 472 result.files_uploaded.append( 473 sbom_blob_path( 474 connector_name=connector_name, 475 version=version, 476 store=store, 477 ), 478 ) 479 _log_progress("Uploaded SBOM to dedicated path: %s", sbom_gcs_uri) 480 481 return result
Publish locally generated artifacts to a GCS registry bucket.
Uses gcsfs.GCSFileSystem to upload the local artifacts_dir to the
versioned path inside the target GCS bucket.
The target GCS path is:
gs://<bucket>/[<prefix>/]metadata/airbyte/<connector>/<version>/
Before uploading, this function validates that the connector_name
(derived from the connector directory) matches the dockerRepository
declared in metadata.yaml. A mismatch would cause the registry
compile step to see duplicate definition-ID entries and fail.
Arguments:
- connector_name: Connector name (e.g.
source-faker). - version: Version string (e.g.
6.2.38). - artifacts_dir: Local directory containing artifacts from
generate. - store: Parsed store target containing bucket, prefix, and stage info.
- dry_run: If
True, report what would be uploaded without writing. - with_validate: If
True(default), validate metadata before uploading. PassFalse(--no-validate) to skip.
Returns:
A
PublishArtifactsResultdescribing what was published.
Raises:
- ValueError: If the connector directory name does not match
dockerRepositoryin the generated metadata.
1749def purge_latest_dirs( 1750 *, 1751 store: RegistryStore, 1752 connector_name: list[str] | None = None, 1753 dry_run: bool = False, 1754) -> PurgeLatestResult: 1755 """Delete all `latest/` directories from the registry store. 1756 1757 Discovers connector directories via glob, then deletes each 1758 `latest/` subdirectory in parallel using a thread pool. 1759 1760 Args: 1761 store: Registry store (bucket + optional prefix). 1762 connector_name: If provided, only purge these connectors. 1763 dry_run: If True, report what would be done without deleting. 1764 1765 Returns: 1766 A `PurgeLatestResult` describing what was done. 1767 """ 1768 result = PurgeLatestResult(target=store.bucket_root, dry_run=dry_run) 1769 1770 token = get_gcs_credentials_token() 1771 fs = gcsfs.GCSFileSystem(token=token) 1772 1773 base = f"{store.bucket_root}/{METADATA_FOLDER}/airbyte" 1774 1775 # Discover latest/ dirs by listing connector directories that contain 1776 # a `latest/` subdirectory. 1777 _log_progress("Discovering latest/ directories...") 1778 base_with_slash = f"{base}/" 1779 if connector_name: 1780 # Check each requested connector for a latest/ dir 1781 seen: set[str] = set() 1782 connectors_with_latest: list[str] = [] 1783 for name in connector_name: 1784 if name in seen: 1785 continue 1786 latest_path = f"{base}/{name}/latest" 1787 if fs.exists(latest_path): 1788 connectors_with_latest.append(name) 1789 seen.add(name) 1790 else: 1791 # Glob for all connectors, then filter to those with latest/ 1792 all_connector_dirs = fs.glob(f"{base}/*/latest") 1793 seen = set() 1794 connectors_with_latest = [] 1795 for path in all_connector_dirs: 1796 # Strip the known base prefix and take the first component 1797 if not path.startswith(base_with_slash): 1798 logger.warning("Could not parse latest path: %s", path) 1799 continue 1800 relative = path[len(base_with_slash) :] 1801 connector = relative.split("/")[0] 1802 if connector and connector not in seen: 1803 connectors_with_latest.append(connector) 1804 seen.add(connector) 1805 1806 result.connectors_found = len(connectors_with_latest) 1807 _log_progress( 1808 "Found %d connectors with latest/ directories", 1809 result.connectors_found, 1810 ) 1811 1812 if not connectors_with_latest: 1813 _log_progress("Nothing to purge.") 1814 _log_progress(result.summary()) 1815 return result 1816 1817 if dry_run: 1818 for connector in sorted(connectors_with_latest): 1819 _log_progress(" [DRY RUN] Would delete %s/latest/", connector) 1820 result.latest_dirs_deleted = len(connectors_with_latest) 1821 _log_progress(result.summary()) 1822 return result 1823 1824 # Delete latest/ dirs in parallel using the shared helper. 1825 def _delete_one(connector: str) -> str | None: 1826 """Delete a single connector's latest/ dir. Returns error string or None.""" 1827 try: 1828 _delete_latest_dir( 1829 fs, 1830 store=store, 1831 connector=connector, 1832 ) 1833 return None 1834 except Exception as exc: 1835 return f"Failed to delete latest/ for {connector}: {exc}" 1836 1837 _log_progress( 1838 "Deleting %d latest/ directories (max_workers=%d)...", 1839 len(connectors_with_latest), 1840 _PURGE_LATEST_MAX_WORKERS, 1841 ) 1842 1843 with ThreadPoolExecutor(max_workers=_PURGE_LATEST_MAX_WORKERS) as pool: 1844 futures = { 1845 pool.submit(_delete_one, c): c for c in sorted(connectors_with_latest) 1846 } 1847 for i, future in enumerate(as_completed(futures), 1): 1848 connector = futures[future] 1849 error = future.result() 1850 if error: 1851 logger.error(error) 1852 result.errors.append(error) 1853 else: 1854 result.latest_dirs_deleted += 1 1855 if i % 100 == 0: 1856 _log_progress(" Deleted %d / %d...", i, len(connectors_with_latest)) 1857 1858 _log_progress(result.summary()) 1859 return result
Delete all latest/ directories from the registry store.
Discovers connector directories via glob, then deletes each
latest/ subdirectory in parallel using a thread pool.
Arguments:
- store: Registry store (bucket + optional prefix).
- connector_name: If provided, only purge these connectors.
- dry_run: If True, report what would be done without deleting.
Returns:
A
PurgeLatestResultdescribing what was done.
147def rebuild_registry( 148 source_bucket: str, 149 output_mode: OutputMode, 150 output_path_root: str | None = None, 151 gcs_bucket: str | None = None, 152 s3_bucket: str | None = None, 153 dry_run: bool = False, 154 connector_name: list[str] | None = None, 155) -> RebuildResult: 156 """Rebuild the entire registry from a source GCS bucket to an output target. 157 158 Reads all connector metadata blobs from the source GCS bucket and copies 159 them to the output target using fsspec for unified filesystem access. 160 161 The output targets are: 162 - local: Write to a local directory tree. 163 - gcs: Copy to a GCS bucket (must not be the prod bucket). 164 - s3: Copy to an S3 bucket. 165 166 Args: 167 source_bucket: The GCS bucket to read from (typically prod). 168 output_mode: Where to write: "local", "gcs", or "s3". 169 output_path_root: Root path/prefix for output. For local mode, if None 170 creates a temp directory. For GCS/S3, prepended to all blob paths. 171 gcs_bucket: Target GCS bucket name (required if output_mode="gcs"). 172 s3_bucket: Target S3 bucket name (required if output_mode="s3"). 173 dry_run: If True, report what would be done without writing. 174 connector_name: If provided, only rebuild these connector names 175 (e.g. ["source-faker", "destination-bigquery"]). If None, rebuilds all. 176 177 Returns: 178 RebuildResult with details of the operation. 179 180 Raises: 181 ValueError: If target is the prod bucket, or required bucket arg is missing. 182 """ 183 if output_mode == "gcs" and gcs_bucket: 184 _validate_not_prod_bucket(gcs_bucket) 185 186 # Resolve output root for local mode 187 effective_output_root = output_path_root or "" 188 if output_mode == "local": 189 effective_output_root = _resolve_local_output_root(output_path_root) 190 191 result = RebuildResult( 192 source_bucket=source_bucket, 193 output_mode=output_mode, 194 output_root=effective_output_root, 195 dry_run=dry_run, 196 ) 197 198 # Create source filesystem (always GCS) 199 source_fs, gcs_token = _make_source_fs() 200 source_base = f"{source_bucket}/{METADATA_FOLDER}" 201 202 # List files under metadata/ in the source bucket 203 if connector_name: 204 _log_progress( 205 "Listing blobs for %d connectors under gs://%s/...", 206 len(connector_name), 207 source_base, 208 ) 209 source_paths: list[str] = [] 210 for name in connector_name: 211 connector_prefix = f"{source_base}/airbyte/{name}" 212 found = source_fs.find(connector_prefix) 213 source_paths.extend(found) 214 _log_progress(" %s: %d blobs", name, len(found)) 215 else: 216 _log_progress("Listing all blobs under gs://%s/...", source_base) 217 source_paths = source_fs.find(source_base) 218 219 if not source_paths: 220 _log_progress("No blobs found under gs://%s/", source_base) 221 return result 222 223 total_blobs = len(source_paths) 224 _log_progress("Found %d blobs to process", total_blobs) 225 226 # Create output filesystem 227 output_fs, output_base = _make_output_fs( 228 output_mode=output_mode, 229 output_root=effective_output_root, 230 gcs_bucket=gcs_bucket, 231 s3_bucket=s3_bucket, 232 gcs_token=gcs_token, 233 ) 234 235 # Collect connector names and compute relative paths for all blobs. 236 bucket_prefix = f"{source_bucket}/" 237 blob_relative_paths: list[str] = [] 238 connector_names: set[str] = set() 239 240 for source_path in source_paths: 241 relative_path = source_path 242 if source_path.startswith(bucket_prefix): 243 relative_path = source_path[len(bucket_prefix) :] 244 blob_relative_paths.append(relative_path) 245 246 parts = relative_path.split("/") 247 if len(parts) >= 3: 248 connector_names.add(parts[2]) 249 250 if dry_run: 251 result.blobs_copied = total_blobs 252 result.connectors_processed = len(connector_names) 253 _log_progress( 254 "[DRY RUN] Would copy %d blobs (%d connectors)", 255 total_blobs, 256 len(connector_names), 257 ) 258 _log_progress(result.summary()) 259 return result 260 261 # Use GCS-native server-side copy for GCS→GCS mirrors. 262 if output_mode == "gcs" and gcs_bucket: 263 result = _gcs_native_copy( 264 source_bucket=source_bucket, 265 dest_bucket_name=gcs_bucket, 266 dest_prefix=effective_output_root, 267 blob_relative_paths=blob_relative_paths, 268 connector_names=connector_names, 269 gcs_token=gcs_token, 270 result=result, 271 connector_name_filter=connector_name, 272 ) 273 return result 274 275 # Fallback: fsspec-based copy for local and S3 output modes. 276 result = _fsspec_copy( 277 source_fs=source_fs, 278 source_paths=source_paths, 279 output_fs=output_fs, 280 output_base=output_base, 281 output_mode=output_mode, 282 source_bucket=source_bucket, 283 blob_relative_paths=blob_relative_paths, 284 connector_names=connector_names, 285 result=result, 286 ) 287 return result
Rebuild the entire registry from a source GCS bucket to an output target.
Reads all connector metadata blobs from the source GCS bucket and copies them to the output target using fsspec for unified filesystem access.
The output targets are:
- local: Write to a local directory tree.
- gcs: Copy to a GCS bucket (must not be the prod bucket).
- s3: Copy to an S3 bucket.
Arguments:
- source_bucket: The GCS bucket to read from (typically prod).
- output_mode: Where to write: "local", "gcs", or "s3".
- output_path_root: Root path/prefix for output. For local mode, if None creates a temp directory. For GCS/S3, prepended to all blob paths.
- gcs_bucket: Target GCS bucket name (required if output_mode="gcs").
- s3_bucket: Target S3 bucket name (required if output_mode="s3").
- dry_run: If True, report what would be done without writing.
- connector_name: If provided, only rebuild these connector names (e.g. ["source-faker", "destination-bigquery"]). If None, rebuilds all.
Returns:
RebuildResult with details of the operation.
Raises:
- ValueError: If target is the prod bucket, or required bucket arg is missing.
263def resolve_registry_store( 264 store: str | None = None, 265 connector_name: str | None = None, 266 cwd: Path | None = None, 267 default_env: str = "dev", 268) -> RegistryStore: 269 """Resolve a `RegistryStore` from CLI inputs. 270 271 All applicable detection methods are evaluated. Explicit sources 272 (`--store`, then the `AIRBYTE_REGISTRY_STORE` env var) take priority 273 and are returned directly. When only auto-detected sources remain, they 274 are compared and a `ValueError` is raised if they disagree. 275 276 Priority (highest → lowest): 277 278 1. **Explicit** `--store` argument (e.g. `"coral:dev"`). 279 2. **Environment variable** -- `AIRBYTE_REGISTRY_STORE`. 280 3. **Auto-detected** -- connector name and/or working directory. 281 If both are present and disagree, a `ValueError` is raised. 282 283 Args: 284 store: Explicit store target string (e.g. `"coral:dev"`). 285 connector_name: Optional connector name for auto-detection. 286 cwd: Working directory for repo-based detection. 287 default_env: Environment to use when auto-detecting (default `"dev"`). 288 289 Returns: 290 A fully resolved `RegistryStore`. 291 292 Raises: 293 ValueError: If no detection method succeeds, or if auto-detected 294 methods produce conflicting store types. 295 """ 296 # -- Collect all detection results ------------------------------------ 297 # Explicit sources (take priority — no conflict checking needed). 298 explicit_target: RegistryStore | None = None 299 explicit_source: str | None = None 300 301 if store is not None: 302 explicit_target = RegistryStore.parse(store) 303 explicit_source = "--store" 304 305 env_store = os.environ.get(REGISTRY_STORE_ENV_VAR) 306 if env_store and explicit_target is None: 307 explicit_target = RegistryStore.parse(env_store) 308 explicit_source = REGISTRY_STORE_ENV_VAR 309 310 # Auto-detected sources (only consulted when no explicit source). 311 auto_detections: dict[str, StoreType] = {} 312 313 if connector_name is not None: 314 auto_detections["connector_name"] = StoreType.get_from_connector_name( 315 connector_name, 316 ) 317 318 dir_type = StoreType.detect_from_repo_dir(cwd) 319 if dir_type is not None: 320 auto_detections["working_directory"] = dir_type 321 322 # -- Return explicit source if present -------------------------------- 323 if explicit_target is not None: 324 logger.debug( 325 "Using explicit store target from %s: %s:%s", 326 explicit_source, 327 explicit_target.store_type.value, 328 explicit_target.env, 329 ) 330 return explicit_target 331 332 # -- No explicit source: resolve from auto-detections ----------------- 333 if not auto_detections: 334 raise ValueError( 335 "Cannot determine registry store. " 336 "Provide --store (e.g. 'coral:dev' or 'sonar:prod'), " 337 f"set ${REGISTRY_STORE_ENV_VAR}, " 338 "or run from a recognized repository directory." 339 ) 340 341 distinct = set(auto_detections.values()) 342 343 if len(distinct) > 1: 344 detail = ", ".join(f"{src}={st.value}" for src, st in auto_detections.items()) 345 raise ValueError( 346 f"Conflicting store types detected: {detail}. " 347 "Provide an explicit --store to resolve the ambiguity." 348 ) 349 350 resolved_type = distinct.pop() 351 logger.info( 352 "Auto-detected store type '%s' (sources: %s)", 353 resolved_type.value, 354 ", ".join(auto_detections), 355 ) 356 return RegistryStore(store_type=resolved_type, env=default_env)
Resolve a RegistryStore from CLI inputs.
All applicable detection methods are evaluated. Explicit sources
(--store, then the AIRBYTE_REGISTRY_STORE env var) take priority
and are returned directly. When only auto-detected sources remain, they
are compared and a ValueError is raised if they disagree.
Priority (highest → lowest):
- Explicit
--storeargument (e.g."coral:dev"). - Environment variable --
AIRBYTE_REGISTRY_STORE. - Auto-detected -- connector name and/or working directory.
If both are present and disagree, a
ValueErroris raised.
Arguments:
- store: Explicit store target string (e.g.
"coral:dev"). - connector_name: Optional connector name for auto-detection.
- cwd: Working directory for repo-based detection.
- default_env: Environment to use when auto-detecting (default
"dev").
Returns:
A fully resolved
RegistryStore.
Raises:
- ValueError: If no detection method succeeds, or if auto-detected methods produce conflicting store types.
322def unyank_connector_version( 323 connector_name: str, 324 version: str, 325 bucket_name: str, 326 dry_run: bool = False, 327) -> YankResult: 328 """Rename the active yank marker to an unyanked audit marker. 329 330 Moves the active version-yank.yml marker at: 331 metadata/airbyte/{connector_name}/{version}/version-yank.yml 332 to: 333 metadata/airbyte/{connector_name}/{version}/version-unyanked-yyyymmdd.yml 334 335 Args: 336 connector_name: The connector name (e.g., "source-faker"). 337 version: The version to unyank (e.g., "1.2.3"). 338 bucket_name: The GCS bucket name. 339 dry_run: If True, report what would be done without writing. 340 341 Returns: 342 YankResult with details of the operation. 343 """ 344 yank_path = _get_yank_blob_path(connector_name, version) 345 346 storage_client = get_gcs_storage_client() 347 bucket = storage_client.bucket(bucket_name) 348 349 # Check if yank marker exists 350 yank_blob = bucket.blob(yank_path) 351 if not yank_blob.exists(): 352 return YankResult( 353 connector_name=connector_name, 354 version=version, 355 bucket_name=bucket_name, 356 action="unyank", 357 success=False, 358 message=f"Version {version} of {connector_name} is not yanked.", 359 dry_run=dry_run, 360 ) 361 362 if dry_run: 363 return YankResult( 364 connector_name=connector_name, 365 version=version, 366 bucket_name=bucket_name, 367 action="unyank", 368 success=True, 369 message=f"[DRY RUN] Would unyank {connector_name} {version}.", 370 dry_run=True, 371 ) 372 373 unyanked_path = _get_yank_blob_path(connector_name, version).replace( 374 YANK_FILE_NAME, 375 unyanked_marker_file(), 376 ) 377 bucket.copy_blob(yank_blob, bucket, new_name=unyanked_path) 378 yank_blob.delete() 379 380 logger.info("Unyanked %s version %s in %s", connector_name, version, bucket_name) 381 382 return YankResult( 383 connector_name=connector_name, 384 version=version, 385 bucket_name=bucket_name, 386 action="unyank", 387 success=True, 388 message=f"Successfully unyanked {connector_name} {version}.", 389 )
Rename the active yank marker to an unyanked audit marker.
Moves the active version-yank.yml marker at: metadata/airbyte/{connector_name}/{version}/version-yank.yml to: metadata/airbyte/{connector_name}/{version}/version-unyanked-yyyymmdd.yml
Arguments:
- connector_name: The connector name (e.g., "source-faker").
- version: The version to unyank (e.g., "1.2.3").
- bucket_name: The GCS bucket name.
- dry_run: If True, report what would be done without writing.
Returns:
YankResult with details of the operation.
270def validate_metadata( 271 metadata_data: dict[str, Any], 272 opts: ValidateOptions | None = None, 273) -> ValidationResult: 274 """Run all pre-publish validators against raw `metadata.data`. 275 276 Args: 277 metadata_data: The `data` section of a parsed `metadata.yaml`. 278 opts: Options influencing validation behaviour. 279 280 Returns: 281 A `ValidationResult` with aggregate pass/fail and error list. 282 """ 283 if opts is None: 284 opts = ValidateOptions() 285 286 result = ValidationResult() 287 288 for validator in PRE_PUBLISH_VALIDATORS: 289 result.validators_run += 1 290 logger.info("Running validator: %s", validator.__name__) 291 passed, error = validator(metadata_data, opts) 292 if not passed and error: 293 logger.error("Validation failed: %s", error) 294 result.add_error(error) 295 296 return result
Run all pre-publish validators against raw metadata.data.
Arguments:
- metadata_data: The
datasection of a parsedmetadata.yaml. - opts: Options influencing validation behaviour.
Returns:
A
ValidationResultwith aggregate pass/fail and error list.
226def yank_connector_version( 227 connector_name: str, 228 version: str, 229 bucket_name: str, 230 reason: str = "", 231 approval_url: str = "", 232 dry_run: bool = False, 233) -> YankResult: 234 """Mark a connector version as yanked by writing a version-yank.yml marker. 235 236 The marker file is placed at: 237 metadata/airbyte/{connector_name}/{version}/version-yank.yml 238 239 Args: 240 connector_name: The connector name (e.g., "source-faker"). 241 version: The version to yank (e.g., "1.2.3"). 242 bucket_name: The GCS bucket name. 243 reason: Optional reason for yanking the version. 244 approval_url: Optional approval evidence URL to record in the marker. 245 dry_run: If True, report what would be done without writing. 246 247 Returns: 248 YankResult with details of the operation. 249 250 Raises: 251 ValueError: If the bucket is the production bucket and no override is set, 252 or if the version does not exist. 253 """ 254 yank_path = _get_yank_blob_path(connector_name, version) 255 metadata_path = _get_metadata_blob_path(connector_name, version) 256 257 storage_client = get_gcs_storage_client() 258 bucket = storage_client.bucket(bucket_name) 259 260 # Verify the version exists 261 metadata_blob = bucket.blob(metadata_path) 262 if not metadata_blob.exists(): 263 return YankResult( 264 connector_name=connector_name, 265 version=version, 266 bucket_name=bucket_name, 267 action="yank", 268 success=False, 269 message=f"Version {version} not found for {connector_name} in {bucket_name}.", 270 dry_run=dry_run, 271 ) 272 273 # Check if already yanked 274 yank_blob = bucket.blob(yank_path) 275 if yank_blob.exists(): 276 return YankResult( 277 connector_name=connector_name, 278 version=version, 279 bucket_name=bucket_name, 280 action="yank", 281 success=False, 282 message=f"Version {version} of {connector_name} is already yanked.", 283 dry_run=dry_run, 284 ) 285 286 if dry_run: 287 return YankResult( 288 connector_name=connector_name, 289 version=version, 290 bucket_name=bucket_name, 291 action="yank", 292 success=True, 293 message=f"[DRY RUN] Would yank {connector_name} {version}.", 294 dry_run=True, 295 ) 296 297 # Write the yank marker file 298 yank_content: dict[str, Any] = { 299 "yanked": True, 300 "yanked_at": datetime.now(tz=timezone.utc).isoformat(), 301 } 302 if reason: 303 yank_content["reason"] = reason 304 if approval_url: 305 yank_content["approval_url"] = approval_url 306 307 yank_yaml = yaml.dump(yank_content, default_flow_style=False) 308 yank_blob.upload_from_string(yank_yaml, content_type="application/x-yaml") 309 310 logger.info("Yanked %s version %s in %s", connector_name, version, bucket_name) 311 312 return YankResult( 313 connector_name=connector_name, 314 version=version, 315 bucket_name=bucket_name, 316 action="yank", 317 success=True, 318 message=f"Successfully yanked {connector_name} {version}.", 319 )
Mark a connector version as yanked by writing a version-yank.yml marker.
The marker file is placed at:
metadata/airbyte/{connector_name}/{version}/version-yank.yml
Arguments:
- connector_name: The connector name (e.g., "source-faker").
- version: The version to yank (e.g., "1.2.3").
- bucket_name: The GCS bucket name.
- reason: Optional reason for yanking the version.
- approval_url: Optional approval evidence URL to record in the marker.
- dry_run: If True, report what would be done without writing.
Returns:
YankResult with details of the operation.
Raises:
- ValueError: If the bucket is the production bucket and no override is set, or if the version does not exist.