Skip to content
File

Blob: src/rust/worker/ffi.rs

rust247 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 
5use std::pin::Pin;
6 
7use kj::http::ConnectSettings;
8use kj_rs::KjDate;
9 
10use crate::CustomEvent;
11use crate::Result;
12use crate::ffi::bridge::AlarmResult;
13use crate::ffi::bridge::CustomEventResult;
14use crate::ffi::bridge::ScheduledResult;
15 
16#[cxx::bridge(namespace = "workerd::rust::worker")]
17pub mod bridge {
18 #[namespace = "kj::rust"]
19 unsafe extern "C++" {
20 type HttpMethod = kj::http::ffi::HttpMethod;
21 type HttpHeaders = kj::http::ffi::HttpHeaders;
22 type HttpServiceResponse = kj::http::ffi::HttpServiceResponse;
23 type AsyncInputStream = kj::io::ffi::AsyncInputStream;
24 type AsyncIoStream = kj::io::ffi::AsyncIoStream;
25 type ConnectResponse = kj::http::ffi::ConnectResponse;
26 }
27 
28 #[derive(Debug, Clone, PartialEq, Eq)]
29 pub enum EventOutcome {
30 Unknown = 0,
31 Ok = 1,
32 Exception = 2,
33 ExceededCpu = 3,
34 KillSwitch = 4,
35 DaemonDown = 5,
36 ScriptNotFound = 6,
37 Canceled = 7,
38 ExceededMemory = 8,
39 LoadShed = 9,
40 ResponseStreamDisconnected = 10,
41 InternalError = 11,
42 }
43 
44 #[derive(Debug, Clone, PartialEq, Eq)]
45 // keep in sync with `src/workerd/io/worker-interface.h`
46 pub struct ScheduledResult {
47 pub retry: bool,
48 pub outcome: EventOutcome,
49 }
50 
51 #[derive(Debug, Clone, PartialEq, Eq)]
52 // keep in sync with `src/workerd/io/worker-interface.h`
53 pub struct AlarmResult {
54 pub retry: bool,
55 pub retry_counts_against_limit: bool,
56 pub outcome: EventOutcome,
57 pub error_description: String,
58 }
59 
60 #[derive(Debug, Clone, PartialEq, Eq)]
61 // keep in sync with `src/workerd/io/worker-interface.h`
62 pub struct CustomEventResult {
63 pub outcome: EventOutcome,
64 }
65 
66 #[namespace = "kj::rust"]
67 extern "C++" {
68 type HttpConnectSettings<'a> = kj::http::ConnectSettings<'a>;
69 }
70 
71 extern "Rust" {
72 type Wrapper;
73 
74 async unsafe fn request<'a>(
75 self: &'a mut Wrapper,
76 method: HttpMethod,
77 url: &'a [u8],
78 headers: &'a HttpHeaders,
79 request_body: Pin<&'a mut AsyncInputStream>,
80 response: Pin<&'a mut HttpServiceResponse>,
81 ) -> Result<()>;
82 
83 async unsafe fn connect<'a>(
84 self: &'a mut Wrapper,
85 host: &'a [u8],
86 headers: &'a HttpHeaders,
87 connection: Pin<&'a mut AsyncIoStream>,
88 response: Pin<&'a mut ConnectResponse>,
89 settings: HttpConnectSettings<'a>,
90 ) -> Result<()>;
91 
92 async unsafe fn prewarm<'a>(self: &'a mut Wrapper, url: &'a [u8]) -> Result<()>;
93 
94 async unsafe fn run_scheduled<'a>(
95 self: &'a mut Wrapper,
96 scheduled_time: KjDate,
97 cron: &'a [u8],
98 ) -> Result<ScheduledResult>;
99 
100 async unsafe fn run_alarm<'a>(
101 self: &'a mut Wrapper,
102 scheduled_time: KjDate,
103 retry_count: u32,
104 ) -> Result<AlarmResult>;
105 
106 async unsafe fn custom_event<'a>(
107 self: &'a mut Wrapper,
108 event: Pin<&'a mut CustomEvent>,
109 ) -> Result<CustomEventResult>;
110 
111 async unsafe fn test<'a>(self: &'a mut Wrapper) -> Result<bool>;
112 }
113 
114 unsafe extern "C++" {
115 include!("workerd/rust/worker/ffi.h");
116 }
117 
118 unsafe extern "C++" {
119 type CustomEvent;
120 }
121 
122 impl Box<Wrapper> {}
123}
124 
125pub struct Wrapper {
126 worker: Box<dyn crate::Interface>,
127}
128 
129impl Wrapper {
130 pub(crate) fn new(worker: Box<dyn crate::Interface>) -> Self {
131 Self { worker }
132 }
133 
134 async fn request<'a>(
135 &'a mut self,
136 method: bridge::HttpMethod,
137 url: &'a [u8],
138 headers: &'a bridge::HttpHeaders,
139 request_body: Pin<&'a mut bridge::AsyncInputStream>,
140 response: Pin<&'a mut bridge::HttpServiceResponse>,
141 ) -> Result<()> {
142 let response = crate::ServiceResponse::from(response);
143 self.worker
144 .request(method, url, headers.into(), request_body, response)
145 .await?;
146 Ok(())
147 }
148 
149 async fn connect<'a>(
150 &'a mut self,
151 host: &'a [u8],
152 headers: &'a bridge::HttpHeaders,
153 connection: Pin<&'a mut bridge::AsyncIoStream>,
154 response: Pin<&'a mut bridge::ConnectResponse>,
155 settings: ConnectSettings<'a>,
156 ) -> Result<()> {
157 let response = crate::ConnectResponse::from(response);
158 self.worker
159 .connect(host, headers.into(), connection, response, settings)
160 .await?;
161 Ok(())
162 }
163 
164 async fn prewarm(&mut self, url: &[u8]) -> Result<()> {
165 self.worker.prewarm(&String::from_utf8_lossy(url)).await?;
166 Ok(())
167 }
168 
169 async fn run_scheduled(
170 &mut self,
171 scheduled_time: KjDate,
172 cron: &[u8],
173 ) -> Result<ScheduledResult> {
174 let scheduled_time = scheduled_time.into();
175 let result = self
176 .worker
177 .run_scheduled(&scheduled_time, &String::from_utf8_lossy(cron))
178 .await?;
179 Ok(result.into())
180 }
181 
182 async fn run_alarm(&mut self, scheduled_time: KjDate, retry_count: u32) -> Result<AlarmResult> {
183 let scheduled_time = scheduled_time.into();
184 let result = self.worker.run_alarm(&scheduled_time, retry_count).await?;
185 Ok(result.into())
186 }
187 
188 async fn custom_event(&mut self, event: Pin<&mut CustomEvent>) -> Result<CustomEventResult> {
189 let result = self.worker.custom_event(event).await?;
190 Ok(result.into())
191 }
192 
193 async fn test(&mut self) -> Result<bool> {
194 let result = self.worker.test().await?;
195 Ok(result)
196 }
197}
198 
199impl From<crate::ScheduledResult> for bridge::ScheduledResult {
200 fn from(value: crate::ScheduledResult) -> Self {
201 Self {
202 retry: value.retry,
203 outcome: value.outcome.into(),
204 }
205 }
206}
207 
208impl From<crate::AlarmResult> for bridge::AlarmResult {
209 fn from(value: crate::AlarmResult) -> Self {
210 Self {
211 retry: value.retry,
212 retry_counts_against_limit: value.retry_counts_against_limit,
213 outcome: value.outcome.into(),
214 error_description: value.error_description.unwrap_or_default(),
215 }
216 }
217}
218 
219impl From<crate::CustomEventResult> for bridge::CustomEventResult {
220 fn from(value: crate::CustomEventResult) -> Self {
221 Self {
222 outcome: value.outcome.into(),
223 }
224 }
225}
226 
227impl From<outcome_capnp::EventOutcome> for bridge::EventOutcome {
228 fn from(value: outcome_capnp::EventOutcome) -> Self {
229 match value {
230 outcome_capnp::EventOutcome::Unknown => Self::Unknown,
231 outcome_capnp::EventOutcome::Ok => Self::Ok,
232 outcome_capnp::EventOutcome::Exception => Self::Exception,
233 outcome_capnp::EventOutcome::ExceededCpu => Self::ExceededCpu,
234 outcome_capnp::EventOutcome::KillSwitch => Self::KillSwitch,
235 outcome_capnp::EventOutcome::DaemonDown => Self::DaemonDown,
236 outcome_capnp::EventOutcome::ScriptNotFound => Self::ScriptNotFound,
237 outcome_capnp::EventOutcome::Canceled => Self::Canceled,
238 outcome_capnp::EventOutcome::ExceededMemory => Self::ExceededMemory,
239 outcome_capnp::EventOutcome::LoadShed => Self::LoadShed,
240 outcome_capnp::EventOutcome::ResponseStreamDisconnected => {
241 Self::ResponseStreamDisconnected
242 }
243 outcome_capnp::EventOutcome::InternalError => Self::InternalError,
244 }
245 }
246}