Skip to content
File

Blob: src/rust/worker/lib.rs

rust165 lines
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 
16pub mod error;
17pub mod exception;
18pub mod ffi;
19pub mod kill_switch;
20pub mod ok;
21 
22use std::pin::Pin;
23use std::time::SystemTime;
24 
25use cxx::KjError;
26pub use kj::http::ConnectResponse;
27pub use kj::http::ConnectSettings;
28pub use kj::http::HeaderId;
29pub use kj::http::Headers;
30pub use kj::http::HeadersRef;
31pub use kj::http::Method;
32pub use kj::http::Service;
33pub use kj::http::ServiceResponse;
34pub use kj::io::AsyncInputStream;
35pub use kj::io::AsyncIoStream;
36use outcome_capnp::EventOutcome;
37 
38pub use crate::ffi::Wrapper;
39pub use crate::ffi::bridge::CustomEvent;
40 
41pub 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)]
46pub 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)]
108pub struct ScheduledResult {
109 pub retry: bool,
110 pub outcome: EventOutcome,
111}
112 
113/// Result of an alarm event.
114#[derive(Debug, Clone, PartialEq, Eq)]
115pub 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)]
124pub struct CustomEventResult {
125 pub outcome: EventOutcome,
126}
127 
128#[cfg(test)]
129mod 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}