1use crate::host::{invoke, with_host, IoTask, JsObj};
21use fusevm::Value;
22use indexmap::IndexMap;
23use once_cell::sync::OnceCell;
24use rustls::pki_types::{CertificateDer, PrivateKeyDer, ServerName, UnixTime};
25use rustls::{
26 ClientConfig, ClientConnection, ConnectionCommon, DigitallySignedStruct, RootCertStore,
27 ServerConfig, ServerConnection, SideData, SignatureScheme, StreamOwned,
28};
29use std::collections::HashMap;
30use std::io::{Read, Write};
31use std::net::TcpStream;
32use std::ops::{Deref, DerefMut};
33use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
34use std::sync::mpsc::{Receiver, Sender};
35use std::sync::Arc;
36
37pub const MODULE_METHODS: &[&str] = &[
39 "connect",
40 "createServer",
41 "createSecureContext",
42 "checkServerIdentity",
43 "convertALPNProtocols",
44 "getCiphers",
45 "getCACertificates",
46 "setDefaultCACertificates",
47 "getCertificateCompressionAlgorithms",
48];
49
50const CIPHERS: &[&str] = &[
54 "aes128-gcm-sha256",
55 "aes128-sha",
56 "aes128-sha256",
57 "aes256-gcm-sha384",
58 "aes256-sha",
59 "aes256-sha256",
60 "dhe-psk-aes128-cbc-sha",
61 "dhe-psk-aes128-cbc-sha256",
62 "dhe-psk-aes128-gcm-sha256",
63 "dhe-psk-aes256-cbc-sha",
64 "dhe-psk-aes256-cbc-sha384",
65 "dhe-psk-aes256-gcm-sha384",
66 "dhe-psk-chacha20-poly1305",
67 "dhe-rsa-aes128-gcm-sha256",
68 "dhe-rsa-aes128-sha",
69 "dhe-rsa-aes128-sha256",
70 "dhe-rsa-aes256-gcm-sha384",
71 "dhe-rsa-aes256-sha",
72 "dhe-rsa-aes256-sha256",
73 "dhe-rsa-chacha20-poly1305",
74 "ecdhe-ecdsa-aes128-gcm-sha256",
75 "ecdhe-ecdsa-aes128-sha",
76 "ecdhe-ecdsa-aes128-sha256",
77 "ecdhe-ecdsa-aes256-gcm-sha384",
78 "ecdhe-ecdsa-aes256-sha",
79 "ecdhe-ecdsa-aes256-sha384",
80 "ecdhe-ecdsa-chacha20-poly1305",
81 "ecdhe-psk-aes128-cbc-sha",
82 "ecdhe-psk-aes128-cbc-sha256",
83 "ecdhe-psk-aes256-cbc-sha",
84 "ecdhe-psk-aes256-cbc-sha384",
85 "ecdhe-psk-chacha20-poly1305",
86 "ecdhe-rsa-aes128-gcm-sha256",
87 "ecdhe-rsa-aes128-sha",
88 "ecdhe-rsa-aes128-sha256",
89 "ecdhe-rsa-aes256-gcm-sha384",
90 "ecdhe-rsa-aes256-sha",
91 "ecdhe-rsa-aes256-sha384",
92 "ecdhe-rsa-chacha20-poly1305",
93 "psk-aes128-cbc-sha",
94 "psk-aes128-cbc-sha256",
95 "psk-aes128-gcm-sha256",
96 "psk-aes256-cbc-sha",
97 "psk-aes256-cbc-sha384",
98 "psk-aes256-gcm-sha384",
99 "psk-chacha20-poly1305",
100 "rsa-psk-aes128-cbc-sha",
101 "rsa-psk-aes128-cbc-sha256",
102 "rsa-psk-aes128-gcm-sha256",
103 "rsa-psk-aes256-cbc-sha",
104 "rsa-psk-aes256-cbc-sha384",
105 "rsa-psk-aes256-gcm-sha384",
106 "rsa-psk-chacha20-poly1305",
107 "srp-aes-128-cbc-sha",
108 "srp-aes-256-cbc-sha",
109 "srp-rsa-aes-128-cbc-sha",
110 "srp-rsa-aes-256-cbc-sha",
111 "tls_aes_128_ccm_8_sha256",
112 "tls_aes_128_ccm_sha256",
113 "tls_aes_128_gcm_sha256",
114 "tls_aes_256_gcm_sha384",
115 "tls_chacha20_poly1305_sha256",
116];
117
118pub const SERVER_METHODS: &[&str] = &["listen", "close", "address"];
121pub const SOCKET_METHODS: &[&str] = &[
122 "write",
123 "end",
124 "destroy",
125 "setEncoding",
126 "setKeepAlive",
127 "setNoDelay",
128 "setTimeout",
129 "ref",
130 "unref",
131 "pause",
132 "resume",
133];
134
135pub type ConnHook = std::rc::Rc<dyn Fn(&Value, &Value, u64) -> Result<(), String>>;
139
140enum WriteCmd {
142 Data(Vec<u8>),
143 Shutdown,
144}
145
146static NEXT_TLS_ID: AtomicU64 = AtomicU64::new(1);
150static NEXT_SERVER_ID: AtomicU64 = AtomicU64::new(1);
151
152fn next_tls_id() -> u64 {
153 NEXT_TLS_ID.fetch_add(1, Ordering::Relaxed)
154}
155fn next_server_id() -> u64 {
156 NEXT_SERVER_ID.fetch_add(1, Ordering::Relaxed)
157}
158
159struct TlsServerRec {
161 emitter: Value,
162 stop: Arc<AtomicBool>,
163 conn_hook: Option<ConnHook>,
164 listener: Option<Value>,
166}
167
168struct TlsSocketRec {
170 emitter: Value,
171 tx: Sender<WriteCmd>,
173}
174
175#[derive(Default)]
176struct TlsState {
177 servers: HashMap<u64, TlsServerRec>,
178 sockets: HashMap<u64, TlsSocketRec>,
179}
180
181thread_local! {
182 static TLS: std::cell::RefCell<TlsState> = std::cell::RefCell::new(TlsState::default());
183 static PENDING_CONFIGS: std::cell::RefCell<Vec<(Value, Arc<ServerConfig>)>> =
186 const { std::cell::RefCell::new(Vec::new()) };
187 static PENDING_HOOKS: std::cell::RefCell<Vec<(Value, ConnHook)>> =
188 const { std::cell::RefCell::new(Vec::new()) };
189 static DEFAULT_CA_CERTS: std::cell::RefCell<Vec<String>> =
192 const { std::cell::RefCell::new(Vec::new()) };
193}
194
195fn get_prop(recv: &Value, key: &str) -> Option<Value> {
198 with_host(|h| match h.get(recv) {
199 Some(JsObj::Object(p)) => p.get(key).cloned(),
200 _ => None,
201 })
202}
203
204fn set_prop(recv: &Value, key: &str, val: Value) {
205 with_host(|h| {
206 if let Some(JsObj::Object(p)) = h.get_mut(recv) {
207 p.insert(key.to_string(), val);
208 }
209 });
210}
211
212fn u64_prop(recv: &Value, key: &str) -> Option<u64> {
213 get_prop(recv, key).map(|v| with_host(|h| h.to_number(&v)) as u64)
214}
215
216fn emitter_dispatch(recv: &Value, method: &str, args: &[Value]) -> Option<Result<Value, String>> {
218 super::events::METHODS
219 .contains(&method)
220 .then(|| super::events::instance_call(recv, method, args.to_vec()))
221}
222
223fn value_bytes(v: Option<&Value>) -> Vec<u8> {
226 let Some(v) = v else { return Vec::new() };
227 let is_buffer =
228 with_host(|h| matches!(h.get(v), Some(JsObj::Object(p)) if p.contains_key("@@bytes")));
229 if is_buffer {
230 return with_host(|h| match h.get(v) {
231 Some(JsObj::Object(p)) => match p.get("@@bytes").and_then(|b| h.get(b)) {
232 Some(JsObj::Array(items)) => items.iter().map(|x| h.to_number(x) as u8).collect(),
233 _ => Vec::new(),
234 },
235 _ => Vec::new(),
236 });
237 }
238 with_host(|h| h.str_of(v)).into_bytes()
239}
240
241pub fn call(method: &str, args: &[Value]) -> Option<Result<Value, String>> {
245 match method {
246 "connect" => Some(connect(args)),
247 "createServer" => Some(create_server(args)),
248 "createSecureContext" => Some(Ok(with_host(|h| {
251 let mut m = IndexMap::new();
252 m.insert("@@native".into(), h.new_str("SecureContext"));
253 if let Some(o) = args.first() {
254 m.insert("@@options".into(), o.clone());
255 }
256 h.new_object(m)
257 }))),
258 "checkServerIdentity" => Some(check_server_identity(args)),
259 "convertALPNProtocols" => Some(convert_alpn_protocols(args)),
260 "getCiphers" => Some(Ok(with_host(|h| {
261 let items = CIPHERS.iter().map(|c| h.new_str(*c)).collect();
262 h.new_array(items)
263 }))),
264 "getCACertificates" => Some(get_ca_certificates(args)),
265 "setDefaultCACertificates" => Some(set_default_ca_certificates(args)),
266 "getCertificateCompressionAlgorithms" => Some(Ok(with_host(|h| h.new_array(Vec::new())))),
270 _ => None,
271 }
272}
273
274fn read_string_array(v: Option<&Value>) -> Vec<String> {
279 let Some(v) = v else { return Vec::new() };
280 with_host(|h| match h.get(v) {
281 Some(JsObj::Array(items)) => items.iter().map(|x| h.str_of(x)).collect(),
282 _ if h.as_str(v).is_some() => vec![h.str_of(v)],
283 _ => Vec::new(),
284 })
285}
286
287fn host_matches(host: &str, pattern: &str) -> bool {
291 let host = host.trim().to_ascii_lowercase();
292 let pattern = pattern.trim().to_ascii_lowercase();
293 if let Some(suffix) = pattern.strip_prefix("*.") {
294 match host.split_once('.') {
296 Some((_, rest)) => rest == suffix,
297 None => false,
298 }
299 } else {
300 host == pattern
301 }
302}
303
304fn check_server_identity(args: &[Value]) -> Result<Value, String> {
309 let host = super::arg_str(args, 0);
310 let cert = args.get(1).cloned().unwrap_or(Value::Undef);
311
312 let altname_str = get_prop(&cert, "subjectaltname")
314 .filter(|v| with_host(|h| h.as_str(v)).is_some())
315 .map(|v| with_host(|h| h.str_of(&v)))
316 .unwrap_or_default();
317 let dns_names: Vec<String> = altname_str
318 .split(',')
319 .filter_map(|e| e.trim().strip_prefix("DNS:").map(|s| s.trim().to_string()))
320 .filter(|s| !s.is_empty())
321 .collect();
322
323 let matched = if !dns_names.is_empty() {
324 dns_names.iter().any(|p| host_matches(&host, p))
325 } else {
326 let cn = get_prop(&cert, "subject")
328 .and_then(|s| get_prop(&s, "CN"))
329 .filter(|v| with_host(|h| h.as_str(v)).is_some())
330 .map(|v| with_host(|h| h.str_of(&v)))
331 .unwrap_or_default();
332 !cn.is_empty() && host_matches(&host, &cn)
333 };
334
335 if matched {
336 return Ok(Value::Undef);
337 }
338 let reason = if !altname_str.is_empty() {
339 format!("Host: {host}. is not in the cert's altnames: {altname_str}")
340 } else {
341 format!("Host: {host}. is not cert's CN")
342 };
343 let message = format!("Hostname/IP does not match certificate's altnames: {reason}");
344 Ok(with_host(|h| {
345 let mut m = IndexMap::new();
346 m.insert("message".into(), h.new_str(message));
347 m.insert("reason".into(), h.new_str(reason));
348 m.insert("host".into(), h.new_str(host));
349 m.insert("code".into(), h.new_str("ERR_TLS_CERT_ALTNAME_INVALID"));
350 m.insert("cert".into(), cert);
351 h.new_object(m)
352 }))
353}
354
355fn convert_alpn_protocols(args: &[Value]) -> Result<Value, String> {
359 let protocols = read_string_array(args.first());
360 let mut wire = Vec::new();
361 for p in &protocols {
362 let bytes = p.as_bytes();
363 wire.push(bytes.len().min(255) as u8);
366 wire.extend_from_slice(&bytes[..bytes.len().min(255)]);
367 }
368 let buf = super::buffer::from_bytes(&wire);
369 if let Some(out) = args.get(1).filter(|v| matches!(v, Value::Obj(_))) {
370 set_prop(out, "ALPNProtocols", buf.clone());
371 }
372 Ok(buf)
373}
374
375fn get_ca_certificates(args: &[Value]) -> Result<Value, String> {
384 let kind = if args.is_empty() {
385 "default".to_string()
386 } else {
387 super::arg_str(args, 0)
388 };
389 let certs: Vec<String> = match kind.as_str() {
390 "default" => DEFAULT_CA_CERTS.with(|c| c.borrow().clone()),
391 _ => Vec::new(),
392 };
393 Ok(with_host(|h| {
394 let items = certs.into_iter().map(|c| h.new_str(c)).collect();
395 h.new_array(items)
396 }))
397}
398
399fn set_default_ca_certificates(args: &[Value]) -> Result<Value, String> {
402 let certs = read_string_array(args.first());
403 DEFAULT_CA_CERTS.with(|c| *c.borrow_mut() = certs);
404 Ok(Value::Undef)
405}
406
407fn verifying_client_config() -> Arc<ClientConfig> {
414 static CFG: OnceCell<Arc<ClientConfig>> = OnceCell::new();
415 CFG.get_or_init(|| {
416 let root_store = RootCertStore {
417 roots: webpki_roots::TLS_SERVER_ROOTS.to_vec(),
418 };
419 let provider = Arc::new(rustls::crypto::aws_lc_rs::default_provider());
420 let cfg = ClientConfig::builder_with_provider(provider)
421 .with_safe_default_protocol_versions()
422 .expect("aws-lc-rs provider supports the default protocol versions")
423 .with_root_certificates(root_store)
424 .with_no_client_auth();
425 Arc::new(cfg)
426 })
427 .clone()
428}
429
430fn insecure_client_config() -> Arc<ClientConfig> {
434 static CFG: OnceCell<Arc<ClientConfig>> = OnceCell::new();
435 CFG.get_or_init(|| {
436 let provider = Arc::new(rustls::crypto::aws_lc_rs::default_provider());
437 let cfg = ClientConfig::builder_with_provider(provider.clone())
438 .with_safe_default_protocol_versions()
439 .expect("aws-lc-rs provider supports the default protocol versions")
440 .dangerous()
441 .with_custom_certificate_verifier(Arc::new(NoVerify(provider)))
442 .with_no_client_auth();
443 Arc::new(cfg)
444 })
445 .clone()
446}
447
448#[derive(Debug)]
451struct NoVerify(Arc<rustls::crypto::CryptoProvider>);
452
453impl rustls::client::danger::ServerCertVerifier for NoVerify {
454 fn verify_server_cert(
455 &self,
456 _end_entity: &CertificateDer<'_>,
457 _intermediates: &[CertificateDer<'_>],
458 _server_name: &ServerName<'_>,
459 _ocsp: &[u8],
460 _now: UnixTime,
461 ) -> Result<rustls::client::danger::ServerCertVerified, rustls::Error> {
462 Ok(rustls::client::danger::ServerCertVerified::assertion())
463 }
464 fn verify_tls12_signature(
465 &self,
466 message: &[u8],
467 cert: &CertificateDer<'_>,
468 dss: &DigitallySignedStruct,
469 ) -> Result<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
470 rustls::crypto::verify_tls12_signature(
471 message,
472 cert,
473 dss,
474 &self.0.signature_verification_algorithms,
475 )
476 }
477 fn verify_tls13_signature(
478 &self,
479 message: &[u8],
480 cert: &CertificateDer<'_>,
481 dss: &DigitallySignedStruct,
482 ) -> Result<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
483 rustls::crypto::verify_tls13_signature(
484 message,
485 cert,
486 dss,
487 &self.0.signature_verification_algorithms,
488 )
489 }
490 fn supported_verify_schemes(&self) -> Vec<SignatureScheme> {
491 self.0.signature_verification_algorithms.supported_schemes()
492 }
493}
494
495pub fn client_config(reject_unauthorized: bool) -> Arc<ClientConfig> {
499 if reject_unauthorized {
500 verifying_client_config()
501 } else {
502 insecure_client_config()
503 }
504}
505
506pub fn build_server_config(cert_pem: &[u8], key_pem: &[u8]) -> Result<Arc<ServerConfig>, String> {
508 let certs: Vec<CertificateDer<'static>> = rustls_pemfile::certs(&mut &cert_pem[..])
509 .collect::<Result<_, _>>()
510 .map_err(|e| format!("Error: tls: bad certificate PEM: {e}"))?;
511 if certs.is_empty() {
512 return Err("Error: tls: no certificates found in `cert`".to_string());
513 }
514 let key: PrivateKeyDer<'static> = rustls_pemfile::private_key(&mut &key_pem[..])
515 .map_err(|e| format!("Error: tls: bad private key PEM: {e}"))?
516 .ok_or_else(|| "Error: tls: no private key found in `key`".to_string())?;
517 let cfg = ServerConfig::builder()
518 .with_no_client_auth()
519 .with_single_cert(certs, key)
520 .map_err(|e| format!("Error: tls: invalid key/cert: {e}"))?;
521 Ok(Arc::new(cfg))
522}
523
524pub fn connect(args: &[Value]) -> Result<Value, String> {
530 let mut port: u16 = 0;
531 let mut host = "localhost".to_string();
532 let mut servername: Option<String> = None;
533 let mut reject_unauthorized = true;
534 let mut cb: Option<Value> = None;
535
536 for a in args {
537 let n = with_host(|h| h.to_number(a));
538 if with_host(|h| crate::host::is_callable(h, a)) {
539 cb = Some(a.clone());
540 } else if !n.is_nan() && matches!(a, Value::Float(_) | Value::Int(_)) {
541 port = n as u16;
542 } else if with_host(|h| h.as_str(a)).is_some() {
543 host = with_host(|h| h.str_of(a));
544 } else if matches!(a, Value::Obj(_)) {
545 if let Some(p) = get_prop(a, "port") {
547 port = with_host(|h| h.to_number(&p)) as u16;
548 }
549 for key in ["host", "hostname"] {
550 if let Some(hv) = get_prop(a, key).filter(|v| with_host(|h| h.as_str(v)).is_some())
551 {
552 host = with_host(|h| h.str_of(&hv));
553 }
554 }
555 if let Some(sv) =
556 get_prop(a, "servername").filter(|v| with_host(|h| h.as_str(v)).is_some())
557 {
558 servername = Some(with_host(|h| h.str_of(&sv)));
559 }
560 if let Some(rv) = get_prop(a, "rejectUnauthorized") {
561 reject_unauthorized = with_host(|h| h.truthy(&rv));
562 }
563 }
564 }
565 let servername = servername.unwrap_or_else(|| host.clone());
566
567 let sock_id = next_tls_id();
569 let (tx, rx) = std::sync::mpsc::channel::<WriteCmd>();
570 let mut extra = IndexMap::new();
571 extra.insert("@@tlsid".into(), Value::Float(sock_id as f64));
572 extra.insert("authorized".into(), Value::Bool(reject_unauthorized));
573 extra.insert("encrypted".into(), Value::Bool(true));
574 let socket = new_emitter_object("TLSSocket", extra);
575 TLS.with(|s| {
576 s.borrow_mut().sockets.insert(
577 sock_id,
578 TlsSocketRec {
579 emitter: socket.clone(),
580 tx,
581 },
582 );
583 });
584 with_host(|h| h.incr_handle());
585 if let Some(cb) = cb {
587 super::events::instance_call(
588 &socket,
589 "once",
590 vec![with_host(|h| h.new_str("secureConnect")), cb],
591 )?;
592 }
593
594 let io_tx = with_host(|h| h.io_sender());
595 let config = if reject_unauthorized {
596 verifying_client_config()
597 } else {
598 insecure_client_config()
599 };
600
601 std::thread::spawn(move || {
602 let server_name = match ServerName::try_from(servername.clone()) {
603 Ok(n) => n,
604 Err(_) => {
605 post_socket_error(
606 &io_tx,
607 sock_id,
608 format!("Error: tls: invalid servername '{servername}'"),
609 );
610 return;
611 }
612 };
613 let mut sock = match TcpStream::connect((host.as_str(), port)) {
614 Ok(s) => s,
615 Err(e) => {
616 post_socket_error(
617 &io_tx,
618 sock_id,
619 format!("Error: connect ECONNREFUSED {host}:{port}: {e}"),
620 );
621 return;
622 }
623 };
624 let mut conn = match ClientConnection::new(config, server_name) {
625 Ok(c) => c,
626 Err(e) => {
627 post_socket_error(&io_tx, sock_id, format!("Error: tls: {e}"));
628 return;
629 }
630 };
631 if let Err(e) = conn.complete_io(&mut sock) {
633 post_socket_error(&io_tx, sock_id, format!("Error: tls handshake failed: {e}"));
634 return;
635 }
636 let _ = io_tx.send(Box::new(move || on_secure_connect(sock_id)));
637 let stream = StreamOwned::new(conn, sock);
638 owner_loop(stream, sock_id, rx, io_tx);
639 });
640
641 Ok(socket)
642}
643
644fn on_secure_connect(sock_id: u64) -> Result<(), String> {
646 let socket = TLS.with(|s| s.borrow().sockets.get(&sock_id).map(|r| r.emitter.clone()));
647 if let Some(socket) = socket {
648 super::events::instance_call(
649 &socket,
650 "emit",
651 vec![with_host(|h| h.new_str("secureConnect"))],
652 )?;
653 }
654 Ok(())
655}
656
657fn post_socket_error(io_tx: &Sender<IoTask>, sock_id: u64, msg: String) {
659 let _ = io_tx.send(Box::new(move || {
660 let socket = TLS.with(|s| s.borrow().sockets.get(&sock_id).map(|r| r.emitter.clone()));
661 if let Some(socket) = socket {
662 let err = with_host(|h| {
663 let mut m = IndexMap::new();
664 m.insert("message".into(), h.new_str(msg.clone()));
665 h.new_object(m)
666 });
667 super::events::instance_call(
668 &socket,
669 "emit",
670 vec![with_host(|h| h.new_str("error")), err],
671 )?;
672 }
673 on_socket_close(sock_id)
674 }));
675}
676
677pub fn create_server(args: &[Value]) -> Result<Value, String> {
683 let mut options: Option<Value> = None;
684 let mut listener: Option<Value> = None;
685 for a in args {
686 if with_host(|h| crate::host::is_callable(h, a)) {
687 listener = Some(a.clone());
688 } else if matches!(a, Value::Obj(_)) {
689 options = Some(a.clone());
690 }
691 }
692 let opts = options.ok_or_else(|| {
693 crate::host::type_error("tls.createServer requires an options object with `key` and `cert`")
694 })?;
695 let cert = value_bytes(get_prop(&opts, "cert").as_ref());
696 let key = value_bytes(get_prop(&opts, "key").as_ref());
697 if cert.is_empty() || key.is_empty() {
698 return Err(crate::host::type_error(
699 "tls.createServer requires `key` and `cert`",
700 ));
701 }
702 let config = build_server_config(&cert, &key)?;
703
704 let mut extra = IndexMap::new();
705 if let Some(cb) = listener {
706 extra.insert("@@connListener".into(), cb);
707 }
708 let server = new_emitter_object("TLSServer", extra);
709 PENDING_CONFIGS.with(|p| p.borrow_mut().push((server.clone(), config)));
710 Ok(server)
711}
712
713pub fn create_server_with_config(
717 config: Arc<ServerConfig>,
718 hook: ConnHook,
719 request_listener: Value,
720) -> Value {
721 let mut extra = IndexMap::new();
722 extra.insert("@@requestListener".into(), request_listener);
723 let server = new_emitter_object("TLSServer", extra);
724 PENDING_CONFIGS.with(|p| p.borrow_mut().push((server.clone(), config)));
725 PENDING_HOOKS.with(|p| p.borrow_mut().push((server.clone(), hook)));
726 server
727}
728
729fn take_pending_config(server: &Value) -> Option<Arc<ServerConfig>> {
730 PENDING_CONFIGS.with(|p| {
731 let mut p = p.borrow_mut();
732 p.iter()
733 .position(|(s, _)| s == server)
734 .map(|pos| p.remove(pos).1)
735 })
736}
737fn take_pending_hook(server: &Value) -> Option<ConnHook> {
738 PENDING_HOOKS.with(|p| {
739 let mut p = p.borrow_mut();
740 p.iter()
741 .position(|(s, _)| s == server)
742 .map(|pos| p.remove(pos).1)
743 })
744}
745
746pub fn instance_call(
749 tag: &str,
750 recv: &Value,
751 method: &str,
752 args: Vec<Value>,
753) -> Result<Value, String> {
754 match tag {
755 "TLSServer" => server_call(recv, method, args),
756 "TLSSocket" => socket_call(recv, method, args),
757 _ => Err(crate::host::type_error(&format!(
758 "{method} is not a function"
759 ))),
760 }
761}
762
763fn server_call(recv: &Value, method: &str, args: Vec<Value>) -> Result<Value, String> {
764 if let Some(r) = emitter_dispatch(recv, method, &args) {
765 return r;
766 }
767 match method {
768 "listen" => server_listen(recv, &args),
769 "close" => server_close(recv, &args),
770 "address" => Ok(get_prop(recv, "@@address").unwrap_or(Value::Undef)),
771 _ => Err(crate::host::type_error(&format!(
772 "server.{method} is not a function"
773 ))),
774 }
775}
776
777fn server_listen(recv: &Value, args: &[Value]) -> Result<Value, String> {
780 let port = with_host(|h| args.first().map(|v| h.to_number(v)).unwrap_or(0.0)) as u16;
781 let mut host = "0.0.0.0".to_string();
782 let mut cb: Option<Value> = None;
783 for a in &args[1.min(args.len())..] {
784 if with_host(|h| h.as_str(a)).is_some() {
785 host = with_host(|h| h.str_of(a));
786 } else if with_host(|h| crate::host::is_callable(h, a)) {
787 cb = Some(a.clone());
788 }
789 }
790
791 let config = take_pending_config(recv)
792 .ok_or_else(|| crate::host::type_error("tls server has no secure context"))?;
793 let listener = std::net::TcpListener::bind((host.as_str(), port))
794 .map_err(|e| format!("Error: listen EADDRINUSE: {e}"))?;
795 let local = listener.local_addr().ok();
796
797 let id = next_server_id();
798 set_prop(recv, "@@serverid", Value::Float(id as f64));
799 if let Some(addr) = local {
800 let mut a = IndexMap::new();
801 a.insert("port".into(), Value::Float(addr.port() as f64));
802 a.insert(
803 "address".into(),
804 with_host(|h| h.new_str(addr.ip().to_string())),
805 );
806 a.insert(
807 "family".into(),
808 with_host(|h| h.new_str(if addr.is_ipv6() { "IPv6" } else { "IPv4" })),
809 );
810 let addr_obj = with_host(|h| h.new_object(a));
811 set_prop(recv, "@@address", addr_obj);
812 }
813 let conn_hook = take_pending_hook(recv);
814 let request_listener = get_prop(recv, "@@requestListener");
815 let plain_listener = get_prop(recv, "@@connListener");
816 let stop = Arc::new(AtomicBool::new(false));
817 TLS.with(|s| {
818 s.borrow_mut().servers.insert(
819 id,
820 TlsServerRec {
821 emitter: recv.clone(),
822 stop: stop.clone(),
823 conn_hook,
824 listener: plain_listener,
825 },
826 );
827 });
828 let _ = request_listener;
831 with_host(|h| h.incr_handle());
832
833 let io_tx = with_host(|h| h.io_sender());
834 listener.set_nonblocking(true).ok();
835 std::thread::spawn(move || {
836 loop {
837 if stop.load(Ordering::Acquire) {
838 break;
839 }
840 match listener.accept() {
841 Ok((stream, _addr)) => {
842 let cfg = config.clone();
843 let tx = io_tx.clone();
844 std::thread::spawn(move || accept_connection(id, stream, cfg, tx));
846 }
847 Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
848 std::thread::sleep(std::time::Duration::from_millis(5));
849 }
850 Err(_) => break,
851 }
852 }
853 });
854
855 let server = recv.clone();
856 let _ = with_host(|h| h.io_sender()).send(Box::new(move || {
857 super::events::instance_call(&server, "emit", vec![with_host(|h| h.new_str("listening"))])?;
858 if let Some(cb) = cb {
859 invoke(&cb, Vec::new(), None)?;
860 }
861 Ok(())
862 }));
863 Ok(recv.clone())
864}
865
866fn accept_connection(
869 server_id: u64,
870 mut sock: TcpStream,
871 config: Arc<ServerConfig>,
872 io_tx: Sender<IoTask>,
873) {
874 sock.set_nonblocking(false).ok();
880 let mut conn = match ServerConnection::new(config) {
881 Ok(c) => c,
882 Err(_) => return,
883 };
884 if conn.complete_io(&mut sock).is_err() {
885 return;
886 }
887 let sock_id = next_tls_id();
888 let (tx, rx) = std::sync::mpsc::channel::<WriteCmd>();
889 let tx_for_main = tx;
890 let _ = io_tx.send(Box::new(move || {
891 on_server_connection(server_id, sock_id, tx_for_main)
892 }));
893 let stream = StreamOwned::new(conn, sock);
894 owner_loop(stream, sock_id, rx, io_tx);
895}
896
897fn on_server_connection(server_id: u64, sock_id: u64, tx: Sender<WriteCmd>) -> Result<(), String> {
900 let server = TLS.with(|s| {
901 s.borrow()
902 .servers
903 .get(&server_id)
904 .map(|r| r.emitter.clone())
905 });
906 let Some(server) = server else { return Ok(()) };
907
908 let mut extra = IndexMap::new();
909 extra.insert("@@tlsid".into(), Value::Float(sock_id as f64));
910 extra.insert("encrypted".into(), Value::Bool(true));
911 let socket = new_emitter_object("TLSSocket", extra);
912 TLS.with(|s| {
913 s.borrow_mut().sockets.insert(
914 sock_id,
915 TlsSocketRec {
916 emitter: socket.clone(),
917 tx,
918 },
919 );
920 });
921 with_host(|h| h.incr_handle());
922
923 super::events::instance_call(
925 &server,
926 "emit",
927 vec![with_host(|h| h.new_str("secureConnection")), socket.clone()],
928 )?;
929 super::events::instance_call(
930 &server,
931 "emit",
932 vec![with_host(|h| h.new_str("connection")), socket.clone()],
933 )?;
934
935 let hook = TLS.with(|s| {
937 s.borrow()
938 .servers
939 .get(&server_id)
940 .and_then(|r| r.conn_hook.clone())
941 });
942 if let Some(hook) = hook {
943 hook(&server, &socket, sock_id)?;
944 } else {
945 let listener = TLS.with(|s| {
946 s.borrow()
947 .servers
948 .get(&server_id)
949 .and_then(|r| r.listener.clone())
950 });
951 if let Some(cb) = listener {
952 invoke(&cb, vec![socket.clone()], None)?;
953 }
954 }
955 Ok(())
956}
957
958fn server_close(recv: &Value, args: &[Value]) -> Result<Value, String> {
959 if let Some(id) = u64_prop(recv, "@@serverid") {
960 let rec = TLS.with(|s| s.borrow_mut().servers.remove(&id));
961 if let Some(rec) = rec {
962 rec.stop.store(true, Ordering::Release);
963 with_host(|h| h.decr_handle());
964 let _ = with_host(|h| h.io_sender()).send(Box::new(|| Ok(())));
965 }
966 }
967 if let Some(cb) = args
968 .first()
969 .filter(|v| with_host(|h| crate::host::is_callable(h, v)))
970 {
971 invoke(cb, Vec::new(), None)?;
972 }
973 super::events::instance_call(recv, "emit", vec![with_host(|h| h.new_str("close"))])?;
974 Ok(recv.clone())
975}
976
977fn socket_call(recv: &Value, method: &str, args: Vec<Value>) -> Result<Value, String> {
980 if let Some(r) = emitter_dispatch(recv, method, &args) {
981 return r;
982 }
983 match method {
984 "write" => {
985 if let Some(id) = u64_prop(recv, "@@tlsid") {
986 socket_write(id, &value_bytes(args.first()));
987 }
988 Ok(Value::Bool(true))
989 }
990 "end" => {
991 if let Some(id) = u64_prop(recv, "@@tlsid") {
992 if let Some(chunk) = args.first().filter(|v| !matches!(v, Value::Undef)) {
993 socket_write(id, &value_bytes(Some(chunk)));
994 }
995 socket_end(id);
996 }
997 Ok(recv.clone())
998 }
999 "destroy" => {
1000 if let Some(id) = u64_prop(recv, "@@tlsid") {
1001 socket_end(id);
1002 }
1003 Ok(recv.clone())
1004 }
1005 "setEncoding" | "setTimeout" | "setNoDelay" | "setKeepAlive" | "ref" | "unref"
1006 | "pause" | "resume" => Ok(recv.clone()),
1007 _ => Err(crate::host::type_error(&format!(
1008 "socket.{method} is not a function"
1009 ))),
1010 }
1011}
1012
1013pub fn socket_write(sock_id: u64, data: &[u8]) {
1016 let tx = TLS.with(|s| s.borrow().sockets.get(&sock_id).map(|r| r.tx.clone()));
1017 if let Some(tx) = tx {
1018 let _ = tx.send(WriteCmd::Data(data.to_vec()));
1019 }
1020}
1021
1022pub fn socket_end(sock_id: u64) {
1024 let tx = TLS.with(|s| s.borrow().sockets.get(&sock_id).map(|r| r.tx.clone()));
1025 if let Some(tx) = tx {
1026 let _ = tx.send(WriteCmd::Shutdown);
1027 }
1028}
1029
1030fn owner_loop<C, S>(
1038 mut stream: StreamOwned<C, TcpStream>,
1039 sock_id: u64,
1040 rx: Receiver<WriteCmd>,
1041 io_tx: Sender<IoTask>,
1042) where
1043 C: DerefMut + Deref<Target = ConnectionCommon<S>>,
1044 S: SideData,
1045{
1046 stream
1047 .sock
1048 .set_read_timeout(Some(std::time::Duration::from_millis(20)))
1049 .ok();
1050 let mut buf = [0u8; 16384];
1051 loop {
1052 let mut shutdown = false;
1054 loop {
1055 match rx.try_recv() {
1056 Ok(WriteCmd::Data(bytes)) => {
1057 if stream
1058 .write_all(&bytes)
1059 .and_then(|_| stream.flush())
1060 .is_err()
1061 {
1062 let _ = io_tx.send(Box::new(move || on_socket_close(sock_id)));
1063 return;
1064 }
1065 }
1066 Ok(WriteCmd::Shutdown) => shutdown = true,
1067 Err(std::sync::mpsc::TryRecvError::Empty) => break,
1068 Err(std::sync::mpsc::TryRecvError::Disconnected) => break,
1069 }
1070 }
1071 if shutdown {
1072 stream.conn.send_close_notify();
1073 let _ = stream.flush();
1074 let _ = stream.sock.shutdown(std::net::Shutdown::Write);
1075 }
1076
1077 match stream.read(&mut buf) {
1079 Ok(0) => {
1080 let _ = io_tx.send(Box::new(move || on_socket_end(sock_id)));
1081 return;
1082 }
1083 Ok(n) => {
1084 let bytes = buf[..n].to_vec();
1085 let _ = io_tx.send(Box::new(move || on_socket_data(sock_id, bytes)));
1086 }
1087 Err(ref e)
1088 if matches!(
1089 e.kind(),
1090 std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
1091 ) =>
1092 {
1093 continue;
1095 }
1096 Err(_) => {
1097 let _ = io_tx.send(Box::new(move || on_socket_close(sock_id)));
1098 return;
1099 }
1100 }
1101 }
1102}
1103
1104fn on_socket_data(sock_id: u64, bytes: Vec<u8>) -> Result<(), String> {
1107 let socket = TLS.with(|s| s.borrow().sockets.get(&sock_id).map(|r| r.emitter.clone()));
1108 let Some(socket) = socket else { return Ok(()) };
1109 super::https::feed(sock_id, &socket, &bytes)?;
1111 let chunk = super::buffer::from_bytes(&bytes);
1112 super::events::instance_call(
1113 &socket,
1114 "emit",
1115 vec![with_host(|h| h.new_str("data")), chunk],
1116 )?;
1117 Ok(())
1118}
1119
1120fn on_socket_end(sock_id: u64) -> Result<(), String> {
1121 let socket = TLS.with(|s| s.borrow().sockets.get(&sock_id).map(|r| r.emitter.clone()));
1122 if let Some(socket) = socket {
1123 super::events::instance_call(&socket, "emit", vec![with_host(|h| h.new_str("end"))])?;
1124 }
1125 on_socket_close(sock_id)
1126}
1127
1128fn on_socket_close(sock_id: u64) -> Result<(), String> {
1129 let rec = TLS.with(|s| s.borrow_mut().sockets.remove(&sock_id));
1130 super::https::drop_conn(sock_id);
1131 if let Some(rec) = rec {
1132 super::events::instance_call(
1133 &rec.emitter,
1134 "emit",
1135 vec![with_host(|h| h.new_str("close"))],
1136 )?;
1137 with_host(|h| h.decr_handle());
1138 let _ = with_host(|h| h.io_sender()).send(Box::new(|| Ok(())));
1139 }
1140 Ok(())
1141}
1142
1143pub fn new_emitter_object(tag: &str, mut extra: IndexMap<String, Value>) -> Value {
1148 with_host(|h| {
1149 let on = h.new_object(IndexMap::new());
1150 let once = h.new_object(IndexMap::new());
1151 let mut m = IndexMap::new();
1152 m.insert("@@native".into(), h.new_str(tag));
1153 m.insert("@@on".into(), on);
1154 m.insert("@@once".into(), once);
1155 for (k, v) in extra.drain(..) {
1156 m.insert(k, v);
1157 }
1158 h.new_object(m)
1159 })
1160}