File
Blob: src/workerd/io/worker-interface.capnp
| 1 | # Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | # Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | # https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | @0xf7958855f6746344; |
| 6 | |
| 7 | using Cxx = import "/capnp/c++.capnp"; |
| 8 | $Cxx.namespace("workerd::rpc"); |
| 9 | # We do not use `$Cxx.allowCancellation` because runAlarm() currently depends on blocking |
| 10 | # cancellation. |
| 11 | |
| 12 | using import "/capnp/compat/http-over-capnp.capnp".HttpMethod; |
| 13 | using import "/capnp/compat/http-over-capnp.capnp".HttpService; |
| 14 | using import "/capnp/compat/byte-stream.capnp".ByteStream; |
| 15 | using import "/workerd/io/outcome.capnp".EventOutcome; |
| 16 | using import "/workerd/io/script-version.capnp".ScriptVersion; |
| 17 | using import "/workerd/io/trace.capnp".TagValue; |
| 18 | using import "/workerd/io/trace.capnp".UserSpanData; |
| 19 | using import "/workerd/io/frankenvalue.capnp".Frankenvalue; |
| 20 | |
| 21 | # A 128-bit trace ID used to identify traces. |
| 22 | struct TraceId { |
| 23 | high @0 :UInt64; |
| 24 | low @1 :UInt64; |
| 25 | } |
| 26 | |
| 27 | # W3C trace flags from an upstream traceparent. |
| 28 | # Bit 0 = sampled. Only meaningful when set to a value. |
| 29 | # When the field is missing or the union is unset, no upstream sampling decision |
| 30 | # was made and the tail worker should make its own sampling decision. |
| 31 | struct TraceFlags { |
| 32 | value :union { |
| 33 | unset @0 :Void; |
| 34 | set @1 :UInt8; |
| 35 | } |
| 36 | } |
| 37 | |
| 38 | # InvocationSpanContext used to identify the current tracing context. Only used internally so far. |
| 39 | struct InvocationSpanContext { |
| 40 | # The 128-bit ID uniquely identifying a trace. |
| 41 | traceId @0 :TraceId; |
| 42 | # The 128-bit ID identifying a worker stage invocation within a trace. |
| 43 | invocationId @1 :TraceId; |
| 44 | # The 64-bit span ID identifying an individual span within a worker stage invocation. |
| 45 | spanId @2 :UInt64; |
| 46 | # W3C trace flags. |
| 47 | traceFlags @3 :TraceFlags; |
| 48 | } |
| 49 | |
| 50 | # Span context for a tail event – this is provided for each tail event. |
| 51 | struct SpanContext { |
| 52 | # The 128-bit ID uniquely identifying a trace. |
| 53 | traceId @0 :TraceId; |
| 54 | # spanId in which this event is handled |
| 55 | # for Onset and SpanOpen events this would be the parent span id |
| 56 | # for Outcome and SpanClose these this would be the span id of the opening Onset and SpanOpen events |
| 57 | # For Hibernate and Mark this would be the span under which they were emitted. |
| 58 | # This is only empty if: |
| 59 | # 1. This is an Onset event |
| 60 | # 2. We are not inheriting any SpanContext. (e.g. this is a cross-account service binding or a new top-level invocation) |
| 61 | info :union { |
| 62 | empty @1 :Void; |
| 63 | spanId @2 :UInt64; |
| 64 | } |
| 65 | # W3C trace flags. |
| 66 | traceFlags @3 :TraceFlags; |
| 67 | } |
| 68 | |
| 69 | struct Trace @0x8e8d911203762d34 { |
| 70 | logs @0 :List(Log); |
| 71 | struct Log { |
| 72 | timestampNs @0 :Int64; |
| 73 | |
| 74 | logLevel @1 :Level; |
| 75 | enum Level { |
| 76 | debug @0 $Cxx.name("debug_"); # avoid collision with macro on Apple platforms |
| 77 | info @1; |
| 78 | log @2; |
| 79 | warn @3; |
| 80 | error @4; |
| 81 | } |
| 82 | |
| 83 | message @2 :Text; |
| 84 | } |
| 85 | |
| 86 | obsolete26 @26 :List(UserSpanData); |
| 87 | # spans are unavailable in full trace objects. |
| 88 | |
| 89 | exceptions @1 :List(Exception); |
| 90 | struct Exception { |
| 91 | timestampNs @0 :Int64; |
| 92 | name @1 :Text; |
| 93 | message @2 :Text; |
| 94 | stack @3 :Text; |
| 95 | } |
| 96 | |
| 97 | outcome @2 :EventOutcome; |
| 98 | scriptName @4 :Text; |
| 99 | scriptVersion @19 :ScriptVersion; |
| 100 | scriptId @23 :Text; |
| 101 | |
| 102 | eventTimestampNs @5 :Int64; |
| 103 | |
| 104 | eventInfo :union { |
| 105 | none @3 :Void; |
| 106 | fetch @6 :FetchEventInfo; |
| 107 | jsRpc @21 :JsRpcEventInfo; |
| 108 | scheduled @7 :ScheduledEventInfo; |
| 109 | alarm @9 :AlarmEventInfo; |
| 110 | queue @15 :QueueEventInfo; |
| 111 | custom @13 :CustomEventInfo; |
| 112 | email @16 :EmailEventInfo; |
| 113 | trace @18 :TraceEventInfo; |
| 114 | hibernatableWebSocket @20 :HibernatableWebSocketEventInfo; |
| 115 | connect @29 :ConnectEventInfo; |
| 116 | } |
| 117 | struct FetchEventInfo { |
| 118 | method @0 :HttpMethod; |
| 119 | url @1 :Text; |
| 120 | cfJson @2 :Text; |
| 121 | # Empty string indicates missing cf blob |
| 122 | headers @3 :List(Header); |
| 123 | struct Header { |
| 124 | name @0 :Text; |
| 125 | value @1 :Text; |
| 126 | } |
| 127 | } |
| 128 | |
| 129 | struct JsRpcEventInfo { |
| 130 | methodName @0 :Text; |
| 131 | } |
| 132 | |
| 133 | struct ConnectEventInfo { |
| 134 | } |
| 135 | |
| 136 | struct ScheduledEventInfo { |
| 137 | scheduledTime @0 :Float64; |
| 138 | cron @1 :Text; |
| 139 | } |
| 140 | |
| 141 | struct AlarmEventInfo { |
| 142 | scheduledTimeMs @0 :Int64; |
| 143 | } |
| 144 | |
| 145 | struct QueueEventInfo { |
| 146 | queueName @0 :Text; |
| 147 | batchSize @1 :UInt32; |
| 148 | } |
| 149 | |
| 150 | struct EmailEventInfo { |
| 151 | mailFrom @0 :Text; |
| 152 | rcptTo @1 :Text; |
| 153 | rawSize @2 :UInt32; |
| 154 | } |
| 155 | |
| 156 | struct TracePreviewInfo { |
| 157 | id @0 :Text; |
| 158 | slug @1 :Text; |
| 159 | name @2 :Text; |
| 160 | } |
| 161 | |
| 162 | struct TraceEventInfo { |
| 163 | struct TraceItem { |
| 164 | scriptName @0 :Text; |
| 165 | } |
| 166 | |
| 167 | traces @0 :List(TraceItem); |
| 168 | } |
| 169 | |
| 170 | struct HibernatableWebSocketEventInfo { |
| 171 | type :union { |
| 172 | message @0 :Void; |
| 173 | close :group { |
| 174 | code @1 :UInt16; |
| 175 | wasClean @2 :Bool; |
| 176 | } |
| 177 | error @3 :Void; |
| 178 | } |
| 179 | } |
| 180 | |
| 181 | struct CustomEventInfo {} |
| 182 | |
| 183 | response @8 :FetchResponseInfo; |
| 184 | struct FetchResponseInfo { |
| 185 | statusCode @0 :UInt16; |
| 186 | } |
| 187 | |
| 188 | cpuTime @10 :UInt64; |
| 189 | wallTime @11 :UInt64; |
| 190 | |
| 191 | dispatchNamespace @12 :Text; |
| 192 | scriptTags @14 :List(Text); |
| 193 | |
| 194 | entrypoint @22 :Text; |
| 195 | preview @30 :TracePreviewInfo; |
| 196 | durableObjectId @27 :Text; |
| 197 | tailAttributes @28 :List(Attribute); |
| 198 | |
| 199 | diagnosticChannelEvents @17 :List(DiagnosticChannelEvent); |
| 200 | struct DiagnosticChannelEvent { |
| 201 | timestampNs @0 :Int64; |
| 202 | channel @1 :Text; |
| 203 | message @2 :Data; |
| 204 | } |
| 205 | |
| 206 | # Indicates how many tail stream events were dropped in total. |
| 207 | struct DroppedEvents { |
| 208 | count @0 :UInt32; |
| 209 | } |
| 210 | |
| 211 | struct StreamDiagnosticsEvent { |
| 212 | # In the future, we plan to support several types of events here, for now only the dropped |
| 213 | # events diagnostic is supported. |
| 214 | diagnostic :union { |
| 215 | undefined @0 :Void; |
| 216 | droppedEvents @1 :DroppedEvents; |
| 217 | } |
| 218 | } |
| 219 | |
| 220 | truncated @24 :Bool; |
| 221 | # Indicates that the trace was truncated due to reaching the maximum size limit. |
| 222 | |
| 223 | enum ExecutionModel { |
| 224 | stateless @0; |
| 225 | durableObject @1; |
| 226 | workflow @2; |
| 227 | } |
| 228 | executionModel @25 :ExecutionModel; |
| 229 | # the execution model of the worker being traced. Can be stateless for a regular worker, |
| 230 | # durableObject for a DO worker or workflow for the upcoming Workflows feature. |
| 231 | |
| 232 | # ===================================================================================== |
| 233 | # Additional types for streaming tail workers |
| 234 | |
| 235 | struct Attribute { |
| 236 | # An Attribute mark is used to add detail to a span over its lifetime. |
| 237 | # The Attribute struct can also be used to provide arbitrary additional |
| 238 | # properties for some other structs. |
| 239 | # Modeled after https://opentelemetry.io/docs/concepts/signals/traces/#attributes |
| 240 | name @0 :Text; |
| 241 | value @1 :List(TagValue); |
| 242 | } |
| 243 | |
| 244 | struct Return { |
| 245 | # A Return mark is used to mark the point at which a span operation returned |
| 246 | # a value. For instance, when a fetch subrequest response is received, or when |
| 247 | # the fetch handler returns a Response. Importantly, it does not signal that the |
| 248 | # span has closed, which may not happen for some period of time after the return |
| 249 | # mark is recorded (e.g. due to things like waitUntils or waiting to fully ready |
| 250 | # the response body payload, etc). Not all spans will have a Return mark. |
| 251 | info :union { |
| 252 | empty @0 :Void; |
| 253 | fetch @1 :FetchResponseInfo; |
| 254 | } |
| 255 | } |
| 256 | |
| 257 | struct SpanOpen { |
| 258 | # Marks the opening of a child span within the streaming tail session. |
| 259 | operationName @0 :Text; |
| 260 | spanId @1 :UInt64; |
| 261 | # id for the span being opened by this SpanOpen event. |
| 262 | info :union { |
| 263 | empty @2 :Void; |
| 264 | custom @3 :List(Attribute); |
| 265 | fetch @4 :FetchEventInfo; |
| 266 | jsRpc @5 :JsRpcEventInfo; |
| 267 | } |
| 268 | } |
| 269 | |
| 270 | struct SpanClose { |
| 271 | # Marks the closing of a child span within the streaming tail session. |
| 272 | # Once emitted, no further mark events should occur within the closed |
| 273 | # span. |
| 274 | outcome @0 :EventOutcome; |
| 275 | } |
| 276 | |
| 277 | struct Onset { |
| 278 | # The Onset and Outcome event types are special forms of SpanOpen and |
| 279 | # SpanClose that explicitly mark the start and end of the root span. |
| 280 | # A streaming tail session will always begin with an Onset event, and |
| 281 | # always end with an Outcome event. |
| 282 | executionModel @0 :ExecutionModel; |
| 283 | scriptName @1 :Text; |
| 284 | scriptVersion @2 :ScriptVersion; |
| 285 | dispatchNamespace @3 :Text; |
| 286 | scriptId @4 :Text; |
| 287 | scriptTags @5 :List(Text); |
| 288 | entryPoint @6 :Text; |
| 289 | preview @10 :TracePreviewInfo; |
| 290 | |
| 291 | struct Info { union { |
| 292 | fetch @0 :FetchEventInfo; |
| 293 | jsRpc @1 :JsRpcEventInfo; |
| 294 | scheduled @2 :ScheduledEventInfo; |
| 295 | alarm @3 :AlarmEventInfo; |
| 296 | queue @4 :QueueEventInfo; |
| 297 | email @5 :EmailEventInfo; |
| 298 | trace @6 :TraceEventInfo; |
| 299 | hibernatableWebSocket @7 :HibernatableWebSocketEventInfo; |
| 300 | connect @9 :ConnectEventInfo; |
| 301 | custom @8 :CustomEventInfo; |
| 302 | } |
| 303 | } |
| 304 | info @7: Info; |
| 305 | spanId @8: UInt64; |
| 306 | # id for the span being opened by this Onset event. |
| 307 | attributes @9 :List(Attribute); |
| 308 | } |
| 309 | |
| 310 | struct Outcome { |
| 311 | outcome @0 :EventOutcome; |
| 312 | cpuTime @1 :UInt64; |
| 313 | wallTime @2 :UInt64; |
| 314 | } |
| 315 | |
| 316 | struct TailEvent { |
| 317 | # A streaming tail worker receives a series of Tail Events. Tail events always occur within an |
| 318 | # InvocationSpanContext. The first TailEvent delivered to a streaming tail session is always an |
| 319 | # Onset. The final TailEvent delivered is always an Outcome. Between those can be any number of |
| 320 | # SpanOpen, SpanClose, and Mark events. Every SpanOpen *must* be associated with a SpanClose |
| 321 | # unless the stream was abruptly terminated. |
| 322 | # Inherited spanContext for this event. |
| 323 | spanContext @0: SpanContext; |
| 324 | # invocation id of the currently invoked worker stage. |
| 325 | # invocation id will always be unique to every Onset event and will be the same until the Outcome event. |
| 326 | invocationId @1: TraceId; |
| 327 | # time for the tail event. This will be provided as I/O time from the perspective of the tail worker. |
| 328 | timestampNs @2 :Int64; |
| 329 | # unique sequence identifier for this tail event, starting at zero. |
| 330 | sequence @3 :UInt32; |
| 331 | event :union { |
| 332 | onset @4 :Onset; |
| 333 | outcome @5 :Outcome; |
| 334 | spanOpen @6 :SpanOpen; |
| 335 | spanClose @7 :SpanClose; |
| 336 | attribute @8 :List(Attribute); |
| 337 | return @9 :Return; |
| 338 | diagnosticChannelEvent @10 :DiagnosticChannelEvent; |
| 339 | exception @11 :Exception; |
| 340 | log @12 :Log; |
| 341 | streamDiagnostics @13 :StreamDiagnosticsEvent; |
| 342 | } |
| 343 | } |
| 344 | } |
| 345 | |
| 346 | struct SendTracesRun @0xde913ebe8e1b82a5 { |
| 347 | outcome @0 :EventOutcome; |
| 348 | } |
| 349 | |
| 350 | struct ScheduledRun @0xd98fc1ae5c8095d0 { |
| 351 | outcome @0 :EventOutcome; |
| 352 | |
| 353 | retry @1 :Bool; |
| 354 | } |
| 355 | |
| 356 | struct AlarmRun @0xfa8ea4e97e23b03d { |
| 357 | outcome @0 :EventOutcome; |
| 358 | |
| 359 | retry @1 :Bool; |
| 360 | retryCountsAgainstLimit @2 :Bool = true; |
| 361 | errorDescription @3 :Text; |
| 362 | } |
| 363 | |
| 364 | struct QueueMessage @0x944adb18c0352295 { |
| 365 | id @0 :Text; |
| 366 | timestampNs @1 :Int64; |
| 367 | data @2 :Data; |
| 368 | contentType @3 :Text; |
| 369 | attempts @4 :UInt16; |
| 370 | } |
| 371 | |
| 372 | struct QueueRetryBatch { |
| 373 | retry @0 :Bool; |
| 374 | union { |
| 375 | undefined @1 :Void; |
| 376 | delaySeconds @2 :Int32; |
| 377 | } |
| 378 | } |
| 379 | |
| 380 | struct QueueRetryMessage { |
| 381 | msgId @0 :Text; |
| 382 | union { |
| 383 | undefined @1 :Void; |
| 384 | delaySeconds @2 :Int32; |
| 385 | } |
| 386 | } |
| 387 | |
| 388 | struct QueueResponse @0x90e98932c0bfc0de { |
| 389 | outcome @0 :EventOutcome; |
| 390 | ackAll @1 :Bool; |
| 391 | retryBatch @2 :QueueRetryBatch; |
| 392 | # Retry options for the batch. |
| 393 | explicitAcks @3 :List(Text); |
| 394 | # List of Message IDs that were explicitly marked as acknowledged. |
| 395 | retryMessages @4 :List(QueueRetryMessage); |
| 396 | # List of retry options for messages that were explicitly marked for retry. |
| 397 | } |
| 398 | |
| 399 | struct MessageBatchMetrics { |
| 400 | backlogCount @0 :Float64; |
| 401 | # Number of messages remaining in the queue backlog. |
| 402 | backlogBytes @1 :Float64; |
| 403 | # Total bytes of messages remaining in the queue backlog. |
| 404 | oldestMessageTimestamp @2 :Float64; |
| 405 | # Timestamp (ms since epoch) of the oldest message in the queue. |
| 406 | } |
| 407 | |
| 408 | struct MessageBatchMetadata { |
| 409 | metrics @0 :MessageBatchMetrics; |
| 410 | # Best effort queue metrics at the time the batch was dispatched. |
| 411 | } |
| 412 | |
| 413 | struct HibernatableWebSocketEventMessage { |
| 414 | payload :union { |
| 415 | text @0 :Text; |
| 416 | data @1 :Data; |
| 417 | close :group { |
| 418 | code @2 :UInt16; |
| 419 | reason @3 :Text; |
| 420 | wasClean @4 :Bool; |
| 421 | } |
| 422 | error @5 :Text; |
| 423 | # TODO(someday): This could be an Exception instead of Text. |
| 424 | } |
| 425 | websocketId @6: Text; |
| 426 | eventTimeoutMs @7: UInt32; |
| 427 | } |
| 428 | |
| 429 | struct HibernatableWebSocketResponse { |
| 430 | outcome @0 :EventOutcome; |
| 431 | } |
| 432 | |
| 433 | interface HibernatableWebSocketEventDispatcher { |
| 434 | hibernatableWebSocketEvent @0 (message: HibernatableWebSocketEventMessage ) |
| 435 | -> (result :HibernatableWebSocketResponse); |
| 436 | # Run a hibernatable websocket event |
| 437 | } |
| 438 | |
| 439 | enum SerializationTag { |
| 440 | # Tag values for all serializable types supported by the Workers API. |
| 441 | |
| 442 | invalid @0; |
| 443 | # Not assigned to anything. Reserved to make things less weird if a zero-valued tag gets written |
| 444 | # by accident. |
| 445 | |
| 446 | jsRpcStub @1; |
| 447 | |
| 448 | writableStream @2; |
| 449 | readableStream @3; |
| 450 | |
| 451 | headers @4; |
| 452 | request @5; |
| 453 | response @6; |
| 454 | |
| 455 | domException @7; |
| 456 | domExceptionV2 @8; |
| 457 | # Keep this value in sync with the DOMException::SERIALIZATION_TAG in |
| 458 | # /src/workerd/jsg/dom-exception (but we can't actually change this value |
| 459 | # without breaking things). |
| 460 | |
| 461 | abortSignal @9; |
| 462 | |
| 463 | nativeError @10; |
| 464 | # A JavaScript native error, such as Error, TypeError, etc. These are typically |
| 465 | # not handled as host objects in V8 but we handle them as such in workers in |
| 466 | # order to preserve additional information that we may attach to them. |
| 467 | |
| 468 | serviceStub @11; |
| 469 | # A ServiceStub aka Fetcher aka Service Binding. |
| 470 | # |
| 471 | # Such stubs are different from jsRpcStub in that they don't point to a single live object, but |
| 472 | # instead represent a service that can be instantiated anywhere. This means that they can be |
| 473 | # passed around the world and instantiated in a different location, as well as persisted in |
| 474 | # long-term storage. |
| 475 | # |
| 476 | # Also because of all this, service stubs can be embedded in the `env` and `ctx.props` of other |
| 477 | # Workers. Regular RPC stubs cannot. |
| 478 | |
| 479 | actorClass @12; |
| 480 | # An actor class reference, aka DurableObjectClass. Can be used to instantiate a facet. |
| 481 | # |
| 482 | # Similar to serviceStub, this refers to the entrypoint of a Worker that can be instantiated |
| 483 | # anywhere and any time, and thus can be persisted and used in `env` and `ctx.props`, etc. |
| 484 | } |
| 485 | |
| 486 | enum StreamEncoding { |
| 487 | # Specifies the internal content-encoding of a ReadableStream or WritableStream. This serves an |
| 488 | # optimization which is not visible to the app: if we end up hooking up streams so that a source |
| 489 | # is pumped to a sink that has the same encoding, we can avoid a decompression/recompression |
| 490 | # round trip. However, if the application reads/writes raw bytes, then we must decode/encode |
| 491 | # them under the hood. |
| 492 | |
| 493 | identity @0; |
| 494 | gzip @1; |
| 495 | brotli @2; |
| 496 | } |
| 497 | |
| 498 | interface Handle { |
| 499 | # Type with no methods, but something happens when you drop it. |
| 500 | } |
| 501 | |
| 502 | struct JsValue { |
| 503 | # A serialized JavaScript value being passed over RPC. |
| 504 | |
| 505 | v8Serialized @0 :Data; |
| 506 | # JS value that has been serialized for network transport. |
| 507 | |
| 508 | externals @1 :List(External); |
| 509 | # The serialized data may contain "externals" -- references to external resources that cannot |
| 510 | # simply be serialized. If so, they are placed in this separate list of externals. |
| 511 | # |
| 512 | # (We could also call these "capabilities", but that word is pretty overloaded already.) |
| 513 | |
| 514 | struct External { |
| 515 | union { |
| 516 | invalid @0 :Void; |
| 517 | # Invalid default value to reduce confusion if an External wasn't initialized properly. |
| 518 | # This should never appear in a real JsValue. |
| 519 | |
| 520 | rpcTarget @1 :JsRpcTarget; |
| 521 | # An object that can be called over RPC. |
| 522 | |
| 523 | writableStream :group { |
| 524 | # A WritableStream. This is much easier to represent that ReadableStream because the bytes |
| 525 | # flow from the receiver to the sender, and therefore a round trip is obviously necessary |
| 526 | # before the bytes can begin flowing. |
| 527 | |
| 528 | byteStream @2 :ByteStream; |
| 529 | encoding @3 :StreamEncoding; |
| 530 | } |
| 531 | |
| 532 | readableStream :group { |
| 533 | # A ReadableStream. The sender of the JsValue will use the associated StreamSink to open a |
| 534 | # stream of type `ByteStream`. |
| 535 | |
| 536 | stream @10 :ExternalPusher.InputStream; |
| 537 | # If present, a stream pushed using the destination isolate's ExternalPusher. |
| 538 | # |
| 539 | # If null (deprecated), then the sender will use the associated StreamSink to open a stream |
| 540 | # of type `ByteStream`. StreamSink is in the process of being replaced by ExternalPusher. |
| 541 | |
| 542 | encoding @4 :StreamEncoding; |
| 543 | # Bytes read from the stream have this encoding. |
| 544 | |
| 545 | expectedLength :union { |
| 546 | # NOTE: This is obsolete when `stream` is set. Instead, the length is passed to |
| 547 | # ExternalPusher.pushByteStream(). |
| 548 | |
| 549 | unknown @5 :Void; |
| 550 | known @6 :UInt64; |
| 551 | } |
| 552 | } |
| 553 | |
| 554 | abortTrigger @7 :Void; |
| 555 | # Indicates that an `AbortTrigger` is being passed, see the `AbortTrigger` interface for the |
| 556 | # mechanism used to trigger the abort later. This is modeled as a stream, since the sender is |
| 557 | # the one that will later on send the abort signal. This external will have an associated |
| 558 | # stream in the corresponding `StreamSink` with type `AbortTrigger`. |
| 559 | # |
| 560 | # TODO(soon): This will be obsolete when we stop using `StreamSink`; `abortSignal` will |
| 561 | # replace it. (The name is wrong anyway -- this is the signal end, not the trigger end.) |
| 562 | |
| 563 | abortSignal @11 :ExternalPusher.AbortSignal; |
| 564 | # Indicates that an `AbortSignal` is being passed. |
| 565 | |
| 566 | subrequestChannelToken @8 :Data; |
| 567 | actorClassChannelToken @9 :Data; |
| 568 | # Encoded ChannelTokens. See channel-token.capnp. |
| 569 | |
| 570 | # TODO(soon): WebSocket, Request, Response |
| 571 | } |
| 572 | } |
| 573 | |
| 574 | interface StreamSink { |
| 575 | # A JsValue may contain streams that flow from the sender to the receiver. We don't want such |
| 576 | # streams to require a network round trip before the stream can begin pumping. So, we need a |
| 577 | # place to start sending bytes right away. |
| 578 | # |
| 579 | # To that end, JsRpcTarget::call() returns a `paramsStreamSink`. Immediately upon sending the |
| 580 | # request, the client can use promise pipelining to begin pushing bytes to this object. |
| 581 | # |
| 582 | # Similarly, the caller passes a `resultsStreamSink` to the callee. If the response contains |
| 583 | # any streams, it can start pushing to this immediately after responding. |
| 584 | # |
| 585 | # TODO(soon): This design is overcomplicated since it requires allocating StreamSinks for every |
| 586 | # request, even when not used, and requires a lot of weird promise magic. The newer |
| 587 | # ExternalPusher design is simpler, and only incurs overhead when used. Once all of |
| 588 | # production has been updated to understand ExternalPusher, then we can flip an autogate to |
| 589 | # use it by default. Once that has rolled out globally, we can remove StreamSink. |
| 590 | |
| 591 | startStream @0 (externalIndex :UInt32) -> (stream :Capability); |
| 592 | # Opens a stream corresponding to the given index in the JsValue's `externals` array. The type |
| 593 | # of capability returned depends on the type of external. E.g. for `readableStream`, it is a |
| 594 | # `ByteStream`. |
| 595 | } |
| 596 | |
| 597 | interface ExternalPusher { |
| 598 | # This object allows "pushing" external objects to a target isolate, so that they can |
| 599 | # sublequently be referenced by a `JsValue.External`. This allows implementing externals where |
| 600 | # the sender might need to send subsequent information to the receiver *before* the receiver |
| 601 | # has had a chance to call back to request it. For example, when a ReadableStream is sent over |
| 602 | # RPC, the sender will immediately start sending body bytes without waiting for a round trip. |
| 603 | # |
| 604 | # The key to ExternalPusher is that it constructs and returns capabilities pointing at objects |
| 605 | # living directly in the target isolate's runtime. These capabilities have empty interfaces, |
| 606 | # but can be passed back to the target in the `External` table of a `JsValue`. Since the |
| 607 | # capabilities point to objects directly in the recipient's memory space, they can then be |
| 608 | # unwrapped to obtain the underlying local object, which the recipient then uses to back the |
| 609 | # external value delivered to the application. |
| 610 | # |
| 611 | # Typically, externals are pushed before the JsValue that uses them is sent. However, this |
| 612 | # is not strictly required, as all pushable externals deserialize into an object that can wait |
| 613 | # for the push as a first step. For example, if a ReadableStream is received in a JsValue |
| 614 | # before pushByteStream() has been called, the resulting stream will simply wait for the push |
| 615 | # as part of its first read() call. This is important as Cap'n Proto doesn't strictly guarantee |
| 616 | # call ordering between calls on different capabilities, or a call vs. a return, so it's |
| 617 | # difficult for the sender to guarantee that a push arrives before the JsValue that refers |
| 618 | # to it. |
| 619 | |
| 620 | pushByteStream @0 (lengthPlusOne :UInt64 = 0) -> (source :InputStream, sink :ByteStream); |
| 621 | # Creates a readable stream within the remote's memory space. `source` should be placed in a |
| 622 | # sublequent `External` of type `readableStream`. The caller should write bytes to `sink`. |
| 623 | # |
| 624 | # `lengthPlusOne` is the expected length of the stream, plus 1, with zero indicating no |
| 625 | # expectation. This is used e.g. when the `ReadableStream` was created with `FixedLengthStream`. |
| 626 | # (The weird "plus one" encoding is used because Cap'n Proto doesn't have a Maybe. Perhaps we |
| 627 | # can fix this eventually.) |
| 628 | |
| 629 | interface InputStream { |
| 630 | # No methods. This will be unwrapped by the recipient to obtain the underlying local value. |
| 631 | } |
| 632 | |
| 633 | pushAbortSignal @1 () -> (signal :AbortSignal, trigger :AbortTrigger); |
| 634 | |
| 635 | interface AbortSignal { |
| 636 | # No methods. This can be unwrapped by the recipient to obtain a Promise<void> which |
| 637 | # rejects when the signal is aborted. |
| 638 | } |
| 639 | |
| 640 | # TODO(soon): |
| 641 | # - Promises |
| 642 | } |
| 643 | } |
| 644 | |
| 645 | interface AbortTrigger $Cxx.allowCancellation { |
| 646 | # When an `AbortSignal` is sent over RPC, the sender initiates a "stream" with this RPC interface |
| 647 | # type which is later used to signal the abort. This is not really a "stream", since only one |
| 648 | # message is sent. But it makes sense to model this way because the message is sent in the same |
| 649 | # direction as the `JsValue` that originally transmitted the `AbortSignal` object. |
| 650 | # When an `AbortSignal` is serialized, the original signal is the client, and the deserialized |
| 651 | # clone is the server. |
| 652 | |
| 653 | abort @0 (reason :JsValue) -> (); |
| 654 | # Allows a cloned abort signal to be triggered over RPC when the original signal is triggered. |
| 655 | # `reason` is an arbitrary JavaScript value which will appear in the resulting `AbortError`s. |
| 656 | |
| 657 | release @1 () -> (); |
| 658 | # Informs a cloned signal that the original signal is being destroyed, and the abort will never |
| 659 | # be triggered. Otherwise, the cloned signal will treat a dropped cabability as an abort. |
| 660 | } |
| 661 | |
| 662 | interface JsRpcTarget extends(JsValue.ExternalPusher) $Cxx.allowCancellation { |
| 663 | # Target on which RPC methods may be invoked. |
| 664 | # |
| 665 | # This is the backing capnp type for a JsRpcStub, as well as used to represent top-level RPC |
| 666 | # events. |
| 667 | # |
| 668 | # JsRpcTarget must implement `JsValue.ExternalPusher` to allow externals to be pushed to the |
| 669 | # target in advance of a call that uses them. |
| 670 | |
| 671 | struct CallParams { |
| 672 | union { |
| 673 | methodName @0 :Text; |
| 674 | # Equivalent to `methodPath` where the list has only one element equal to this. |
| 675 | |
| 676 | methodPath @2 :List(Text); |
| 677 | # Path of properties to follow from the JsRpcTarget itself to find the method being called. |
| 678 | # E.g. if the application does: |
| 679 | # |
| 680 | # myRpcTarget.foo.bar.baz() |
| 681 | # |
| 682 | # Then the path is ["foo", "bar", "baz"]. |
| 683 | # |
| 684 | # The path can also be empty, which means that the JsRpcTarget itself is being invoked as a |
| 685 | # function. |
| 686 | } |
| 687 | |
| 688 | operation :union { |
| 689 | callWithArgs @1 :JsValue; |
| 690 | # Call the property as a function. This is a JsValue that always encodes a JavaScript Array |
| 691 | # containing the arguments to the call. |
| 692 | # |
| 693 | # If `callWithArgs` is null (but is still the active member of the union), this indicates |
| 694 | # that the argument list is empty. |
| 695 | |
| 696 | getProperty @3 :Void; |
| 697 | # This indicates that we are not actually calling a method at all, but rather retrieving the |
| 698 | # value of a property. RPC classes are allowed to define properties that can be fetched |
| 699 | # asynchronously, although more commonly properties will be RPC targets themselves and their |
| 700 | # methods will be invoked by sending a `methodPath` with more than one element. That is, |
| 701 | # imagine you have: |
| 702 | # |
| 703 | # myRpcTarget.foo.bar(); |
| 704 | # |
| 705 | # This code makes a single RPC call with a path of ["foo", "bar"]. However, you could also |
| 706 | # write: |
| 707 | # |
| 708 | # let foo = await myRpcTarget.foo; |
| 709 | # foo.bar(); |
| 710 | # |
| 711 | # This will make two separate calls. The first call is to "foo" and `getProperty` is used. |
| 712 | # This returns a new JsRpcTarget. The second call is on that target, invoking the method |
| 713 | # "bar". |
| 714 | } |
| 715 | |
| 716 | resultsStreamHandler :union { |
| 717 | # We're in the process of switching from `StreamSink` to `ExternalPusher`. A caller will only |
| 718 | # offer one or the other, and expect the callee to use that. (Initially, callers will still |
| 719 | # send StreamSink for backwards-compatibility, but once all recipients are able to understand |
| 720 | # ExternalPusher, we'll flip an autogate to make callers send it.) |
| 721 | |
| 722 | streamSink @4 :JsValue.StreamSink; |
| 723 | # StreamSink used for ReadableStreams found in the results. |
| 724 | |
| 725 | externalPusher @5 :JsValue.ExternalPusher; |
| 726 | # ExternalPusher object which will push into the caller's isolate. Use this to push externals |
| 727 | # that will be included in the results. |
| 728 | } |
| 729 | } |
| 730 | |
| 731 | struct CallResults { |
| 732 | result @0 :JsValue; |
| 733 | # The returned value. |
| 734 | |
| 735 | callPipeline @1 :JsRpcTarget; |
| 736 | # Enables promise pipelining on the eventual call result. This is a JsRpcTarget wrapping the |
| 737 | # result of the call, even if the result itself is a serializable object that would not |
| 738 | # normally be treated as an RPC target. The caller may use this to initiate speculative calls |
| 739 | # on this result without waiting for the initial call to complete (using promise pipelining). |
| 740 | |
| 741 | hasDisposer @2 :Bool; |
| 742 | # If `hasDisposer` is true, the server side returned a serializable object (not a stub) with a |
| 743 | # disposer (Symbol.dispose). The disposer itself is not included in the object's serialization, |
| 744 | # but dropping the `callPipeline` will invoke it. |
| 745 | # |
| 746 | # On the client side, when an RPC returns a plain object, a disposer is added to it. In order |
| 747 | # to avoid confusion, we want the server-side disposer to be invoked only after the client-side |
| 748 | # disposer is invoked. To that end, when `hasDisposer` is true, the client should hold on to |
| 749 | # `callPipeline` until the disposer is invoked. If `hasDisposer` is false, `callPipeline` can |
| 750 | # safely be dropped immediately. |
| 751 | |
| 752 | paramsStreamSink @3 :JsValue.StreamSink; |
| 753 | # StreamSink used for ReadableStreams found in the params. The caller begins sending bytes for |
| 754 | # these streams immediately using promise pipelining. |
| 755 | } |
| 756 | |
| 757 | call @0 CallParams -> CallResults; |
| 758 | # Runs a Worker/DO's RPC method. |
| 759 | } |
| 760 | |
| 761 | interface JsRpcSession { |
| 762 | # Represents an ongoing JSRPC session. To cancel the session (revoking all capabilities), drop |
| 763 | # this object. The `JsRpcSession` object will resolve itself to a null capability when the |
| 764 | # session is complete (all stubs have been droppped); the caller can await `whenResolved()` to |
| 765 | # find out when this happens. |
| 766 | |
| 767 | # Currently there are no methods. This handle exists solely to be dropped. |
| 768 | } |
| 769 | |
| 770 | interface TailStreamTarget $Cxx.allowCancellation { |
| 771 | # Interface used to deliver streaming tail events to a tail worker. |
| 772 | struct TailStreamParams { |
| 773 | events @0 :List(Trace.TailEvent); |
| 774 | } |
| 775 | |
| 776 | struct TailStreamResults { |
| 777 | stop @0 :Bool; |
| 778 | # For an initial tailStream call, the stop flag indicates that the tail worker does |
| 779 | # not wish to continue receiving events. If the stop field is not set, or the value |
| 780 | # is false, then events will be delivered to the tail worker until stop is indicated. |
| 781 | } |
| 782 | |
| 783 | report @0 TailStreamParams -> TailStreamResults; |
| 784 | # Report one or more streaming tail events to a tail worker. |
| 785 | } |
| 786 | |
| 787 | interface EventDispatcher @0xf20697475ec1752d { |
| 788 | # Interface used to deliver events to a Worker's global event handlers. |
| 789 | |
| 790 | getHttpService @0 () -> (http :HttpService) $Cxx.allowCancellation; |
| 791 | # Gets the HTTP interface to this worker (to trigger FetchEvents). |
| 792 | |
| 793 | sendTraces @1 (traces :List(Trace)) -> (result :SendTracesRun) $Cxx.allowCancellation; |
| 794 | # Deliver a trace event to a trace worker. This always completes immediately; the trace handler |
| 795 | # runs as a "waitUntil" task. |
| 796 | |
| 797 | prewarm @2 (url :Text) $Cxx.allowCancellation; |
| 798 | |
| 799 | runScheduled @3 (scheduledTime :Int64, cron :Text) -> (result :ScheduledRun) |
| 800 | $Cxx.allowCancellation; |
| 801 | # Runs a scheduled worker. Returns a ScheduledRun, detailing information about the run such as |
| 802 | # the outcome and whether the run should be retried. This does not complete immediately. |
| 803 | |
| 804 | |
| 805 | runAlarm @4 (scheduledTime :Int64, retryCount :UInt32) -> (result :AlarmRun); |
| 806 | # Runs a worker's alarm. |
| 807 | # scheduledTime is a unix timestamp in milliseconds for when the alarm should be run |
| 808 | # retryCount indicates the retry count, if it's a retry. Else it'll be 0. |
| 809 | # Returns an AlarmRun, detailing information about the run such as |
| 810 | # the outcome and whether the run should be retried. This does not complete immediately. |
| 811 | # |
| 812 | # TODO(cleanup): runAlarm()'s implementation currently relies on *not* allowing cancellation. |
| 813 | # It would be cleaner to handle that inside the implementation so we could mark the entire |
| 814 | # interface (and file) with allowCancellation. |
| 815 | |
| 816 | queue @8 (messages :List(QueueMessage), queueName :Text, metadata :MessageBatchMetadata) |
| 817 | -> (result :QueueResponse) |
| 818 | $Cxx.allowCancellation; |
| 819 | # Delivers a batch of queue messages to a worker's queue event handler. Returns information about |
| 820 | # the success of the batch, including which messages should be considered acknowledged and which |
| 821 | # should be retried. The optional metadata field carries queue metrics at the time the batch was |
| 822 | # dispatched; it is safe for the sender to omit this field (the consumer sees it as absent). |
| 823 | |
| 824 | jsRpcSession @9 () -> (topLevel :JsRpcTarget, session :JsRpcSession) $Cxx.allowCancellation; |
| 825 | # Opens a JS rpc "session". The call does not return until the session is complete. |
| 826 | # |
| 827 | # `topLevel` is the top-level RPC target, on which exactly one method call can be made. This |
| 828 | # call should be made using pipelining to avoid a round trip at startup, and to properly handle |
| 829 | # the old semantics while they still exist in production (see below). |
| 830 | # |
| 831 | # The exact return semantics of this method are currently in flux. Both an old approach and a |
| 832 | # new approach may be live in production: |
| 833 | # * Old approach: `jsRpcSession()` does not return until (1) exactly one call has been made on |
| 834 | # `topLevel`, and (2) any stubs passed over that call (in either direction) have been dropped. |
| 835 | # The session can be canceled by cancelling the call. When the call returns, `session` is null, |
| 836 | # which is consistent with the session being complete. |
| 837 | # * New approach: `jsRpcSession()` returns immediately. The returned `session` capability keeps |
| 838 | # the session alive. Dropping `session` cancels the session. `session` resolves itself to a |
| 839 | # null capability when `topLevel` and all stubs introduced through it have been dropped; the |
| 840 | # caller may await `whenResolved()` to find out when this happens. |
| 841 | # |
| 842 | # The transition will take place in three phases: |
| 843 | # 1. Caller is adjusted to support both approaches. |
| 844 | # 2. Automate is rolled out to switch the callee to the new approach. |
| 845 | # 3. Remove code to support old approach. |
| 846 | # |
| 847 | # In C++, we use `WorkerInterface::customEvent()` to dispatch this event. |
| 848 | |
| 849 | tailStreamSession @10 () -> (topLevel :TailStreamTarget, result :EventOutcome) $Cxx.allowCancellation; |
| 850 | # Opens a streaming tail session. The call does not return until the session is complete. |
| 851 | # |
| 852 | # `topLevel` is the top-level tail session target, on which exactly one method call can |
| 853 | # be made. This call must be made using pipelining since `tailStreamSession()` won't return |
| 854 | # until after the call completes. result is accessed after the session is complete. |
| 855 | |
| 856 | obsolete5 @5(); |
| 857 | obsolete6 @6(); |
| 858 | obsolete7 @7(); |
| 859 | # Deleted methods, do not reuse these numbers. |
| 860 | |
| 861 | abandonAlarm @11 (scheduledTimeMs :Int64) -> (storedAlarmTimeMs :Int64) $Cxx.allowCancellation; |
| 862 | # Called by AlarmManager when it has given up retrying an alarm after too many counted failures. |
| 863 | # The actor should clear its alarm state (ActorCache / SQLite) so that getAlarm() correctly |
| 864 | # reflects that no alarm will ever fire for this scheduled time again. |
| 865 | # |
| 866 | # Returns the actor's current alarm time (in ms since epoch) if it differs from scheduledTimeMs, |
| 867 | # indicating the user set a new alarm that should still fire. Returns 0 if the alarm was cleared |
| 868 | # or no alarm was stored. |
| 869 | |
| 870 | # Other methods might be added to handle other kinds of events, e.g. TCP connections, or maybe |
| 871 | # even native Cap'n Proto RPC eventually. |
| 872 | } |
| 873 | |
| 874 | interface WorkerdBootstrap { |
| 875 | # Bootstrap interface exposed by workerd when serving Cap'n Proto RPC. |
| 876 | |
| 877 | startEvent @0 (cfBlobJson :Text) -> (dispatcher :EventDispatcher); |
| 878 | # Start a new event. Exactly one event should be delivered to the returned EventDispatcher. |
| 879 | # |
| 880 | # If the event is an HTTP request, `cfBlobJson` optionally carries the JSON-encoded `request.cf` |
| 881 | # object. The dispatcher will pass it through to the worker via SubrequestMetadata. |
| 882 | } |
| 883 | |
| 884 | interface WorkerdDebugPort { |
| 885 | # Bootstrap interface exposed on the debug RPC port, if one is configured. This exposes access |
| 886 | # to all services in the process, with the ability for the client to specify arbitrary props, so |
| 887 | # this interface should be considered privileged, and should probably only be used for testing |
| 888 | # purposes. |
| 889 | # |
| 890 | # This interface is subject to change. It is intended for use by miniflare. |
| 891 | |
| 892 | getEntrypoint @0 (service :Text, entrypoint :Text, props :Frankenvalue) |
| 893 | -> (entrypoint :WorkerdBootstrap); |
| 894 | # Get direct access to a stateless entrypoint. |
| 895 | |
| 896 | getActor @1 (service :Text, entrypoint :Text, actorId :Text) -> (actor :WorkerdBootstrap); |
| 897 | # Get an actor (Durable Object) stub. |
| 898 | # The actorId should be a hex string for Durable Objects or a plain string for ephemeral actors. |
| 899 | } |