File
Blob: src/rust/worker/ffi.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 | use std::pin::Pin; |
| 6 | |
| 7 | use kj::http::ConnectSettings; |
| 8 | use kj_rs::KjDate; |
| 9 | |
| 10 | use crate::CustomEvent; |
| 11 | use crate::Result; |
| 12 | use crate::ffi::bridge::AlarmResult; |
| 13 | use crate::ffi::bridge::CustomEventResult; |
| 14 | use crate::ffi::bridge::ScheduledResult; |
| 15 | |
| 16 | #[cxx::bridge(namespace = "workerd::rust::worker")] |
| 17 | pub 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 | |
| 125 | pub struct Wrapper { |
| 126 | worker: Box<dyn crate::Interface>, |
| 127 | } |
| 128 | |
| 129 | impl 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 | |
| 199 | impl 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 | |
| 208 | impl 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 | |
| 219 | impl From<crate::CustomEventResult> for bridge::CustomEventResult { |
| 220 | fn from(value: crate::CustomEventResult) -> Self { |
| 221 | Self { |
| 222 | outcome: value.outcome.into(), |
| 223 | } |
| 224 | } |
| 225 | } |
| 226 | |
| 227 | impl 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 | } |