From d08d7c2ecf72f14ee8d351eb6c36a9299c0929c2 Mon Sep 17 00:00:00 2001 From: Lukas Bindreiter Date: Mon, 5 Oct 2026 12:32:19 +0200 Subject: [PATCH] Streaming logs worker rpc --- apis/workflows/v1/storage_location.proto | 14 +++++++++++--- apis/workflows/v1/worker.proto | 21 ++++++++++++++++++--- 2 files changed, 29 insertions(+), 6 deletions(-) diff --git a/apis/workflows/v1/storage_location.proto b/apis/workflows/v1/storage_location.proto index 216aa5c..0facb2b 100644 --- a/apis/workflows/v1/storage_location.proto +++ b/apis/workflows/v1/storage_location.proto @@ -150,7 +150,7 @@ message UpdateStorageLocationRequest { // StorageSubscriptionType identifies the notification delivery provider. enum StorageSubscriptionType { STORAGE_SUBSCRIPTION_TYPE_UNSPECIFIED = 0; - // Amazon SNS delivering notifications from an AWS S3 bucket. + // Amazon SNS delivering object notifications from an AWS S3 bucket, or custom SNS message formats configured via message_format. STORAGE_SUBSCRIPTION_TYPE_AWS_SNS = 1; // Google Cloud Pub/Sub delivering notifications from a GCS bucket. STORAGE_SUBSCRIPTION_TYPE_GOOGLE_PUBSUB = 2; @@ -165,6 +165,13 @@ message AWSSNSStorageSubscription { min_bytes: 1 max_bytes: 2048 }]; + // Payload format: s3, earthsearch_stac, or usgs_landsat_c2. If omitted, s3 is assumed. + string message_format = 2 [ + (buf.validate.field).string.in = "", + (buf.validate.field).string.in = "s3", + (buf.validate.field).string.in = "earthsearch_stac", + (buf.validate.field).string.in = "usgs_landsat_c2" + ]; } // GooglePubSubStorageSubscription identifies the delivery subscription and authenticated push identity. @@ -228,8 +235,9 @@ message StorageSubscription { 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. +// CreateStorageSubscriptionRequest registers notification delivery for a storage location. It enables the +// Tilebox webhook endpoint to receive object notifications from the provider. Infrastructure on the provider side +// must then be configured to send notifications to the returned endpoint. message CreateStorageSubscriptionRequest { option (buf.validate.message).oneof = { fields: [ diff --git a/apis/workflows/v1/worker.proto b/apis/workflows/v1/worker.proto index d7366eb..0fb1d35 100644 --- a/apis/workflows/v1/worker.proto +++ b/apis/workflows/v1/worker.proto @@ -4,6 +4,7 @@ package workflows.v1; import "buf/validate/validate.proto"; import "google/protobuf/empty.proto"; +import "opentelemetry/proto/logs/v1/logs.proto"; import "tilebox/v1/id.proto"; import "workflows/v1/core.proto"; import "workflows/v1/task.proto"; @@ -28,8 +29,8 @@ message InitializeRunnerRequest { // TileboxAPIConnection is a message containing the information needed for a worker runtime to connect to the Tilebox // API, to e.g. export observability data to (logs and traces). It's the credentials that the tilebox-cli locally // uses to connect to the Tilebox API, and is passed down to the worker runtimes so that they can also connect to the -// Tilebox API in the same way. Since this only is done on a local machine, via a unix socket, and not in a shared -// environment or over the internet, this is safe to do. +// Tilebox API in the same way. These credentials must only be sent over a trusted local runtime connection, +// using a Unix socket or an authenticated loopback TCP connection, never an exposed network listener. message TileboxAPIConnection { // The API url to which the worker runtime can connect, to e.g. export observability data to (logs and traces). string url = 1; @@ -72,7 +73,21 @@ service WorkerService { // as failed with an unknown error. rpc ExecuteTask(workflows.v1.Task) returns (ExecuteTaskResponse) {} + // WatchLogs streams structured runtime logs to the local runner, independently of API log exports and raw + // stdout/stderr. It is available before workflow import completes and before InitializeWorker is called. + // The runner sends one request; the worker sends buffered startup records followed by live records in enqueue + // order. The stream stays open until shutdown, cancellation, or a connection failure, even while no tasks run. + // Only one subscriber is allowed per runtime; concurrent subscriptions fail with ALREADY_EXISTS. + // Delivery is best-effort with bounded buffering: slow or disconnected subscribers must not block tasks. + // Buffer overflow is reported as a warning record when delivery resumes. There is no acknowledgement or replay + // of delivered records. Canceling this stream does not shut down the worker or disable its API log exports. + // Messages use the log body, severity, trace/span IDs, and attributes from the OpenTelemetry log model; + // exception details use exception.type, exception.message, and exception.stacktrace attributes. + // buf:lint:ignore RPC_NO_SERVER_STREAMING + rpc WatchLogs(google.protobuf.Empty) returns (stream opentelemetry.proto.logs.v1.LogRecord); + // Gracefully shuts down the worker runtime. After receiving this request, the worker runtime will - // cleanly shut down. + // finish task cleanup, flush API logs, and drain buffered WatchLogs records within a bounded deadline. + // The log stream ends before the worker stops its RPC server; an open subscription must not prevent shutdown. rpc ShutdownWorker(google.protobuf.Empty) returns (google.protobuf.Empty) {} }