From ea62b3611bc5cabe3c845e8903a4cc0d17570ee9 Mon Sep 17 00:00:00 2001 From: Tommy Healy Date: Thu, 1 Oct 2026 17:29:20 +0200 Subject: [PATCH] Add GCS sink destination support to define_sink_pipe Matches the Forward CLI's gcs_hmac export service: GCSSinkConfig mirrors S3SinkConfig's fields (bucket_uri/file_template/format/ schedule/strategy/compression) but references a GCSConnectionDefinition. Since GCS and S3 sinks share identical EXPORT_* directives, the generator emits an explicit EXPORT_SERVICE gcs_hmac line to disambiguate on reverse parse, and the migrate runner's sink/connection compatibility check now maps gcs_hmac sinks to "gcs" connections. Co-Authored-By: Claude Sonnet 5 --- CHANGELOG.md | 6 + README.md | 26 ++- src/tinybird_sdk/generator/pipe.py | 13 ++ src/tinybird_sdk/migrate/__init__.py | 1 + src/tinybird_sdk/migrate/emit_ts.py | 3 +- src/tinybird_sdk/migrate/parse_pipe.py | 53 +++-- src/tinybird_sdk/migrate/runner.py | 5 +- src/tinybird_sdk/migrate/types.py | 14 +- src/tinybird_sdk/schema/__init__.py | 1 + src/tinybird_sdk/schema/pipe.py | 27 ++- tests/test_gcs_sink_destination.py | 262 +++++++++++++++++++++++++ 11 files changed, 389 insertions(+), 22 deletions(-) create mode 100644 tests/test_gcs_sink_destination.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 0c1ecb0..98bc061 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,12 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project follows [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Added + +- `define_sink_pipe` now supports GCS as a sink/export destination, matching the Forward CLI's `gcs_hmac` export service. Add a `GCSConnectionDefinition` as the sink's `connection` and the same `bucket_uri`/`file_template`/`format`/`schedule`/`strategy`/`compression` options already used for S3 sinks; the datafile emitter writes an explicit `EXPORT_SERVICE gcs_hmac` directive (since GCS and S3 sinks otherwise share identical `EXPORT_*` directives) and the migration parser/emitter round-trip it accordingly. + ## [0.4.0] - 2026-06-29 ### Added diff --git a/README.md b/README.md index e42f34b..321f413 100644 --- a/README.md +++ b/README.md @@ -772,7 +772,7 @@ manual_report = define_copy_pipe( ### Sink Pipes Use sink pipes to publish query results to external systems. -The SDK supports Kafka and S3 sinks. +The SDK supports Kafka, S3, and GCS sinks. ```python from tinybird_sdk import define_sink_pipe, node @@ -820,6 +820,30 @@ s3_events_sink = define_sink_pipe( ], }, ) + +# GCS sink +gcs_events_sink = define_sink_pipe( + "gcs_events_sink", + { + "sink": { + "connection": landing_gcs, + "bucket_uri": "gs://my-bucket/exports/", + "file_template": "events_{date}", + "format": "csv", + "schedule": "@once", + "strategy": "create_new", + "compression": "gzip", + }, + "nodes": [ + node( + { + "name": "export", + "sql": "SELECT timestamp, session_id FROM gcs_landing", + } + ) + ], + }, +) ``` ### Static Tokens diff --git a/src/tinybird_sdk/generator/pipe.py b/src/tinybird_sdk/generator/pipe.py index 940856f..1980c61 100644 --- a/src/tinybird_sdk/generator/pipe.py +++ b/src/tinybird_sdk/generator/pipe.py @@ -6,6 +6,7 @@ from ..schema.pipe import ( CopyConfig, EndpointConfig, + GCSSinkConfig, KafkaSinkConfig, MaterializedConfig, PipeDefinition, @@ -83,6 +84,18 @@ def _generate_sink(config: SinkConfig) -> str: parts.append(f"EXPORT_STRATEGY {config.strategy}") if config.compression: parts.append(f"EXPORT_COMPRESSION {config.compression}") + elif isinstance(config, GCSSinkConfig): + # GCS and S3 sinks share the same EXPORT_* directive shape, so EXPORT_SERVICE + # must be emitted explicitly to disambiguate gcs_hmac from s3 on re-parse. + parts.append("EXPORT_SERVICE gcs_hmac") + parts.append(f"EXPORT_BUCKET_URI {config.bucket_uri}") + parts.append(f"EXPORT_FILE_TEMPLATE {config.file_template}") + parts.append(f"EXPORT_SCHEDULE {config.schedule}") + parts.append(f"EXPORT_FORMAT {config.format}") + if config.strategy: + parts.append(f"EXPORT_STRATEGY {config.strategy}") + if config.compression: + parts.append(f"EXPORT_COMPRESSION {config.compression}") return "\n".join(parts) diff --git a/src/tinybird_sdk/migrate/__init__.py b/src/tinybird_sdk/migrate/__init__.py index b35cb58..e1d9610 100644 --- a/src/tinybird_sdk/migrate/__init__.py +++ b/src/tinybird_sdk/migrate/__init__.py @@ -14,6 +14,7 @@ PipeModel, SinkKafkaModel, SinkS3Model, + SinkGCSModel, SinkModel, KafkaConnectionModel, S3ConnectionModel, diff --git a/src/tinybird_sdk/migrate/emit_ts.py b/src/tinybird_sdk/migrate/emit_ts.py index 7b24142..a397618 100644 --- a/src/tinybird_sdk/migrate/emit_ts.py +++ b/src/tinybird_sdk/migrate/emit_ts.py @@ -14,6 +14,7 @@ ParsedResource, PipeModel, S3ConnectionModel, + SinkGCSModel, SinkKafkaModel, SinkS3Model, ) @@ -427,7 +428,7 @@ def _emit_pipe(pipe: PipeModel) -> str: if isinstance(pipe.sink, SinkKafkaModel): lines.append(f" 'topic': {_escape_string(pipe.sink.topic)},") lines.append(f" 'schedule': {_escape_string(pipe.sink.schedule)},") - elif isinstance(pipe.sink, SinkS3Model): + elif isinstance(pipe.sink, (SinkS3Model, SinkGCSModel)): lines.append(f" 'bucket_uri': {_escape_string(pipe.sink.bucket_uri)},") lines.append(f" 'file_template': {_escape_string(pipe.sink.file_template)},") lines.append(f" 'schedule': {_escape_string(pipe.sink.schedule)},") diff --git a/src/tinybird_sdk/migrate/parse_pipe.py b/src/tinybird_sdk/migrate/parse_pipe.py index 52e9201..b11051b 100644 --- a/src/tinybird_sdk/migrate/parse_pipe.py +++ b/src/tinybird_sdk/migrate/parse_pipe.py @@ -19,6 +19,7 @@ PipeTokenModel, PipeTypeModel, ResourceFile, + SinkGCSModel, SinkKafkaModel, SinkModel, SinkS3Model, @@ -600,7 +601,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel: copy_mode = value elif key == "EXPORT_SERVICE": normalized = parse_quoted_value(value).lower() - if normalized not in {"kafka", "s3"}: + if normalized not in {"kafka", "s3", "gcs_hmac"}: raise MigrationParseError( resource.file_path, "pipe", @@ -721,7 +722,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel: ) has_kafka_directives = export_topic is not None - has_s3_directives = any( + has_blob_directives = any( value is not None for value in ( export_bucket_uri, @@ -731,7 +732,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel: ) ) - if has_kafka_directives and has_s3_directives: + if has_kafka_directives and has_blob_directives: raise MigrationParseError( resource.file_path, "pipe", @@ -739,8 +740,11 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel: "Sink pipe cannot mix Kafka and S3 export directives.", ) + # S3 and GCS sinks share the same EXPORT_* directive shape, so when EXPORT_SERVICE + # is omitted we default to "s3" (matching the SDK's historical behavior); GCS + # sinks must set EXPORT_SERVICE gcs_hmac explicitly to be recognized. inferred_service = export_service or ( - "kafka" if has_kafka_directives else "s3" if has_s3_directives else None + "kafka" if has_kafka_directives else "s3" if has_blob_directives else None ) if not inferred_service: raise MigrationParseError( @@ -751,7 +755,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel: ) if inferred_service == "kafka": - if has_s3_directives: + if has_blob_directives: raise MigrationParseError( resource.file_path, "pipe", @@ -793,7 +797,7 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel: topic=export_topic, schedule=export_schedule, ) - else: + elif inferred_service in {"s3", "gcs_hmac"}: if has_kafka_directives: raise MigrationParseError( resource.file_path, @@ -814,15 +818,34 @@ def parse_pipe_file(resource: ResourceFile) -> PipeModel: "S3 sinks require EXPORT_BUCKET_URI, EXPORT_FILE_TEMPLATE, EXPORT_FORMAT, and EXPORT_SCHEDULE.", ) - sink = SinkS3Model( - service="s3", - connection_name=export_connection_name, - bucket_uri=export_bucket_uri, - file_template=export_file_template, - format=export_format, - schedule=export_schedule, - strategy=export_strategy, # type: ignore[arg-type] - compression=export_compression, # type: ignore[arg-type] + if inferred_service == "gcs_hmac": + sink = SinkGCSModel( + service="gcs_hmac", + connection_name=export_connection_name, + bucket_uri=export_bucket_uri, + file_template=export_file_template, + format=export_format, + schedule=export_schedule, + strategy=export_strategy, # type: ignore[arg-type] + compression=export_compression, # type: ignore[arg-type] + ) + else: + sink = SinkS3Model( + service="s3", + connection_name=export_connection_name, + bucket_uri=export_bucket_uri, + file_template=export_file_template, + format=export_format, + schedule=export_schedule, + strategy=export_strategy, # type: ignore[arg-type] + compression=export_compression, # type: ignore[arg-type] + ) + else: + raise MigrationParseError( + resource.file_path, + "pipe", + resource.name, + f'Unsupported EXPORT_SERVICE in strict mode: "{inferred_service}"', ) params: list[PipeParamModel] diff --git a/src/tinybird_sdk/migrate/runner.py b/src/tinybird_sdk/migrate/runner.py index c24c4bc..682f5fa 100644 --- a/src/tinybird_sdk/migrate/runner.py +++ b/src/tinybird_sdk/migrate/runner.py @@ -254,7 +254,10 @@ def run_migrate(options: MigrateOptions | dict[str, Any]) -> MigrationResult: ) continue - if pipe.sink and sink_connection_type != pipe.sink.service: + # The "gcs_hmac" sink service exports via a plain "gcs" connection. + sink_service = pipe.sink.service if pipe.sink else None + expected_connection_type = "gcs" if sink_service == "gcs_hmac" else sink_service + if pipe.sink and sink_connection_type != expected_connection_type: errors.append( MigrationError( file_path=pipe.file_path, diff --git a/src/tinybird_sdk/migrate/types.py b/src/tinybird_sdk/migrate/types.py index 5b9650e..35a6333 100644 --- a/src/tinybird_sdk/migrate/types.py +++ b/src/tinybird_sdk/migrate/types.py @@ -160,7 +160,19 @@ class SinkS3Model: compression: Literal["none", "gzip", "snappy"] | None = None -SinkModel = SinkKafkaModel | SinkS3Model +@dataclass(frozen=True, slots=True) +class SinkGCSModel: + service: Literal["gcs_hmac"] + connection_name: str + bucket_uri: str + file_template: str + format: str + schedule: str + strategy: Literal["create_new", "replace"] | None = None + compression: Literal["none", "gzip", "snappy"] | None = None + + +SinkModel = SinkKafkaModel | SinkS3Model | SinkGCSModel @dataclass(frozen=True, slots=True) diff --git a/src/tinybird_sdk/schema/__init__.py b/src/tinybird_sdk/schema/__init__.py index fa2ad9f..aa983cc 100644 --- a/src/tinybird_sdk/schema/__init__.py +++ b/src/tinybird_sdk/schema/__init__.py @@ -81,6 +81,7 @@ CopyConfig, KafkaSinkConfig, S3SinkConfig, + GCSSinkConfig, SinkConfig, NodeDefinition, ) diff --git a/src/tinybird_sdk/schema/pipe.py b/src/tinybird_sdk/schema/pipe.py index 1606bf5..083e652 100644 --- a/src/tinybird_sdk/schema/pipe.py +++ b/src/tinybird_sdk/schema/pipe.py @@ -4,7 +4,7 @@ from dataclasses import dataclass, field from typing import Any -from .connection import KafkaConnectionDefinition, S3ConnectionDefinition +from .connection import GCSConnectionDefinition, KafkaConnectionDefinition, S3ConnectionDefinition from .datasource import ColumnDefinition, DatasourceDefinition, SchemaDefinition, get_column_type from .params import ParamValidator from .token import TokenDefinition @@ -70,7 +70,18 @@ class S3SinkConfig: compression: SinkCompression | None = None -SinkConfig = KafkaSinkConfig | S3SinkConfig +@dataclass(frozen=True, slots=True) +class GCSSinkConfig: + connection: GCSConnectionDefinition + bucket_uri: str + file_template: str + format: str + schedule: str + strategy: SinkStrategy | None = None + compression: SinkCompression | None = None + + +SinkConfig = KafkaSinkConfig | S3SinkConfig | GCSSinkConfig @dataclass(frozen=True, slots=True) @@ -299,7 +310,17 @@ def _normalize_sink_config(raw: dict[str, Any]) -> SinkConfig: strategy=raw.get("strategy"), compression=raw.get("compression"), ) - raise ValueError("Sink connection must be a Kafka or S3 connection definition.") + if isinstance(connection, GCSConnectionDefinition): + return GCSSinkConfig( + connection=connection, + bucket_uri=raw["bucket_uri"], + file_template=raw["file_template"], + format=raw["format"], + schedule=raw["schedule"], + strategy=raw.get("strategy"), + compression=raw.get("compression"), + ) + raise ValueError("Sink connection must be a Kafka, S3, or GCS connection definition.") def define_sink_pipe(name: str, options: dict[str, Any]) -> PipeDefinition: diff --git a/tests/test_gcs_sink_destination.py b/tests/test_gcs_sink_destination.py new file mode 100644 index 0000000..1c2c149 --- /dev/null +++ b/tests/test_gcs_sink_destination.py @@ -0,0 +1,262 @@ +from __future__ import annotations + +from pathlib import Path + +import pytest + +from tinybird_sdk.generator.pipe import generate_pipe +from tinybird_sdk.migrate.emit_ts import emit_migration_file_content +from tinybird_sdk.migrate.parse_pipe import parse_pipe_file +from tinybird_sdk.migrate.parser_utils import MigrationParseError +from tinybird_sdk.migrate.runner import run_migrate +from tinybird_sdk.migrate.types import ResourceFile, SinkGCSModel +from tinybird_sdk.schema.connection import define_gcs_connection +from tinybird_sdk.schema.pipe import GCSSinkConfig, define_sink_pipe, get_sink_config, node + + +def _resource(name: str, content: str) -> ResourceFile: + return ResourceFile( + kind="pipe", + name=name, + file_path=f"{name}.pipe", + absolute_path=f"/tmp/{name}.pipe", + content=content, + ) + + +def _gcs_connection(): + return define_gcs_connection( + "landing_gcs", {"service_account_credentials_json": '{"project":"demo"}'} + ) + + +def test_define_sink_pipe_accepts_gcs_connection() -> None: + connection = _gcs_connection() + pipe = define_sink_pipe( + "gcs_events_sink", + { + "sink": { + "connection": connection, + "bucket_uri": "gs://my-bucket/exports/", + "file_template": "events_{date}", + "format": "csv", + "schedule": "@once", + "strategy": "create_new", + "compression": "gzip", + }, + "nodes": [node({"name": "export", "sql": "SELECT 1"})], + }, + ) + + sink = get_sink_config(pipe) + assert isinstance(sink, GCSSinkConfig) + assert sink.connection is connection + assert sink.bucket_uri == "gs://my-bucket/exports/" + assert sink.strategy == "create_new" + assert sink.compression == "gzip" + + +def test_generate_pipe_emits_export_service_gcs_hmac_for_gcs_sink() -> None: + connection = _gcs_connection() + pipe = define_sink_pipe( + "gcs_events_sink", + { + "sink": { + "connection": connection, + "bucket_uri": "gs://my-bucket/exports/", + "file_template": "events_{date}", + "format": "csv", + "schedule": "@once", + }, + "nodes": [node({"name": "export", "sql": "SELECT 1"})], + }, + ) + + generated = generate_pipe(pipe).content + + assert "TYPE sink" in generated + assert "EXPORT_SERVICE gcs_hmac" in generated + assert f"EXPORT_CONNECTION_NAME {connection._name}" in generated + assert "EXPORT_BUCKET_URI gs://my-bucket/exports/" in generated + assert "EXPORT_FILE_TEMPLATE events_{date}" in generated + assert "EXPORT_FORMAT csv" in generated + assert "EXPORT_SCHEDULE @once" in generated + + +def test_parse_pipe_requires_explicit_export_service_for_gcs() -> None: + parsed = parse_pipe_file( + _resource( + "gcs_sink", + "\n".join( + [ + "TYPE sink", + "EXPORT_SERVICE gcs_hmac", + "EXPORT_CONNECTION_NAME landing_gcs", + "EXPORT_BUCKET_URI gs://bucket/path", + "EXPORT_FILE_TEMPLATE {date}.ndjson", + "EXPORT_FORMAT ndjson", + "EXPORT_SCHEDULE @hourly", + "EXPORT_COMPRESSION gzip", + "NODE export", + "SQL >", + " SELECT id FROM events", + ] + ), + ) + ) + + assert parsed.sink is not None + assert isinstance(parsed.sink, SinkGCSModel) + assert parsed.sink.service == "gcs_hmac" + assert parsed.sink.bucket_uri == "gs://bucket/path" + assert parsed.sink.compression == "gzip" + + +def test_parse_pipe_without_export_service_defaults_blob_sink_to_s3() -> None: + parsed = parse_pipe_file( + _resource( + "ambiguous_sink", + "\n".join( + [ + "TYPE sink", + "EXPORT_CONNECTION_NAME archive", + "EXPORT_BUCKET_URI s3://bucket/path", + "EXPORT_FILE_TEMPLATE {date}.ndjson", + "EXPORT_FORMAT ndjson", + "EXPORT_SCHEDULE @hourly", + "NODE export", + "SQL >", + " SELECT id FROM events", + ] + ), + ) + ) + + assert parsed.sink is not None + assert parsed.sink.service == "s3" + + +def test_parse_pipe_rejects_unsupported_export_service() -> None: + with pytest.raises(MigrationParseError, match="Unsupported EXPORT_SERVICE"): + parse_pipe_file( + _resource( + "bad_sink", + "\n".join( + [ + "TYPE sink", + "EXPORT_SERVICE azure_blob", + "EXPORT_CONNECTION_NAME archive", + "EXPORT_BUCKET_URI blob://bucket/path", + "NODE export", + "SQL >", + " SELECT 1", + ] + ), + ) + ) + + +def test_run_migrate_accepts_gcs_hmac_sink_against_gcs_connection(tmp_path: Path) -> None: + (tmp_path / "landing_gcs.connection").write_text( + "\n".join( + [ + "TYPE gcs", + 'GCS_SERVICE_ACCOUNT_CREDENTIALS_JSON \'{"project":"demo"}\'', + ] + ), + encoding="utf-8", + ) + (tmp_path / "gcs_sink.pipe").write_text( + "\n".join( + [ + "TYPE sink", + "EXPORT_SERVICE gcs_hmac", + "EXPORT_CONNECTION_NAME landing_gcs", + "EXPORT_BUCKET_URI gs://bucket/path", + "EXPORT_FILE_TEMPLATE {date}.ndjson", + "EXPORT_FORMAT ndjson", + "EXPORT_SCHEDULE @hourly", + "NODE export", + "SQL >", + " SELECT id FROM events", + ] + ), + encoding="utf-8", + ) + + result = run_migrate( + { + "cwd": str(tmp_path), + "patterns": ["*.connection", "*.pipe"], + "dry_run": True, + } + ) + + assert result.success is True, result.errors + + +def test_run_migrate_rejects_gcs_hmac_sink_against_s3_connection(tmp_path: Path) -> None: + (tmp_path / "archive.connection").write_text( + "\n".join( + [ + "TYPE s3", + "S3_REGION us-east-1", + "S3_ARN arn:aws:iam::123456789012:role/demo", + ] + ), + encoding="utf-8", + ) + (tmp_path / "gcs_sink.pipe").write_text( + "\n".join( + [ + "TYPE sink", + "EXPORT_SERVICE gcs_hmac", + "EXPORT_CONNECTION_NAME archive", + "EXPORT_BUCKET_URI gs://bucket/path", + "EXPORT_FILE_TEMPLATE {date}.ndjson", + "EXPORT_FORMAT ndjson", + "EXPORT_SCHEDULE @hourly", + "NODE export", + "SQL >", + " SELECT id FROM events", + ] + ), + encoding="utf-8", + ) + + result = run_migrate( + { + "cwd": str(tmp_path), + "patterns": ["*.connection", "*.pipe"], + "dry_run": True, + } + ) + + assert result.success is False + assert any("is incompatible with connection" in error.message for error in result.errors) + + +def test_emit_migration_emits_gcs_sink_fields() -> None: + from tinybird_sdk.migrate.types import PipeModel, PipeNodeModel + + pipe = PipeModel( + kind="pipe", + name="gcs_sink", + file_path="gcs_sink.pipe", + type="sink", + nodes=[PipeNodeModel(name="export", sql="SELECT 1")], + sink=SinkGCSModel( + service="gcs_hmac", + connection_name="landing_gcs", + bucket_uri="gs://bucket/path", + file_template="{date}.ndjson", + format="ndjson", + schedule="@hourly", + compression="gzip", + ), + ) + + emitted = emit_migration_file_content([pipe]) + + assert "'bucket_uri': \"gs://bucket/path\"" in emitted + assert "'compression': \"gzip\"" in emitted