Skip to content

Commit d65c158

Browse files
committed
fix: server tx dies after getting
1 parent ed619b4 commit d65c158

11 files changed

Lines changed: 90 additions & 61 deletions

File tree

Cargo.lock

Lines changed: 1 addition & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

crates/svld-binding-kv/src/main.rs

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,8 @@ async fn main() -> Result<(), Box<dyn core::error::Error>> {
6767
}
6868
};
6969

70+
tracing::info!("worker {:?}: {}(...)", worker_name, function_name);
71+
7072
let tree = match db.open_tree(&worker_name) {
7173
Ok(t) => t,
7274
Err(e) => {
@@ -122,14 +124,15 @@ async fn main() -> Result<(), Box<dyn core::error::Error>> {
122124

123125
match &payload {
124126
KvPayload::Get { key } => {
127+
tracing::info!("sending...");
125128
match tree.get(&key) {
126129
Ok(Some(data)) => {
127130
let s = String::from_utf8_lossy(&data);
128131
client.send_ok(id, s).await?;
129132
}
130133

131134
Ok(None) => {
132-
client.send_ok(id, payload).await?;
135+
client.send_ok(id, ijson::ijson!(null)).await?;
133136
}
134137

135138
Err(e) => {

crates/svld-ipc-client/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,3 +12,4 @@ serde = "1.0.228"
1212
serde_json = "1.0.150"
1313
thiserror = "2.0.18"
1414
tokio = { version = "1.52.3" }
15+
tracing = "0.1.44"

crates/svld-ipc-client/src/core.rs

Lines changed: 15 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -4,12 +4,7 @@
44
//! let mut client = BindingClient::connect(PathBuf::from(".serverlessd/bindings.sock")).await?;
55
//! ```
66
7-
use std::{
8-
io::{self, IoSlice},
9-
marker::PhantomData,
10-
mem,
11-
path::PathBuf,
12-
};
7+
use std::{io, marker::PhantomData, mem, path::PathBuf};
138

149
use interprocess::local_socket::{
1510
ConnectOptions, GenericFilePath,
@@ -75,21 +70,17 @@ impl BindingClient<Initialized> {
7570
&mut self,
7671
message: ClientMessage<Result<T, String>>,
7772
) -> Result<(), BindingClientError> {
78-
let id_raw = message.id.to_le_bytes();
7973
let payload_raw = serde_json::to_vec(&match message.payload {
8074
Ok(t) => serde_json::json!({"data": &t}),
8175
Err(e) => serde_json::json!({"error": &e}),
8276
})?;
8377

84-
let len_raw = (payload_raw.len() as u32).to_le_bytes();
85-
86-
let slices = [
87-
IoSlice::new(&id_raw),
88-
IoSlice::new(&len_raw),
89-
IoSlice::new(&payload_raw),
90-
];
78+
let mut buf = Vec::with_capacity(size_of::<u32>() * 2 + payload_raw.len());
79+
buf.extend_from_slice(&message.id.to_le_bytes());
80+
buf.extend_from_slice(&(payload_raw.len() as u32).to_le_bytes());
81+
buf.extend_from_slice(&payload_raw);
9182

92-
self.send.write_vectored(&slices).await?;
83+
self.send.write_all(&buf).await?;
9384

9485
Ok(())
9586
}
@@ -129,14 +120,19 @@ impl BindingClient<Initialized> {
129120
&mut self,
130121
) -> Result<ServerMessage<T>, BindingClientError> {
131122
let id = self.recv.read_u32_le().await?;
123+
tracing::info!("got id={id}");
132124

133125
let data_len = self.recv.read_u32_le().await? as usize;
134126
let mut data = Box::<[u8]>::new_uninit_slice(data_len);
135127
let slice =
136128
unsafe { std::slice::from_raw_parts_mut(data.as_mut_ptr() as *mut u8, data_len) };
137129
self.recv.read_exact(slice).await?;
138130

139-
let UnprocessedServerMessage { func, data, worker } = serde_json::from_slice(slice)?;
131+
let UnprocessedServerMessage {
132+
func,
133+
args: data,
134+
worker,
135+
} = serde_json::from_slice(slice)?;
140136

141137
Ok(ServerMessage {
142138
id,
@@ -145,6 +141,8 @@ impl BindingClient<Initialized> {
145141
args: ijson::from_value::<T>(&data)?,
146142
})
147143
}
144+
145+
// pub async fn shutdown(&mut self) {}
148146
}
149147

150148
#[derive(Debug, thiserror::Error)]
@@ -173,7 +171,7 @@ pub struct ClientMessage<T> {
173171
#[derive(Deserialize)]
174172
struct UnprocessedServerMessage {
175173
func: String,
176-
data: ijson::IValue,
174+
args: ijson::IValue,
177175
worker: String,
178176
}
179177

crates/svld-rt/src/bindings/backend.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -28,7 +28,7 @@ pub trait BindingBackend: Send + Sync + 'static {
2828
fn get_tx(&self) -> BindingBackendTx;
2929

3030
/// Creates a client from the env name for the worker.
31-
fn create_client(&self, binding_name: &str) -> Box<dyn BindingClient>;
31+
fn create_client(&self) -> Box<dyn BindingClient>;
3232
}
3333

3434
/// Creates a binding backend channel.

crates/svld-rt/src/bindings/ipc/core.rs

Lines changed: 61 additions & 35 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@ impl IpcBindingsServer {
3030
let binding_backends = binding_types
3131
.into_iter()
3232
.map(|type_| {
33-
let backend = Arc::new(binding_backend::IpcBindingBackend::new());
33+
let backend = Arc::new(binding_backend::IpcBindingBackend::new(type_.clone()));
3434
binding_store.push_binding(&type_, backend.clone());
3535

3636
(type_, backend)
@@ -73,6 +73,7 @@ pub mod binding_backend {
7373
use crate::bindings::{BindingBackend, BindingBackendTx, backend::BindingClient};
7474

7575
pub struct IpcBindingBackend {
76+
type_: String,
7677
maybe_tx: (
7778
AtomicBool, // has data?
7879
AtomicUsize, // BindingBackendTx
@@ -86,8 +87,9 @@ pub mod binding_backend {
8687

8788
impl IpcBindingBackend {
8889
#[inline(always)]
89-
pub const fn new() -> Self {
90+
pub fn new(type_: String) -> Self {
9091
Self {
92+
type_,
9193
maybe_tx: (AtomicBool::new(false), AtomicUsize::new(0)),
9294
}
9395
}
@@ -137,7 +139,10 @@ pub mod binding_backend {
137139

138140
// SAFETY: we asserted BindingBackendTx is usize-sized;
139141
// the bool guard ensures this was previously set via set_tx.
140-
Some(unsafe { mem::transmute::<usize, BindingBackendTx>(raw) })
142+
Some(unsafe {
143+
let tx = mem::ManuallyDrop::new(mem::transmute::<usize, BindingBackendTx>(raw));
144+
(*tx).clone()
145+
})
141146
}
142147

143148
#[inline(always)]
@@ -175,8 +180,8 @@ pub mod binding_backend {
175180
}
176181

177182
#[inline]
178-
fn create_client(&self, binding_name: &str) -> Box<dyn BindingClient> {
179-
Box::new(IpcBindingClient::new(binding_name.to_string()))
183+
fn create_client(&self) -> Box<dyn BindingClient> {
184+
Box::new(IpcBindingClient::new(self.type_.clone()))
180185
}
181186
}
182187

@@ -204,13 +209,13 @@ mod binding_client {
204209
};
205210

206211
pub struct IpcBindingClient {
207-
binding_name: String,
212+
binding_type: String,
208213
}
209214

210215
impl IpcBindingClient {
211216
#[inline(always)]
212-
pub fn new(binding_name: String) -> Self {
213-
Self { binding_name }
217+
pub fn new(binding_type: String) -> Self {
218+
Self { binding_type }
214219
}
215220
}
216221

@@ -232,7 +237,7 @@ mod binding_client {
232237
mut rv: v8::ReturnValue| {
233238
let arr = args.data().cast::<v8::Array>();
234239

235-
let binding_name = arr
240+
let binding_type = arr
236241
.get_index(scope, 0)
237242
.unwrap()
238243
.cast::<v8::String>()
@@ -244,7 +249,7 @@ mod binding_client {
244249
to_rust_string_lossy(scope);
245250

246251
let state = WorkerState::get_from_isolate(scope);
247-
let tx = state.get_binding(&binding_name).unwrap();
252+
let tx = state.get_binding_tx(&binding_type).unwrap();
248253

249254
let args_len = args.length();
250255
let arr = v8::Array::new(scope, args_len);
@@ -331,7 +336,7 @@ mod binding_client {
331336
rv.set(fnk.cast());
332337
},
333338
)
334-
.data(v8::String::new(scope, &self.binding_name)?.cast())
339+
.data(v8::String::new(scope, &self.binding_type)?.cast())
335340
.build(scope)?;
336341

337342
handler.set(scope, v8::String::new(scope, "get")?.cast(), get_fn.cast());
@@ -348,7 +353,7 @@ mod binding_client {
348353
mod task {
349354
use std::{
350355
collections::HashMap,
351-
io::{self, IoSlice},
356+
io,
352357
string::FromUtf8Error,
353358
sync::{
354359
Arc,
@@ -432,6 +437,28 @@ mod task {
432437
let (tx, mut rx) = binding_backend_channel();
433438
backend.set_tx(tx)?;
434439

440+
// spawn dedicated reader so it's never cancelled
441+
let (external_tx, mut external_rx) = tokio::sync::mpsc::channel::<Message>(32);
442+
443+
tokio::spawn(async move {
444+
loop {
445+
match recv.read_parse_message().await {
446+
Ok(msg) => {
447+
if external_tx.send(msg).await.is_err() {
448+
tracing::error!("errored on sending external tx");
449+
break; // main task dropped, exit
450+
}
451+
}
452+
Err(e) => {
453+
tracing::error!(
454+
"failed to read & parse binding client message, breaking: {e:?}"
455+
);
456+
break;
457+
}
458+
}
459+
}
460+
});
461+
435462
let mut resolutions = HashMap::new();
436463
let mut roll_id = 0_u32;
437464

@@ -445,22 +472,20 @@ mod task {
445472
}
446473

447474
let event = tokio::select! {
448-
msg = rx.recv() => {
449-
match msg {
450-
Some(t) => Event::Internal(t),
451-
None => break,
475+
msg = rx.recv() => match msg {
476+
Some(t) => Event::Internal(t),
477+
None => {
478+
tracing::error!("internal channel closed: backend dropped?");
479+
break;
452480
}
453481
},
454-
455-
msg = recv.read_parse_message() => {
456-
match msg {
457-
Ok(t) => Event::External(t),
458-
Err(e) => {
459-
tracing::error!("failed to read & parse binding sclient message, breaking: {e:?}");
460-
break;
461-
}
482+
msg = external_rx.recv() => match msg {
483+
Some(t) => Event::External(t),
484+
None => {
485+
tracing::error!("reader task died");
486+
break;
462487
}
463-
}
488+
},
464489
};
465490

466491
match event {
@@ -480,6 +505,7 @@ mod task {
480505
}
481506

482507
Event::External(Message { id, payload }) => {
508+
tracing::debug!("server finished reading payload for id={id}");
483509
if let Some(replier) = resolutions.remove(&id) {
484510
let _ = replier.send(payload);
485511
}
@@ -536,6 +562,7 @@ mod task {
536562
// get this over with, and we're just reading from them
537563

538564
let len = self.read_u32_le().await? as usize;
565+
539566
let mut buf = Box::<[u8]>::new_uninit_slice(len);
540567

541568
let slice = unsafe { std::slice::from_raw_parts_mut(buf.as_mut_ptr() as *mut u8, len) };
@@ -557,6 +584,7 @@ mod task {
557584

558585
// then we'll get the payload
559586
let raw_payload = self.read_parse_to_boxed_arr().await?;
587+
560588
let payload = serde_json::from_slice::<ijson::IValue>(&raw_payload)?;
561589

562590
// good. fuck you and eat it up
@@ -572,17 +600,15 @@ mod task {
572600
#[async_trait]
573601
impl SendExt for SendHalf {
574602
async fn send_message(&mut self, message: Message) -> Result<(), SingleTaskError> {
575-
let id_raw = message.id.to_le_bytes();
576-
577603
let payload_raw = serde_json::to_vec(&message.payload)?;
578-
let len_raw = (payload_raw.len() as u32).to_le_bytes();
579-
580-
let slices = [
581-
IoSlice::new(&id_raw),
582-
IoSlice::new(&len_raw),
583-
IoSlice::new(&payload_raw),
584-
];
585-
self.write_vectored(&slices).await?;
604+
605+
let mut buf = Vec::with_capacity(size_of::<u32>() * 2 + payload_raw.len());
606+
607+
buf.extend_from_slice(&message.id.to_le_bytes());
608+
buf.extend_from_slice(&(payload_raw.len() as u32).to_le_bytes());
609+
buf.extend_from_slice(&payload_raw);
610+
611+
self.write_all(&buf).await?;
586612

587613
// fah
588614
Ok(())

crates/svld-rt/src/bindings/store.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -106,8 +106,8 @@ pub struct BindingItem {
106106

107107
impl BindingItem {
108108
#[inline(always)]
109-
pub fn create_client(&self, name: &str) -> Box<dyn BindingClient> {
110-
self.backend.create_client(name)
109+
pub fn create_client(&self) -> Box<dyn BindingClient> {
110+
self.backend.create_client()
111111
}
112112

113113
#[inline(always)]

crates/svld-rt/src/env.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -16,7 +16,7 @@ pub fn create_js_env<'s>(
1616
tracing::info!("js env got name={name:?}, type={type_:?}");
1717
let binding = state.binding_store.get_binding(&type_)?;
1818

19-
let client = binding.create_client(&name);
19+
let client = binding.create_client();
2020

2121
let interface = client.create_interface(scope)?;
2222

crates/svld-rt/src/worker/state.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -201,7 +201,7 @@ impl WorkerState {
201201
}
202202

203203
#[inline(always)]
204-
pub fn get_binding(&self, name: &str) -> Option<BindingBackendTx> {
205-
self.binding_store.get_binding_tx(name)
204+
pub fn get_binding_tx(&self, type_: &str) -> Option<BindingBackendTx> {
205+
self.binding_store.get_binding_tx(type_)
206206
}
207207
}

crates/svld-rt/src/worker/task.rs

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -295,7 +295,6 @@ async fn create_task(rx: &mut WorkerRx, args: InitWorkerArgs<'_>) -> Result<bool
295295
tracing::info!("resolving event loop");
296296

297297
let isolate = unsafe { state.get_isolate() };
298-
tracing::info!("creating scope with context");
299298
scope_with_context!(
300299
isolate: isolate,
301300
let &mut scope,

0 commit comments

Comments
 (0)