1pub trait ExtensionProvider {
2 fn name(&self) -> &str;
3 fn install(&self, protocols: &mut ProtocolRegistry);
4 fn construct(&self, type_name: &str, arguments: &[Value]) -> Result<Value, String>;
5}
6
7#[derive(Default, Clone)]
8pub struct ExtensionRegistry {
9 providers: HashMap<String, Rc<dyn ExtensionProvider>>,
10 loaded: HashSet<String>,
11}
12
13impl ExtensionRegistry {
14 pub fn new() -> Self {
15 Self::default()
16 }
17
18 pub fn install<P: ExtensionProvider + 'static>(&mut self, provider: P) {
19 self.providers
20 .insert(provider.name().to_string(), Rc::new(provider));
21 }
22
23 pub fn contains(&self, name: &str) -> bool {
24 self.providers.contains_key(name)
25 }
26
27 pub fn require(
28 &mut self,
29 name: &str,
30 protocols: &mut ProtocolRegistry,
31 ) -> Result<String, String> {
32 let provider = self
33 .providers
34 .get(name)
35 .cloned()
36 .ok_or_else(|| format!("extension/not-found: {name}"))?;
37 if self.loaded.insert(name.to_string()) {
38 provider.install(protocols);
39 }
40 Ok(if self.loaded.len() == 1 {
41 ":loaded".into()
42 } else {
43 ":loaded".into()
44 })
45 }
46
47 pub fn construct(
48 &self,
49 provider: &str,
50 type_name: &str,
51 arguments: &[Value],
52 ) -> Result<Value, String> {
53 self.providers
54 .get(provider)
55 .ok_or_else(|| format!("extension/not-found: {provider}"))?
56 .construct(type_name, arguments)
57 }
58}
59
60#[cfg(any(not(target_arch = "wasm32"), target_os = "wasi"))]
61pub use crate::file::NativeFileProvider;
62pub use crate::file::{
63 CopyOptions, DeleteOptions, FileEntry, FileError, FileProvider, FileType, MemoryFileProvider,
64 MkdirOptions, MoveOptions, TempDirectoryOptions, TempFileOptions, UnsupportedFileProvider,
65 WriteMode, WriteOptions,
66};
67
68#[derive(Debug, Clone, PartialEq)]
69pub enum SocketError {
70 Unsupported,
71 Denied,
72 Invalid(String),
73}
74
75impl SocketError {
76 pub fn code(&self) -> &'static str {
77 match self {
78 Self::Unsupported => "unsupported",
79 Self::Denied => "denied",
80 Self::Invalid(_) => "invalid",
81 }
82 }
83}
84
85pub type SocketHandle = u64;
86pub type SocketCallback = Rc<dyn Fn(SocketEvent)>;
87
88#[derive(Debug, Clone, PartialEq)]
89pub enum SocketEvent {
90 Connected(SocketHandle),
91 Data(SocketHandle, Vec<u8>),
92 Closed(SocketHandle),
93 Failed(SocketHandle, String),
94}
95
96pub type SocketServerCallback = Rc<dyn Fn(SocketServerEvent)>;
97
98#[derive(Debug, Clone, PartialEq)]
99pub enum SocketServerEvent {
100 Open {
101 server: SocketHandle,
102 connection: SocketHandle,
103 },
104 Data {
105 server: SocketHandle,
106 connection: SocketHandle,
107 bytes: Vec<u8>,
108 },
109 Closed {
110 server: SocketHandle,
111 connection: SocketHandle,
112 },
113 Failed {
114 server: SocketHandle,
115 connection: SocketHandle,
116 error: String,
117 },
118}
119
120fn socket_server_event_value(event: SocketServerEvent) -> Value {
121 let mut entries = Vec::new();
122 match event {
123 SocketServerEvent::Open { server, connection } => {
124 entries.push((Value::Keyword("type".into()), Value::Keyword("open".into())));
125 entries.push((
126 Value::Keyword("server".into()),
127 Value::Number(server as i64),
128 ));
129 entries.push((
130 Value::Keyword("connection".into()),
131 Value::Number(connection as i64),
132 ));
133 }
134 SocketServerEvent::Data {
135 server,
136 connection,
137 bytes,
138 } => {
139 entries.push((Value::Keyword("type".into()), Value::Keyword("data".into())));
140 entries.push((
141 Value::Keyword("server".into()),
142 Value::Number(server as i64),
143 ));
144 entries.push((
145 Value::Keyword("connection".into()),
146 Value::Number(connection as i64),
147 ));
148 entries.push((Value::Keyword("bytes".into()), Value::Bytes(bytes)));
149 }
150 SocketServerEvent::Closed { server, connection } => {
151 entries.push((
152 Value::Keyword("type".into()),
153 Value::Keyword("close".into()),
154 ));
155 entries.push((
156 Value::Keyword("server".into()),
157 Value::Number(server as i64),
158 ));
159 entries.push((
160 Value::Keyword("connection".into()),
161 Value::Number(connection as i64),
162 ));
163 }
164 SocketServerEvent::Failed {
165 server,
166 connection,
167 error,
168 } => {
169 entries.push((
170 Value::Keyword("type".into()),
171 Value::Keyword("error".into()),
172 ));
173 entries.push((
174 Value::Keyword("server".into()),
175 Value::Number(server as i64),
176 ));
177 entries.push((
178 Value::Keyword("connection".into()),
179 Value::Number(connection as i64),
180 ));
181 entries.push((Value::Keyword("error".into()), Value::String(error)));
182 }
183 }
184 Value::Map(PMap::from_iter(entries))
185}
186
187pub trait SocketProvider {
188 fn connect(
189 &self,
190 host: &str,
191 port: u16,
192 callback: SocketCallback,
193 ) -> Result<SocketHandle, SocketError>;
194 fn send(&self, socket: SocketHandle, bytes: &[u8]) -> Result<usize, SocketError>;
195 fn close(&self, socket: SocketHandle) -> Result<(), SocketError>;
196 fn listen(
197 &self,
198 _host: &str,
199 _port: u16,
200 _callback: SocketServerCallback,
201 ) -> Result<SocketHandle, SocketError> {
202 Err(SocketError::Unsupported)
203 }
204 fn endpoint(&self, _server: SocketHandle) -> Result<(String, u16), SocketError> {
205 Err(SocketError::Unsupported)
206 }
207 fn events(&self, _handle: SocketHandle) -> Result<SocketHandle, SocketError> {
208 Err(SocketError::Unsupported)
209 }
210 fn next(&self, _stream: SocketHandle) -> Result<Promise, SocketError> {
211 Err(SocketError::Unsupported)
212 }
213}
214
215#[cfg(not(target_arch = "wasm32"))]
216use std::io::{Read, Write};
217#[cfg(not(target_arch = "wasm32"))]
218use std::net::{Shutdown, TcpListener, TcpStream};
219#[cfg(not(target_arch = "wasm32"))]
220use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
221#[cfg(not(target_arch = "wasm32"))]
222use std::sync::{mpsc, Arc, Mutex};
223
224#[cfg(not(target_arch = "wasm32"))]
225#[derive(Clone)]
226enum RawSocketEvent {
227 Open {
228 server: SocketHandle,
229 connection: SocketHandle,
230 stream: Arc<Mutex<TcpStream>>,
231 },
232 Data {
233 server: SocketHandle,
234 connection: SocketHandle,
235 bytes: Vec<u8>,
236 },
237 Closed {
238 server: SocketHandle,
239 connection: SocketHandle,
240 },
241 Failed {
242 server: SocketHandle,
243 connection: SocketHandle,
244 error: String,
245 },
246}
247
248#[cfg(not(target_arch = "wasm32"))]
249struct NativeServer {
250 host: String,
251 port: u16,
252 alive: Arc<AtomicBool>,
253}
254
255#[cfg(not(target_arch = "wasm32"))]
256struct NativeSocketStream {
257 handle: SocketHandle,
258 queue: VecDeque<Value>,
259 queued_bytes: usize,
260 pending: Option<Promise>,
261 closed: bool,
262}
263
264#[cfg(not(target_arch = "wasm32"))]
265struct NativeSocketState {
266 next_handle: Arc<AtomicU64>,
267 sockets: HashMap<SocketHandle, TcpStream>,
268 callbacks: HashMap<SocketHandle, SocketCallback>,
269 servers: HashMap<SocketHandle, NativeServer>,
270 connections: HashMap<SocketHandle, Arc<Mutex<TcpStream>>>,
271 connection_servers: HashMap<SocketHandle, SocketHandle>,
272 server_callbacks: HashMap<SocketHandle, SocketServerCallback>,
273 streams: HashMap<SocketHandle, NativeSocketStream>,
274 sender: mpsc::Sender<RawSocketEvent>,
275 receiver: mpsc::Receiver<RawSocketEvent>,
276}
277
278#[cfg(not(target_arch = "wasm32"))]
279#[derive(Clone)]
280pub struct NativeSocketProvider {
281 state: Rc<RefCell<NativeSocketState>>,
282}
283
284#[cfg(not(target_arch = "wasm32"))]
285impl Default for NativeSocketProvider {
286 fn default() -> Self {
287 let (sender, receiver) = mpsc::channel();
288 Self {
289 state: Rc::new(RefCell::new(NativeSocketState {
290 next_handle: Arc::new(AtomicU64::new(1)),
291 sockets: HashMap::new(),
292 callbacks: HashMap::new(),
293 servers: HashMap::new(),
294 connections: HashMap::new(),
295 connection_servers: HashMap::new(),
296 server_callbacks: HashMap::new(),
297 streams: HashMap::new(),
298 sender,
299 receiver,
300 })),
301 }
302 }
303}
304
305#[cfg(not(target_arch = "wasm32"))]
306impl NativeSocketProvider {
307 fn next_handle(&self) -> SocketHandle {
308 self.state
309 .borrow()
310 .next_handle
311 .fetch_add(1, Ordering::Relaxed)
312 }
313
314 fn pump(&self) {
315 loop {
316 let event = { self.state.borrow().receiver.try_recv().ok() };
317 let Some(event) = event else {
318 break;
319 };
320 self.dispatch(event);
321 }
322 }
323
324 fn wait_and_pump(&self) {
325 let event = { self.state.borrow().receiver.recv().ok() };
326 if let Some(event) = event {
327 self.dispatch(event);
328 }
329 self.pump();
330 }
331
332 fn dispatch(&self, raw: RawSocketEvent) {
333 let event = match raw {
334 RawSocketEvent::Open {
335 server,
336 connection,
337 stream,
338 } => {
339 let mut state = self.state.borrow_mut();
340 state.connections.insert(connection, stream);
341 state.connection_servers.insert(connection, server);
342 SocketServerEvent::Open { server, connection }
343 }
344 RawSocketEvent::Data {
345 server,
346 connection,
347 bytes,
348 } => SocketServerEvent::Data {
349 server,
350 connection,
351 bytes,
352 },
353 RawSocketEvent::Closed { server, connection } => {
354 self.state.borrow_mut().connections.remove(&connection);
355 SocketServerEvent::Closed { server, connection }
356 }
357 RawSocketEvent::Failed {
358 server,
359 connection,
360 error,
361 } => SocketServerEvent::Failed {
362 server,
363 connection,
364 error,
365 },
366 };
367 let callback = {
368 self.state
369 .borrow()
370 .server_callbacks
371 .get(&match &event {
372 SocketServerEvent::Open { server, .. }
373 | SocketServerEvent::Data { server, .. }
374 | SocketServerEvent::Closed { server, .. }
375 | SocketServerEvent::Failed { server, .. } => *server,
376 })
377 .cloned()
378 };
379 if let Some(callback) = callback {
380 callback(event.clone());
381 }
382 let (server, connection, bytes) = match &event {
383 SocketServerEvent::Open { server, connection }
384 | SocketServerEvent::Closed { server, connection }
385 | SocketServerEvent::Failed {
386 server, connection, ..
387 } => (*server, *connection, 0),
388 SocketServerEvent::Data {
389 server,
390 connection,
391 bytes,
392 } => (*server, *connection, bytes.len()),
393 };
394 let value = socket_server_event_value(event);
395 let overflow = {
396 let mut state = self.state.borrow_mut();
397 let mut overflow = false;
398 for stream in state
399 .streams
400 .values_mut()
401 .filter(|stream| stream.handle == server || stream.handle == connection)
402 {
403 if stream.closed {
404 continue;
405 }
406 if stream.queue.len() >= 256
407 || stream.queued_bytes.saturating_add(bytes) > 1_048_576
408 {
409 stream.closed = true;
410 if let Some(promise) = stream.pending.take() {
411 promise.resolve(Value::Map(PMap::from_iter([
412 (
413 Value::Keyword("type".into()),
414 Value::Keyword("error".into()),
415 ),
416 (
417 Value::Keyword("error".into()),
418 Value::String("buffer-overflow".into()),
419 ),
420 ])));
421 }
422 overflow = true;
423 continue;
424 }
425 if let Some(promise) = stream.pending.take() {
426 promise.resolve(value.clone());
427 } else {
428 stream.queued_bytes += bytes;
429 stream.queue.push_back(value.clone());
430 }
431 }
432 overflow
433 };
434 if overflow {
435 let _ = self.close(connection);
436 }
437 }
438}
439
440#[cfg(not(target_arch = "wasm32"))]
441impl SocketProvider for NativeSocketProvider {
442 fn connect(
443 &self,
444 host: &str,
445 port: u16,
446 callback: SocketCallback,
447 ) -> Result<SocketHandle, SocketError> {
448 if host.is_empty() || port == 0 {
449 return Err(SocketError::Invalid("host and port are required".into()));
450 }
451 let stream = TcpStream::connect((host, port))
452 .map_err(|error| SocketError::Invalid(error.to_string()))?;
453 let handle = self.next_handle();
454 self.state.borrow_mut().sockets.insert(handle, stream);
455 self.state
456 .borrow_mut()
457 .callbacks
458 .insert(handle, callback.clone());
459 callback(SocketEvent::Connected(handle));
460 Ok(handle)
461 }
462
463 fn send(&self, socket: SocketHandle, bytes: &[u8]) -> Result<usize, SocketError> {
464 let mut state = self.state.borrow_mut();
465 if let Some(stream) = state.sockets.get_mut(&socket) {
466 stream
467 .write_all(bytes)
468 .map_err(|error| SocketError::Invalid(error.to_string()))?;
469 drop(state);
470 if let Some(callback) = self.state.borrow().callbacks.get(&socket).cloned() {
471 callback(SocketEvent::Data(socket, bytes.to_vec()));
472 }
473 return Ok(bytes.len());
474 }
475 let accepted = state.connections.get(&socket).cloned();
476 drop(state);
477 let accepted = accepted.ok_or_else(|| SocketError::Invalid("unknown socket".into()))?;
478 accepted
479 .lock()
480 .map_err(|_| SocketError::Invalid("socket lock poisoned".into()))?
481 .write_all(bytes)
482 .map_err(|error| SocketError::Invalid(error.to_string()))?;
483 Ok(bytes.len())
484 }
485
486 fn close(&self, socket: SocketHandle) -> Result<(), SocketError> {
487 if self.state.borrow_mut().sockets.remove(&socket).is_some() {
488 if let Some(callback) = self.state.borrow_mut().callbacks.remove(&socket) {
489 callback(SocketEvent::Closed(socket));
490 }
491 return Ok(());
492 }
493 let server = { self.state.borrow_mut().servers.remove(&socket) };
494 if let Some(server) = server {
495 server.alive.store(false, Ordering::Relaxed);
496 self.state.borrow_mut().server_callbacks.remove(&socket);
497 return Ok(());
498 }
499 let (stream, server, sender) = {
500 let mut state = self.state.borrow_mut();
501 (
502 state.connections.remove(&socket),
503 state.connection_servers.remove(&socket).unwrap_or(0),
504 state.sender.clone(),
505 )
506 };
507 if let Some(stream) = stream {
508 let _ = stream.lock().map(|stream| stream.shutdown(Shutdown::Both));
509 let _ = sender.send(RawSocketEvent::Closed {
510 server,
511 connection: socket,
512 });
513 self.pump();
514 return Ok(());
515 }
516 Err(SocketError::Invalid("unknown socket".into()))
517 }
518
519 fn listen(
520 &self,
521 host: &str,
522 port: u16,
523 callback: SocketServerCallback,
524 ) -> Result<SocketHandle, SocketError> {
525 if host.is_empty() {
526 return Err(SocketError::Invalid("host is required".into()));
527 }
528 let listener = TcpListener::bind((host, port))
529 .map_err(|error| SocketError::Invalid(error.to_string()))?;
530 let endpoint = listener
531 .local_addr()
532 .map_err(|error| SocketError::Invalid(error.to_string()))?;
533 listener
534 .set_nonblocking(true)
535 .map_err(|error| SocketError::Invalid(error.to_string()))?;
536 let server = self.next_handle();
537 let alive = Arc::new(AtomicBool::new(true));
538 let sender = self.state.borrow().sender.clone();
539 let next_handle = self.state.borrow().next_handle.clone();
540 let thread_alive = alive.clone();
541 std::thread::Builder::new()
542 .name(format!("hara-socket-{server}"))
543 .spawn(move || {
544 while thread_alive.load(Ordering::Relaxed) {
545 match listener.accept() {
546 Ok((stream, _)) => {
547 let connection = next_handle.fetch_add(1, Ordering::Relaxed);
548 if let Err(error) = stream.set_nonblocking(false) {
549 let _ = sender.send(RawSocketEvent::Failed {
550 server,
551 connection,
552 error: error.to_string(),
553 });
554 continue;
555 }
556 let mut reader = match stream.try_clone() {
557 Ok(reader) => reader,
558 Err(error) => {
559 let _ = sender.send(RawSocketEvent::Failed {
560 server,
561 connection,
562 error: error.to_string(),
563 });
564 continue;
565 }
566 };
567 let shared = Arc::new(Mutex::new(stream));
568 let _ = sender.send(RawSocketEvent::Open {
569 server,
570 connection,
571 stream: shared.clone(),
572 });
573 let reader_sender = sender.clone();
574 std::thread::spawn(move || {
575 let mut buffer = [0u8; 8192];
576 loop {
577 let read = reader.read(&mut buffer);
578 match read {
579 Ok(0) => {
580 let _ = reader_sender.send(RawSocketEvent::Closed {
581 server,
582 connection,
583 });
584 break;
585 }
586 Ok(count) => {
587 let _ = reader_sender.send(RawSocketEvent::Data {
588 server,
589 connection,
590 bytes: buffer[..count].to_vec(),
591 });
592 }
593 Err(error) => {
594 let _ = reader_sender.send(RawSocketEvent::Failed {
595 server,
596 connection,
597 error: error.to_string(),
598 });
599 break;
600 }
601 }
602 }
603 });
604 }
605 Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {
606 std::thread::sleep(std::time::Duration::from_millis(5))
607 }
608 Err(error) => {
609 let _ = sender.send(RawSocketEvent::Failed {
610 server,
611 connection: 0,
612 error: error.to_string(),
613 });
614 break;
615 }
616 }
617 }
618 })
619 .map_err(|error| SocketError::Invalid(error.to_string()))?;
620 let mut state = self.state.borrow_mut();
621 state.servers.insert(
622 server,
623 NativeServer {
624 host: endpoint.ip().to_string(),
625 port: endpoint.port(),
626 alive,
627 },
628 );
629 state.server_callbacks.insert(server, callback);
630 Ok(server)
631 }
632
633 fn endpoint(&self, server: SocketHandle) -> Result<(String, u16), SocketError> {
634 let state = self.state.borrow();
635 let server = state
636 .servers
637 .get(&server)
638 .ok_or_else(|| SocketError::Invalid("unknown socket server".into()))?;
639 Ok((server.host.clone(), server.port))
640 }
641
642 fn events(&self, handle: SocketHandle) -> Result<SocketHandle, SocketError> {
643 let mut state = self.state.borrow_mut();
644 if !state.servers.contains_key(&handle) && !state.connections.contains_key(&handle) {
645 return Err(SocketError::Invalid("unknown socket handle".into()));
646 }
647 let stream = state.next_handle.fetch_add(1, Ordering::Relaxed);
648 state.streams.insert(
649 stream,
650 NativeSocketStream {
651 handle,
652 queue: VecDeque::new(),
653 queued_bytes: 0,
654 pending: None,
655 closed: false,
656 },
657 );
658 Ok(stream)
659 }
660
661 fn next(&self, stream: SocketHandle) -> Result<Promise, SocketError> {
662 self.pump();
663 let promise = Promise::new();
664 {
665 let mut state = self.state.borrow_mut();
666 let stream = state
667 .streams
668 .get_mut(&stream)
669 .ok_or_else(|| SocketError::Invalid("unknown socket stream".into()))?;
670 if let Some(event) = stream.queue.pop_front() {
671 stream.queued_bytes = 0;
672 promise.resolve(event);
673 return Ok(promise);
674 }
675 if stream.closed {
676 promise.resolve(Value::Map(PMap::from_iter([(
677 Value::Keyword("type".into()),
678 Value::Keyword("close".into()),
679 )])));
680 return Ok(promise);
681 }
682 if stream.pending.is_some() {
683 return Err(SocketError::Invalid(
684 "socket stream already has a pending next".into(),
685 ));
686 }
687 stream.pending = Some(promise.clone());
688 }
689 let provider = self.clone();
690 promise.set_poller(Rc::new(move || provider.pump()));
691 let provider = self.clone();
692 promise.set_waiter(Rc::new(move || provider.wait_and_pump()));
693 Ok(promise)
694 }
695}
696
697#[derive(Debug, Clone, Copy, PartialEq, Eq)]
698pub struct ProviderCapabilities {
699 pub file: bool,
700 pub socket: bool,
701 pub process: bool,
702}
703
704pub struct ProviderRegistry {
705 file: Option<Rc<dyn FileProvider>>,
706 socket: Option<Rc<dyn SocketProvider>>,
707 kernel: Option<Rc<KernelProvider>>,
708 promise: Rc<dyn PromiseProvider>,
709 process: bool,
710}
711
712impl Default for ProviderRegistry {
713 fn default() -> Self {
714 Self {
715 file: None,
716 socket: None,
717 kernel: None,
718 promise: Rc::new(LocalPromiseProvider),
719 process: false,
720 }
721 }
722}
723
724impl ProviderRegistry {
725 pub fn new() -> Self {
726 Self::default()
727 }
728
729 pub fn install_file<P: FileProvider + 'static>(&mut self, provider: P) {
730 self.file = Some(Rc::new(provider));
731 }
732 pub fn set_file(&mut self, provider: Option<Rc<dyn FileProvider>>) {
733 self.file = provider;
734 }
735 pub fn install_socket<P: SocketProvider + 'static>(&mut self, provider: P) {
736 self.socket = Some(Rc::new(provider));
737 }
738 pub fn install_kernel(&mut self, provider: Rc<KernelProvider>) {
739 self.kernel = Some(provider);
740 }
741 pub fn install_process(&mut self) {
742 self.process = true;
743 }
744 pub fn install_promise<P: PromiseProvider + 'static>(&mut self, provider: P) {
745 self.promise = Rc::new(provider);
746 }
747 pub fn promise(&self) -> Rc<dyn PromiseProvider> {
748 self.promise.clone()
749 }
750 pub fn file(&self) -> Option<Rc<dyn FileProvider>> {
751 self.file.clone()
752 }
753 pub fn socket(&self) -> Option<Rc<dyn SocketProvider>> {
754 self.socket.clone()
755 }
756 pub fn kernel(&self) -> Option<Rc<KernelProvider>> {
757 self.kernel.clone()
758 }
759 pub fn process(&self) -> bool {
760 self.process
761 }
762 pub fn capabilities(&self) -> ProviderCapabilities {
763 ProviderCapabilities {
764 file: self.file.is_some(),
765 socket: self.socket.is_some(),
766 process: self.process,
767 }
768 }
769}
770
771#[derive(Clone)]
772pub struct LoopbackSocketProvider {
773 next_handle: Rc<Cell<SocketHandle>>,
774 callbacks: Rc<RefCell<HashMap<SocketHandle, SocketCallback>>>,
775 streams: Rc<RefCell<HashMap<SocketHandle, LoopbackSocketStream>>>,
776}
777
778struct LoopbackSocketStream { socket: SocketHandle, queue: VecDeque<Value>, pending: Option<Promise>, closed: bool }
779
780impl Default for LoopbackSocketProvider {
781 fn default() -> Self {
782 Self {
783 next_handle: Rc::new(Cell::new(1)),
784 callbacks: Rc::new(RefCell::new(HashMap::new())),
785 streams: Rc::new(RefCell::new(HashMap::new())),
786 }
787 }
788}
789
790impl SocketProvider for LoopbackSocketProvider {
791 fn connect(
792 &self,
793 host: &str,
794 port: u16,
795 callback: SocketCallback,
796 ) -> Result<SocketHandle, SocketError> {
797 if host.is_empty() || port == 0 {
798 return Err(SocketError::Invalid("host and port are required".into()));
799 }
800 let handle = self.next_handle.get();
801 self.next_handle.set(handle + 1);
802 self.callbacks.borrow_mut().insert(handle, callback.clone());
803 callback(SocketEvent::Connected(handle));
804 Ok(handle)
805 }
806
807 fn send(&self, socket: SocketHandle, bytes: &[u8]) -> Result<usize, SocketError> {
808 let callback = self
809 .callbacks
810 .borrow()
811 .get(&socket)
812 .cloned()
813 .ok_or_else(|| SocketError::Invalid("unknown socket".into()))?;
814 callback(SocketEvent::Data(socket, bytes.to_vec()));
815 let event = socket_server_event_value(SocketServerEvent::Data { server: 0, connection: socket, bytes: bytes.to_vec() });
816 for stream in self.streams.borrow_mut().values_mut().filter(|s| s.socket == socket) {
817 if let Some(promise) = stream.pending.take() { promise.resolve(event.clone()); } else { stream.queue.push_back(event.clone()); }
818 }
819 Ok(bytes.len())
820 }
821
822 fn close(&self, socket: SocketHandle) -> Result<(), SocketError> {
823 let callback = self
824 .callbacks
825 .borrow_mut()
826 .remove(&socket)
827 .ok_or_else(|| SocketError::Invalid("unknown socket".into()))?;
828 callback(SocketEvent::Closed(socket));
829 let event = socket_server_event_value(SocketServerEvent::Closed { server: 0, connection: socket });
830 for stream in self.streams.borrow_mut().values_mut().filter(|s| s.socket == socket) {
831 stream.closed = true;
832 if let Some(promise) = stream.pending.take() { promise.resolve(event.clone()); } else { stream.queue.push_back(event.clone()); }
833 }
834 Ok(())
835 }
836
837 fn events(&self, socket: SocketHandle) -> Result<SocketHandle, SocketError> {
838 if !self.callbacks.borrow().contains_key(&socket) { return Err(SocketError::Invalid("unknown socket".into())); }
839 let handle = self.next_handle.get(); self.next_handle.set(handle + 1);
840 self.streams.borrow_mut().insert(handle, LoopbackSocketStream { socket, queue: VecDeque::new(), pending: None, closed: false });
841 Ok(handle)
842 }
843
844 fn next(&self, handle: SocketHandle) -> Result<Promise, SocketError> {
845 let promise = Promise::new();
846 let mut streams = self.streams.borrow_mut();
847 let stream = streams.get_mut(&handle).ok_or_else(|| SocketError::Invalid("unknown socket stream".into()))?;
848 if let Some(event) = stream.queue.pop_front() { promise.resolve(event); }
849 else if stream.closed { promise.resolve(socket_server_event_value(SocketServerEvent::Closed { server: 0, connection: stream.socket })); }
850 else if stream.pending.is_some() { return Err(SocketError::Invalid("socket stream already has a pending next".into())); }
851 else { stream.pending = Some(promise.clone()); }
852 Ok(promise)
853 }
854}
855
856#[derive(Debug, Default, Clone, Copy)]
857pub struct UnsupportedSocketProvider;
858
859impl SocketProvider for UnsupportedSocketProvider {
860 fn connect(
861 &self,
862 _host: &str,
863 _port: u16,
864 _callback: SocketCallback,
865 ) -> Result<SocketHandle, SocketError> {
866 Err(SocketError::Unsupported)
867 }
868 fn send(&self, _socket: SocketHandle, _bytes: &[u8]) -> Result<usize, SocketError> {
869 Err(SocketError::Unsupported)
870 }
871 fn close(&self, _socket: SocketHandle) -> Result<(), SocketError> {
872 Err(SocketError::Unsupported)
873 }
874}
875
876pub fn portable_type_name(value: &Value) -> &str {
877 match value {
878 Value::Nil => "nil",
879 Value::Number(_) => "long",
880 Value::Float(_) => "float",
881 Value::BigInteger(_) if crate::numeric::is_long_value(value) => "long",
882 Value::BigInteger(_) => "bigint",
883 Value::Character(_) => "character",
884 Value::Regex(_) => "pattern",
885 Value::Tagged(_) => "tagged-literal",
886 Value::Bool(_) => "boolean",
887 Value::String(_) => "string",
888 Value::Keyword(_) => "keyword",
889 Value::Symbol(_) => "symbol",
890 Value::Pointer(_) => "pointer",
891 Value::Function(_) => "function",
892 Value::Bytes(_) => "bytes",
893 Value::ByteBuffer(_) => "byte-buffer",
894 Value::Array(_) => "array",
895 Value::Object(_) => "object",
896 Value::Promise(_) => "promise",
897 Value::Atom(_) => "atom",
898 Value::Recur(_) => "recur",
899 Value::List(_) => "list",
900 Value::Cons(_) => "cons",
901 Value::Queue(_) => "queue",
902 Value::Deque(_) => "deque",
903 Value::Tuple(_) => "vector",
904 Value::Vector(_) => "vector",
905 Value::MapEntry(_) => "map-entry",
906 Value::MutableCollection(_) => "mutable-collection",
907 Value::Seq(_) => "seq",
908 Value::Map(_) => "hash-map",
909 Value::OrderedMap(_) => "ordered-map",
910 Value::SortedMap(_) => "sorted-map",
911 Value::Trie(_) => "trie",
912 Value::PriorityMap(_) => "priority-map",
913 Value::Set(_) => "hash-set",
914 Value::OrderedSet(_) => "ordered-set",
915 Value::SortedSet(_) => "sorted-set",
916 Value::Iterator(_) => "iterator",
917 Value::Var(_) => "var",
918 Value::Namespace(_) => "namespace",
919 Value::Extension(_) => "extension",
920 Value::StructType(_) => "struct-type",
921 Value::Struct(_) => "struct",
922 Value::MutableType(_) => "mutable-type",
923 Value::Mutable(_) => "mutable",
924 Value::Protocol(_) => "protocol",
925 Value::NativeType(_) => "native-type",
926 Value::Schema(_) => "schema",
927 Value::Coroutine(_) => "coroutine",
928 Value::Stream(_) => "stream",
929 Value::Result(_) => "result",
930 Value::ExceptionInfo(_) => "exception",
931 }
932}