Skip to main content

hara_native/core/
provider.rs

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}