From 0daaa706a6d16fae6db73050d70b01dd08aa4ea0 Mon Sep 17 00:00:00 2001 From: Stefan Amberger <1277330+snamber@users.noreply.github.com> Date: Fri, 2 Oct 2026 12:13:13 +0000 Subject: [PATCH] Add storage location and notification subscription API Amp-Thread-ID: https://ampcode.com/threads/T-01a0f838-2233-701d-9912-13798be93e4f Co-authored-by: Amp --- apis/workflows/v1/automation.proto | 58 +--- apis/workflows/v1/storage_location.proto | 349 +++++++++++++++++++++++ 2 files changed, 350 insertions(+), 57 deletions(-) create mode 100644 apis/workflows/v1/storage_location.proto diff --git a/apis/workflows/v1/automation.proto b/apis/workflows/v1/automation.proto index 84fad90..ef7c871 100644 --- a/apis/workflows/v1/automation.proto +++ b/apis/workflows/v1/automation.proto @@ -10,41 +10,10 @@ import "google/protobuf/empty.proto"; import "google/protobuf/timestamp.proto"; import "tilebox/v1/id.proto"; import "workflows/v1/core.proto"; +import "workflows/v1/storage_location.proto"; option features.field_presence = IMPLICIT; -// StorageType specifies a kind of storage bucket that we support. -enum StorageType { - STORAGE_TYPE_UNSPECIFIED = 0; - // Google Cloud Storage - STORAGE_TYPE_GCS = 1; - // Amazon Web Services S3 - STORAGE_TYPE_S3 = 2; - // Local filesystem - STORAGE_TYPE_FS = 3; -} - -// Storage location is some kind of storage that can contain data files or objects and be used as a trigger source. -message StorageLocation { - // Unique identifier for the storage location - tilebox.v1.ID id = 1 [(buf.validate.field).required = true]; - // A unique identifier for the storage location in the storage system - string location = 2 [(buf.validate.field).string = { - min_bytes: 1 - max_bytes: 512 - }]; - // The type of the storage location, e.g. GCS, S3, FS - StorageType type = 3 [(buf.validate.field).enum = { - defined_only: true - not_in: [0] - }]; -} - -// Buckets is a list of storage buckets -message StorageLocations { - repeated StorageLocation locations = 1; -} - // AutomationPrototype is a task prototype that can result in many submitted tasks. Task submissions are triggered by // NRT triggers, such as bucket triggers or cron triggers. message AutomationPrototype { @@ -113,22 +82,6 @@ message Automation { bytes args = 2; } -// StorageEventType specifies the type of event that triggered the task. -enum StorageEventType { - STORAGE_EVENT_TYPE_UNSPECIFIED = 0; - STORAGE_EVENT_TYPE_CREATED = 1; -} - -// TriggeredStorageEvent contains the details of the concrete event that triggered a storage event trigger. -message TriggeredStorageEvent { - // The storage location that triggered the task - tilebox.v1.ID storage_location_id = 1; - // The type of the storage event, e.g. created - StorageEventType type = 2; - // The object that triggered the task, e.g. a file name in a directory or object name in a bucket - string location = 3; -} - // TriggeredCronEvent contains the details of a concrete event that triggered a cron trigger. message TriggeredCronEvent { // The time the cron trigger fired @@ -147,15 +100,6 @@ message DeleteAutomationRequest { // - Bucket triggers, which triggers tasks when an object is uploaded to a storage bucket that matches a glob pattern // - Cron triggers, which triggers tasks on a schedule service AutomationService { - // ListStorageLocations lists all the storage buckets that are available for use as bucket triggers. - rpc ListStorageLocations(google.protobuf.Empty) returns (StorageLocations); - // GetStorageLocation gets a storage location by its ID. - rpc GetStorageLocation(tilebox.v1.ID) returns (StorageLocation); - // CreateStorageLocation creates a new storage bucket. - rpc CreateStorageLocation(StorageLocation) returns (StorageLocation); - // DeleteStorageLocation deletes a storage location. - rpc DeleteStorageLocation(tilebox.v1.ID) returns (google.protobuf.Empty); - // ListAutomations lists all the automations that are currently registered in a namespace. rpc ListAutomations(google.protobuf.Empty) returns (Automations); // GetAutomation gets an automation by its ID. diff --git a/apis/workflows/v1/storage_location.proto b/apis/workflows/v1/storage_location.proto new file mode 100644 index 0000000..216aa5c --- /dev/null +++ b/apis/workflows/v1/storage_location.proto @@ -0,0 +1,349 @@ +// The external API for managing storage locations and object notifications. + +edition = "2023"; + +package workflows.v1; + +import "buf/validate/validate.proto"; +import "google/api/field_behavior.proto"; +import "google/protobuf/empty.proto"; +import "google/protobuf/timestamp.proto"; +import "tilebox/v1/id.proto"; +import "tilebox/v1/query.proto"; + +option features.field_presence = IMPLICIT; + +// StorageType identifies the storage provider and resource kind. +enum StorageType { + STORAGE_TYPE_UNSPECIFIED = 0; + // Google Cloud Storage bucket + STORAGE_TYPE_GCS_BUCKET = 1; + // Amazon Web Services S3 bucket + STORAGE_TYPE_AWS_S3_BUCKET = 2; + // Local filesystem + STORAGE_TYPE_FILESYSTEM = 3; + // Azure Blob Storage container + STORAGE_TYPE_AZURE_BLOB = 4; +} + +// StorageLocation identifies a bucket, container, or filesystem directory whose objects can trigger automations. +message StorageLocation { + // Unique identifier of the storage location. + tilebox.v1.ID id = 1; + // Use reference for the provider coordinates. + string location = 2 [deprecated = true]; + // Use reference.type for the storage type. + StorageType type = 3 [deprecated = true]; + // Display name of the storage location. + string name = 4 [(buf.validate.field).string.max_bytes = 1024]; + // Immutable provider coordinates of the storage location. + StorageLocationReference reference = 5; +} + +// StorageLocationReference identifies an existing storage resource. Set the field corresponding to type. +message StorageLocationReference { + option (buf.validate.message).oneof = { + fields: [ + "aws_s3_bucket", + "gcs_bucket", + "azure_blob", + "filesystem" + ] + required: true + }; + option (buf.validate.message).cel = { + id: "storage_location_reference.type" + message: "The storage type must match the provider reference." + expression: "(this.type == 1 && has(this.gcs_bucket)) || (this.type == 2 && has(this.aws_s3_bucket)) || (this.type == 3 && has(this.filesystem)) || (this.type == 4 && has(this.azure_blob))" + }; + + AWSS3BucketReference aws_s3_bucket = 1; + GCSBucketReference gcs_bucket = 2; + AzureBlobReference azure_blob = 3; + FilesystemReference filesystem = 4; + StorageType type = 5 [(buf.validate.field).enum = { + defined_only: true + not_in: [0] + }]; +} + +// AWSS3BucketReference identifies an Amazon Web Services S3 bucket and its region. +message AWSS3BucketReference { + string bucket = 1 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 63 + }]; + string region = 2 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 128 + }]; +} + +// GCSBucketReference identifies a Google Cloud Storage bucket and its project. +message GCSBucketReference { + string bucket = 1 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 222 + }]; + // ID of the Google Cloud project that owns the bucket. + string project_id = 2 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 128 + }]; + // Optional bucket location, such as europe-west2, EUR4, or EU. May be a region, dual-region, or multi-region. + string location = 3 [(buf.validate.field).string.max_bytes = 128]; +} + +// AzureBlobReference identifies a container in an Azure storage account. +message AzureBlobReference { + // Full ARM resource ID of the storage account, not its URL or a connection string. + string storage_account_resource_id = 1 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 2048 + }]; + string container = 2 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 63 + }]; + // Optional primary region of the storage account, such as westeurope. Containers use their account's region. + string region = 3 [(buf.validate.field).string.max_bytes = 128]; +} + +// FilesystemReference identifies a directory monitored by a filesystem notifier. +message FilesystemReference { + string path = 1 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 512 + }]; +} + +// TriggeredJob identifies a job submitted by an automation for a matching object path. +message TriggeredJob { + tilebox.v1.ID job_id = 1; + tilebox.v1.ID automation_id = 2; + string glob_pattern = 3; +} + +// TriggeredJobs contains the jobs submitted for an event. +message TriggeredJobs { + repeated TriggeredJob triggered_jobs = 1; +} + +// CreateStorageLocationRequest registers an existing bucket, container, or directory as a storage location. +message CreateStorageLocationRequest { + string name = 1 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 1024 + }]; + StorageLocationReference reference = 2 [(buf.validate.field).required = true]; +} + +// UpdateStorageLocationRequest changes the display name of a storage location. +message UpdateStorageLocationRequest { + tilebox.v1.ID storage_location_id = 1 [(buf.validate.field).required = true]; + string name = 2 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 1024 + }]; +} + +// StorageSubscriptionType identifies the notification delivery provider. +enum StorageSubscriptionType { + STORAGE_SUBSCRIPTION_TYPE_UNSPECIFIED = 0; + // Amazon SNS delivering notifications from an AWS S3 bucket. + STORAGE_SUBSCRIPTION_TYPE_AWS_SNS = 1; + // Google Cloud Pub/Sub delivering notifications from a GCS bucket. + STORAGE_SUBSCRIPTION_TYPE_GOOGLE_PUBSUB = 2; + // Azure Event Grid delivering notifications from an Azure Blob Storage container. + STORAGE_SUBSCRIPTION_TYPE_AZURE_EVENT_GRID = 3; +} + +// AWSSNSStorageSubscription identifies the SNS topic whose signed notifications are accepted. +// Configure an HTTPS subscription with raw message delivery disabled. Tilebox verifies and confirms signed SNS requests. +message AWSSNSStorageSubscription { + string topic_arn = 1 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 2048 + }]; +} + +// GooglePubSubStorageSubscription identifies the delivery subscription and authenticated push identity. +// Use authenticated push with wrapped JSON and the returned audience. +message GooglePubSubStorageSubscription { + // Full resource name: projects/{project}/subscriptions/{subscription}. + string subscription = 1 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 1024 + }]; + // Expected Google-authenticated service account email. No service account key is accepted. + string service_account_email = 2 [(buf.validate.field).string = { + email: true + max_bytes: 320 + }]; + // OIDC audience to configure for authenticated push, equal to the subscription's endpoint. + string audience = 3 [(google.api.field_behavior) = OUTPUT_ONLY]; +} + +// AzureEventGridStorageSubscription configures Event Grid delivery with a secret header. +// Subscribe to Microsoft.Storage.BlobCreated with the location's container subject prefix. +// Use Event Grid schema, not CloudEvents schema. +message AzureEventGridStorageSubscription { + // Delivery header to configure in Event Grid: X-Tilebox-Webhook-Secret. + string webhook_secret_header = 1 [(google.api.field_behavior) = OUTPUT_ONLY]; + // Secret delivery header value, returned only when the subscription is created. + // Copy it into Event Grid's secret delivery header configuration; it cannot be retrieved later. + string webhook_secret = 2 [(google.api.field_behavior) = OUTPUT_ONLY]; + // Server-managed verifier; never returned by the API. + bytes webhook_secret_hash = 3 [(google.api.field_behavior) = OUTPUT_ONLY]; +} + +// StorageSubscription receives object notifications for a storage location. +// To change delivery configuration, create a replacement subscription and delete the old one. +message StorageSubscription { + option (buf.validate.message).oneof = { + fields: [ + "aws_sns", + "google_pubsub", + "azure_event_grid" + ] + required: true + }; + option (buf.validate.message).cel = { + id: "storage_subscription.type" + message: "The subscription type must match the delivery configuration." + expression: "(this.type == 1 && has(this.aws_sns)) || (this.type == 2 && has(this.google_pubsub)) || (this.type == 3 && has(this.azure_event_grid))" + }; + + tilebox.v1.ID id = 1 [(google.api.field_behavior) = OUTPUT_ONLY]; + tilebox.v1.ID storage_location_id = 2; + StorageSubscriptionType type = 3 [(buf.validate.field).enum = { + defined_only: true + not_in: [0] + }]; + // HTTPS endpoint to configure as the provider's notification delivery URL. + string endpoint = 4 [(google.api.field_behavior) = OUTPUT_ONLY]; + google.protobuf.Timestamp created_at = 6 [(google.api.field_behavior) = OUTPUT_ONLY]; + AWSSNSStorageSubscription aws_sns = 7; + GooglePubSubStorageSubscription google_pubsub = 8; + AzureEventGridStorageSubscription azure_event_grid = 9; +} + +// CreateStorageSubscriptionRequest registers notification delivery for a storage location. +// Set the configuration corresponding to type. The provider must match the storage location. +message CreateStorageSubscriptionRequest { + option (buf.validate.message).oneof = { + fields: [ + "aws_sns", + "google_pubsub", + "azure_event_grid" + ] + required: true + }; + option (buf.validate.message).cel = { + id: "create_storage_subscription_request.type" + message: "The subscription type must match the delivery configuration." + expression: "(this.type == 1 && has(this.aws_sns)) || (this.type == 2 && has(this.google_pubsub)) || (this.type == 3 && has(this.azure_event_grid))" + }; + + tilebox.v1.ID storage_location_id = 1 [(buf.validate.field).required = true]; + StorageSubscriptionType type = 2 [(buf.validate.field).enum = { + defined_only: true + not_in: [0] + }]; + AWSSNSStorageSubscription aws_sns = 3; + GooglePubSubStorageSubscription google_pubsub = 4; + AzureEventGridStorageSubscription azure_event_grid = 5; +} + +// ListStorageSubscriptionsRequest lists subscriptions for a storage location. +message ListStorageSubscriptionsRequest { + tilebox.v1.ID storage_location_id = 1 [(buf.validate.field).required = true]; +} + +// StorageSubscriptions contains notification subscriptions for a storage location. +message StorageSubscriptions { + repeated StorageSubscription subscriptions = 1; +} + +// StorageSubscriptionEvent is an object notification received for a storage location. +message StorageSubscriptionEvent { + // Unique identifier of the received event. + tilebox.v1.ID id = 1; + // Object key within the bucket or container. + string object_key = 2 [(buf.validate.field).string = { + min_bytes: 1 + max_bytes: 4096 + }]; + StorageEventType type = 3 [(buf.validate.field).enum = { + defined_only: true + not_in: [0] + }]; + // Provider event time, absent when unavailable. + google.protobuf.Timestamp event_time = 4; + // Time Tilebox received the notification. + google.protobuf.Timestamp received_at = 5; + // Jobs triggered by this event. + repeated TriggeredJob triggered_jobs = 6; + // Subscription through which the event was received. + tilebox.v1.ID storage_subscription_id = 7; +} + +// ListStorageSubscriptionEventsRequest requests events received for a storage location. +message ListStorageSubscriptionEventsRequest { + tilebox.v1.ID storage_location_id = 1 [(buf.validate.field).required = true]; + // Pagination in descending receipt order; the default and maximum limit is 100. + tilebox.v1.Pagination page = 2 [features.field_presence = EXPLICIT]; +} + +// StorageSubscriptionEvents contains one page of received events, newest first. +message StorageSubscriptionEvents { + repeated StorageSubscriptionEvent events = 1; + // Pagination parameters for the next page, absent when there are no more events. + tilebox.v1.Pagination next_page = 2 [features.field_presence = EXPLICIT]; +} + +// StorageLocations contains storage locations available for automation triggers. +message StorageLocations { + repeated StorageLocation locations = 1; +} + +// StorageEventType specifies the type of event that triggered the task. +enum StorageEventType { + STORAGE_EVENT_TYPE_UNSPECIFIED = 0; + STORAGE_EVENT_TYPE_CREATED = 1; +} + +// TriggeredStorageEvent contains the details of the concrete event that triggered a storage event trigger. +message TriggeredStorageEvent { + // The storage location that triggered the task + tilebox.v1.ID storage_location_id = 1; + // The type of the storage event, e.g. created + StorageEventType type = 2; + // The object key/path that triggered the task, not the storage location's provider coordinates. + string location = 3; +} + +// StorageLocationService manages storage locations, notification subscriptions, and received events. +service StorageLocationService { + // ListStorageLocations lists storage locations available for automation triggers. + rpc ListStorageLocations(google.protobuf.Empty) returns (StorageLocations); + // GetStorageLocation gets a storage location by its ID. + rpc GetStorageLocation(tilebox.v1.ID) returns (StorageLocation); + // CreateStorageLocation registers an existing bucket, container, or directory for automation triggers. + rpc CreateStorageLocation(CreateStorageLocationRequest) returns (StorageLocation); + // UpdateStorageLocation changes the display name of a storage location. + rpc UpdateStorageLocation(UpdateStorageLocationRequest) returns (StorageLocation); + // DeleteStorageLocation removes the location and its subscriptions from Tilebox without deleting cloud resources. + rpc DeleteStorageLocation(tilebox.v1.ID) returns (google.protobuf.Empty); + + // CreateStorageSubscription registers notification delivery and returns the provider setup values. + rpc CreateStorageSubscription(CreateStorageSubscriptionRequest) returns (StorageSubscription); + // ListStorageSubscriptions lists notification subscriptions for a storage location. + rpc ListStorageSubscriptions(ListStorageSubscriptionsRequest) returns (StorageSubscriptions); + // GetStorageSubscription gets a subscription by its ID. + rpc GetStorageSubscription(tilebox.v1.ID) returns (StorageSubscription); + // DeleteStorageSubscription stops accepting notifications at its endpoint without deleting cloud resources. + rpc DeleteStorageSubscription(tilebox.v1.ID) returns (google.protobuf.Empty); + // ListStorageSubscriptionEvents lists received object notifications and the jobs they triggered. + rpc ListStorageSubscriptionEvents(ListStorageSubscriptionEventsRequest) returns (StorageSubscriptionEvents); +}