File
Blob: src/rust/worker/lib.rs
| 1 | // Copyright (c) 2026 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 | //! Worker interface crate. |
| 6 | //! |
| 7 | //! This crate contains Rust's counterpart of `workerd::WorkerInterface`, necessary ffi bindings |
| 8 | //! and several simple interface implementations. |
| 9 | //! |
| 10 | //! The intended way to use this crate is to `use worker` and reference `worker::Interface` trait |
| 11 | //! in your code. |
| 12 | //! |
| 13 | //! Current limitations: |
| 14 | //! - most parameters are opaque c++ types |
| 15 | |
| 16 | pub mod error; |
| 17 | pub mod exception; |
| 18 | pub mod ffi; |
| 19 | pub mod kill_switch; |
| 20 | pub mod ok; |
| 21 | |
| 22 | use std::pin::Pin; |
| 23 | use std::time::SystemTime; |
| 24 | |
| 25 | use cxx::KjError; |
| 26 | pub use kj::http::ConnectResponse; |
| 27 | pub use kj::http::ConnectSettings; |
| 28 | pub use kj::http::HeaderId; |
| 29 | pub use kj::http::Headers; |
| 30 | pub use kj::http::HeadersRef; |
| 31 | pub use kj::http::Method; |
| 32 | pub use kj::http::Service; |
| 33 | pub use kj::http::ServiceResponse; |
| 34 | pub use kj::io::AsyncInputStream; |
| 35 | pub use kj::io::AsyncIoStream; |
| 36 | use outcome_capnp::EventOutcome; |
| 37 | |
| 38 | pub use crate::ffi::Wrapper; |
| 39 | pub use crate::ffi::bridge::CustomEvent; |
| 40 | |
| 41 | pub type Result<T> = std::result::Result<T, KjError>; |
| 42 | |
| 43 | /// An interface representing the services made available by a worker/pipeline to handle a request. |
| 44 | /// Corresponds to `workerd::WorkerInterface` |
| 45 | #[async_trait::async_trait(?Send)] |
| 46 | pub trait Interface: kj::http::Service { |
| 47 | /// Hints that this worker will likely be invoked in the near future, so should be warmed up now. |
| 48 | /// This method should also call `prewarm()` on any subsequent pipeline stages that are expected |
| 49 | /// to be invoked. |
| 50 | /// |
| 51 | /// If `prewarm()` has to do anything asynchronous, it should use "waitUntil" tasks. |
| 52 | async fn prewarm(&mut self, _url: &str) -> Result<()> { |
| 53 | Ok(()) |
| 54 | } |
| 55 | |
| 56 | /// Trigger a scheduled event with the given scheduled (unix timestamp) time and cron string. |
| 57 | /// The cron string must be valid until the returned promise completes. |
| 58 | /// Async work is queued in a "waitUntil" task set. |
| 59 | async fn run_scheduled( |
| 60 | &mut self, |
| 61 | scheduled_time: &SystemTime, |
| 62 | cron: &str, |
| 63 | ) -> Result<ScheduledResult>; |
| 64 | |
| 65 | /// Trigger an alarm event with the given scheduled (unix timestamp) time. |
| 66 | async fn run_alarm( |
| 67 | &mut self, |
| 68 | scheduled_time: &SystemTime, |
| 69 | retry_count: u32, |
| 70 | ) -> Result<AlarmResult>; |
| 71 | |
| 72 | /// Run the test handler. The returned promise resolves to true or false to indicate that the test |
| 73 | /// passed or failed. In the case of a failure, information should have already been written to |
| 74 | /// stderr and to the devtools; there is no need for the caller to write anything further. (If the |
| 75 | /// promise rejects, this indicates a bug in the test harness itself.) |
| 76 | async fn test(&mut self) -> Result<bool> { |
| 77 | Ok(false) |
| 78 | } |
| 79 | |
| 80 | /// Allows delivery of a variety of event types by implementing a callback that delivers the |
| 81 | /// event to a particular isolate. If and when the event is delivered to an isolate, |
| 82 | /// `callback->run()` will be called inside a fresh `IoContext::IncomingRequest` to begin the |
| 83 | /// event. |
| 84 | /// |
| 85 | /// If the event needs to return some sort of result, it's the responsibility of the callback to |
| 86 | /// store that result in a side object that the event's invoker can inspect after the promise has |
| 87 | /// resolved. |
| 88 | /// |
| 89 | /// Note that it is guaranteed that if the returned promise is canceled, `event` will be dropped |
| 90 | /// immediately; if its callbacks have not run yet, they will not run at all. So, a `CustomEvent` |
| 91 | /// implementation can hold references to objects it doesn't own as long as the returned promise |
| 92 | /// will be canceled before those objects go away. |
| 93 | async fn custom_event(&mut self, event: Pin<&mut CustomEvent>) -> Result<CustomEventResult>; |
| 94 | |
| 95 | /// Convert `self` into the structure suitable for passing to C++ through FFI layer |
| 96 | /// To obtain `workerd::WorkerInterface` on C++ side finish wrapping with `fromRust()` |
| 97 | /// call defined in `bridge.h`. |
| 98 | fn into_ffi(self) -> Box<ffi::Wrapper> |
| 99 | where |
| 100 | Self: Sized + 'static, |
| 101 | { |
| 102 | Box::new(ffi::Wrapper::new(Box::new(self))) |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | /// Result of a scheduled event. |
| 107 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 108 | pub struct ScheduledResult { |
| 109 | pub retry: bool, |
| 110 | pub outcome: EventOutcome, |
| 111 | } |
| 112 | |
| 113 | /// Result of an alarm event. |
| 114 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 115 | pub struct AlarmResult { |
| 116 | pub retry: bool, |
| 117 | pub retry_counts_against_limit: bool, |
| 118 | pub outcome: EventOutcome, |
| 119 | pub error_description: Option<String>, |
| 120 | } |
| 121 | |
| 122 | /// Result of a custom event. |
| 123 | #[derive(Debug, Clone, PartialEq, Eq)] |
| 124 | pub struct CustomEventResult { |
| 125 | pub outcome: EventOutcome, |
| 126 | } |
| 127 | |
| 128 | #[cfg(test)] |
| 129 | mod tests { |
| 130 | use kj::http::HeaderId; |
| 131 | |
| 132 | #[expect(dead_code, reason = "HeaderOverrideProxy exists as compilation test.")] |
| 133 | struct HeaderOverrideProxy(Box<dyn crate::Interface>); |
| 134 | |
| 135 | #[async_trait::async_trait(?Send)] |
| 136 | impl kj::http::Service for HeaderOverrideProxy { |
| 137 | async fn request<'a>( |
| 138 | &'a mut self, |
| 139 | method: kj::http::Method, |
| 140 | url: &'a [u8], |
| 141 | headers: kj::http::HeadersRef<'a>, |
| 142 | request_body: std::pin::Pin<&'a mut kj::io::AsyncInputStream>, |
| 143 | response: kj::http::ServiceResponse<'a>, |
| 144 | ) -> kj::Result<()> { |
| 145 | let mut headers = headers.clone_shallow(); |
| 146 | headers.set(HeaderId::HOST, "example.com"); |
| 147 | self.0 |
| 148 | .request(method, url, headers.as_ref(), request_body, response) |
| 149 | .await?; |
| 150 | Ok(()) |
| 151 | } |
| 152 | |
| 153 | async fn connect<'a>( |
| 154 | &'a mut self, |
| 155 | _host: &'a [u8], |
| 156 | _headers: kj::http::HeadersRef<'a>, |
| 157 | _connection: std::pin::Pin<&'a mut kj::io::AsyncIoStream>, |
| 158 | _response: kj::http::ConnectResponse<'a>, |
| 159 | _settings: kj::http::ConnectSettings<'a>, |
| 160 | ) -> kj::Result<()> { |
| 161 | todo!() |
| 162 | } |
| 163 | } |
| 164 | } |