Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 11 additions & 3 deletions apis/workflows/v1/storage_location.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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.
Expand Down Expand Up @@ -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: [
Expand Down
21 changes: 18 additions & 3 deletions apis/workflows/v1/worker.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -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;
Expand Down Expand Up @@ -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) {}
}
Loading