1use crate::host::{invoke, with_host, JsObj};
25use fusevm::Value;
26use indexmap::IndexMap;
27use rustls::pki_types::ServerName;
28use rustls::{ClientConnection, StreamOwned};
29use std::collections::HashMap;
30use std::io::{Read, Write};
31use std::net::TcpStream;
32
33pub const MODULE_METHODS: &[&str] = &["createServer", "request", "get"];
35
36pub const RESPONSE_METHODS: &[&str] = &[
39 "writeHead",
40 "setHeader",
41 "getHeader",
42 "getHeaderNames",
43 "getHeaders",
44 "hasHeader",
45 "removeHeader",
46 "write",
47 "end",
48 "flushHeaders",
49];
50pub const CLIENT_REQUEST_METHODS: &[&str] = &[
51 "write",
52 "end",
53 "setHeader",
54 "getHeader",
55 "removeHeader",
56 "abort",
57 "destroy",
58 "setTimeout",
59];
60
61pub fn constant(name: &str) -> Option<Value> {
67 match name {
68 "Agent" => Some(with_host(|h| h.alloc(JsObj::Builtin("https.Agent".into())))),
69 "globalAgent" => Some(with_host(|h| {
70 let mut m = IndexMap::new();
71 m.insert("@@native".into(), h.new_str("Agent"));
72 m.insert("maxSockets".into(), Value::Float(f64::INFINITY));
73 m.insert("protocol".into(), h.new_str("https:"));
74 h.new_object(m)
75 })),
76 _ => None,
77 }
78}
79
80pub fn call(method: &str, args: &[Value]) -> Option<Result<Value, String>> {
82 match method {
83 "createServer" => Some(create_server(args)),
84 "request" => Some(request(args, false)),
85 "get" => Some(request(args, true)),
86 _ => None,
87 }
88}
89
90fn get_prop(recv: &Value, key: &str) -> Option<Value> {
93 with_host(|h| match h.get(recv) {
94 Some(JsObj::Object(p)) => p.get(key).cloned(),
95 _ => None,
96 })
97}
98fn set_prop(recv: &Value, key: &str, val: Value) {
99 with_host(|h| {
100 if let Some(JsObj::Object(p)) = h.get_mut(recv) {
101 p.insert(key.to_string(), val);
102 }
103 });
104}
105fn u64_prop(recv: &Value, key: &str) -> Option<u64> {
106 get_prop(recv, key).map(|v| with_host(|h| h.to_number(&v)) as u64)
107}
108
109fn value_bytes(v: Option<&Value>) -> Vec<u8> {
111 let Some(v) = v else { return Vec::new() };
112 let is_buffer =
113 with_host(|h| matches!(h.get(v), Some(JsObj::Object(p)) if p.contains_key("@@bytes")));
114 if is_buffer {
115 return with_host(|h| match h.get(v) {
116 Some(JsObj::Object(p)) => match p.get("@@bytes").and_then(|b| h.get(b)) {
117 Some(JsObj::Array(items)) => items.iter().map(|x| h.to_number(x) as u8).collect(),
118 _ => Vec::new(),
119 },
120 _ => Vec::new(),
121 });
122 }
123 with_host(|h| h.str_of(v)).into_bytes()
124}
125
126struct HttpsConn {
129 listener: Value,
130 buf: Vec<u8>,
131}
132
133struct ResState {
134 sock_id: u64,
135 status: u16,
136 message: Option<String>,
137 headers: Vec<(String, String)>,
138 body: Vec<u8>,
139}
140
141thread_local! {
142 static CONNS: std::cell::RefCell<HashMap<u64, HttpsConn>> =
143 std::cell::RefCell::new(HashMap::new());
144 static RESPONSES: std::cell::RefCell<HashMap<u64, ResState>> =
145 std::cell::RefCell::new(HashMap::new());
146 static NEXT_RESID: std::cell::Cell<u64> = const { std::cell::Cell::new(1) };
147 static CLIENT_REQS: std::cell::RefCell<HashMap<u64, ClientReq>> =
148 std::cell::RefCell::new(HashMap::new());
149 static NEXT_REQID: std::cell::Cell<u64> = const { std::cell::Cell::new(1) };
150}
151
152fn next_resid() -> u64 {
153 NEXT_RESID.with(|c| {
154 let id = c.get();
155 c.set(id + 1);
156 id
157 })
158}
159fn next_reqid() -> u64 {
160 NEXT_REQID.with(|c| {
161 let id = c.get();
162 c.set(id + 1);
163 id
164 })
165}
166
167pub fn create_server(args: &[Value]) -> Result<Value, String> {
170 let mut options: Option<Value> = None;
171 let mut listener = Value::Undef;
172 for a in args {
173 if with_host(|h| crate::host::is_callable(h, a)) {
174 listener = a.clone();
175 } else if matches!(a, Value::Obj(_)) {
176 options = Some(a.clone());
177 }
178 }
179 let opts = options.ok_or_else(|| {
180 crate::host::type_error(
181 "https.createServer requires an options object with `key` and `cert`",
182 )
183 })?;
184 let cert = value_bytes(get_prop(&opts, "cert").as_ref());
185 let key = value_bytes(get_prop(&opts, "key").as_ref());
186 if cert.is_empty() || key.is_empty() {
187 return Err(crate::host::type_error(
188 "https.createServer requires `key` and `cert`",
189 ));
190 }
191 let config = super::tls::build_server_config(&cert, &key)?;
192
193 let listener_for_hook = listener.clone();
195 let hook: super::tls::ConnHook =
196 std::rc::Rc::new(move |_server: &Value, _socket: &Value, sock_id: u64| {
197 CONNS.with(|c| {
198 c.borrow_mut().insert(
199 sock_id,
200 HttpsConn {
201 listener: listener_for_hook.clone(),
202 buf: Vec::new(),
203 },
204 );
205 });
206 Ok(())
207 });
208 Ok(super::tls::create_server_with_config(
209 config, hook, listener,
210 ))
211}
212
213pub fn drop_conn(sock_id: u64) {
215 CONNS.with(|c| {
216 c.borrow_mut().remove(&sock_id);
217 });
218}
219
220pub fn feed(sock_id: u64, _socket: &Value, bytes: &[u8]) -> Result<(), String> {
225 let is_https = CONNS.with(|c| c.borrow().contains_key(&sock_id));
226 if !is_https {
227 return Ok(());
228 }
229 CONNS.with(|c| {
230 c.borrow_mut()
231 .get_mut(&sock_id)
232 .unwrap()
233 .buf
234 .extend_from_slice(bytes)
235 });
236
237 loop {
238 let (listener, parsed) = CONNS.with(|c| {
239 let mut c = c.borrow_mut();
240 let conn = c.get_mut(&sock_id).unwrap();
241 match parse_request(&conn.buf) {
242 Some((req, consumed)) => {
243 conn.buf.drain(..consumed);
244 (conn.listener.clone(), Some(req))
245 }
246 None => (Value::Undef, None),
247 }
248 });
249 let Some(parsed) = parsed else { break };
250
251 let req = build_incoming(&parsed);
252 let res = build_response(sock_id);
253 if with_host(|h| crate::host::is_callable(h, &listener)) {
254 invoke(&listener, vec![req.clone(), res], None)?;
255 }
256 if !parsed.body.is_empty() {
257 let chunk = super::buffer::from_bytes(&parsed.body);
258 super::events::instance_call(
259 &req,
260 "emit",
261 vec![with_host(|h| h.new_str("data")), chunk],
262 )?;
263 }
264 super::events::instance_call(&req, "emit", vec![with_host(|h| h.new_str("end"))])?;
265 }
266 Ok(())
267}
268
269struct ParsedReq {
271 method: String,
272 url: String,
273 http_version: String,
274 headers: Vec<(String, String)>,
275 body: Vec<u8>,
276}
277
278fn parse_request(buf: &[u8]) -> Option<(ParsedReq, usize)> {
281 let head_end = find_subslice(buf, b"\r\n\r\n")?;
282 let head = &buf[..head_end];
283 let body_start = head_end + 4;
284
285 let head_str = String::from_utf8_lossy(head);
286 let mut lines = head_str.split("\r\n");
287 let request_line = lines.next()?;
288 let mut parts = request_line.split(' ');
289 let method = parts.next()?.to_string();
290 let url = parts.next()?.to_string();
291 let version = parts.next().unwrap_or("HTTP/1.1");
292 let http_version = version.strip_prefix("HTTP/").unwrap_or("1.1").to_string();
293
294 let mut headers: Vec<(String, String)> = Vec::new();
295 let mut content_length = 0usize;
296 for line in lines {
297 if line.is_empty() {
298 continue;
299 }
300 if let Some((k, v)) = line.split_once(':') {
301 let name = k.trim().to_ascii_lowercase();
302 let value = v.trim().to_string();
303 if name == "content-length" {
304 content_length = value.parse().unwrap_or(0);
305 }
306 headers.push((name, value));
307 }
308 }
309 if buf.len() < body_start + content_length {
310 return None;
311 }
312 let body = buf[body_start..body_start + content_length].to_vec();
313 Some((
314 ParsedReq {
315 method,
316 url,
317 http_version,
318 headers,
319 body,
320 },
321 body_start + content_length,
322 ))
323}
324
325fn find_subslice(haystack: &[u8], needle: &[u8]) -> Option<usize> {
326 haystack.windows(needle.len()).position(|w| w == needle)
327}
328
329fn build_incoming(req: &ParsedReq) -> Value {
330 let headers_obj = with_host(|h| {
331 let mut m = IndexMap::new();
332 for (k, v) in &req.headers {
333 m.insert(k.clone(), h.new_str(v.clone()));
334 }
335 h.new_object(m)
336 });
337 let mut extra = IndexMap::new();
338 extra.insert(
339 "method".into(),
340 with_host(|h| h.new_str(req.method.clone())),
341 );
342 extra.insert("url".into(), with_host(|h| h.new_str(req.url.clone())));
343 extra.insert(
344 "httpVersion".into(),
345 with_host(|h| h.new_str(req.http_version.clone())),
346 );
347 extra.insert("headers".into(), headers_obj);
348 super::tls::new_emitter_object("IncomingMessage", extra)
351}
352
353fn build_response(sock_id: u64) -> Value {
354 let resid = next_resid();
355 RESPONSES.with(|r| {
356 r.borrow_mut().insert(
357 resid,
358 ResState {
359 sock_id,
360 status: 200,
361 message: None,
362 headers: Vec::new(),
363 body: Vec::new(),
364 },
365 );
366 });
367 let mut extra = IndexMap::new();
368 extra.insert("@@resid".into(), Value::Float(resid as f64));
369 extra.insert("statusCode".into(), Value::Float(200.0));
370 super::tls::new_emitter_object("HTTPSServerResponse", extra)
371}
372
373pub fn instance_call(
376 tag: &str,
377 recv: &Value,
378 method: &str,
379 args: Vec<Value>,
380) -> Result<Value, String> {
381 if matches!(
382 method,
383 "on" | "addListener"
384 | "prependListener"
385 | "once"
386 | "prependOnceListener"
387 | "emit"
388 | "removeListener"
389 | "off"
390 | "removeAllListeners"
391 | "listenerCount"
392 | "eventNames"
393 | "setMaxListeners"
394 | "getMaxListeners"
395 | "listeners"
396 ) {
397 return super::events::instance_call(recv, method, args);
398 }
399 match tag {
400 "HTTPSServerResponse" => response_call(recv, method, args),
401 "HTTPSClientRequest" => client_request_call(recv, method, args),
402 _ => Err(crate::host::type_error(&format!(
403 "{method} is not a function"
404 ))),
405 }
406}
407
408fn resid_of(res: &Value) -> Option<u64> {
409 u64_prop(res, "@@resid")
410}
411
412fn response_call(res: &Value, method: &str, args: Vec<Value>) -> Result<Value, String> {
413 let Some(resid) = resid_of(res) else {
414 return Err(crate::host::type_error("invalid ServerResponse"));
415 };
416 match method {
417 "writeHead" => {
418 let status =
419 with_host(|h| args.first().map(|v| h.to_number(v)).unwrap_or(200.0)) as u16;
420 let mut message: Option<String> = None;
421 let mut headers_arg: Option<Value> = None;
422 if let Some(a) = args.get(1) {
423 if with_host(|h| h.as_str(a)).is_some() {
424 message = Some(with_host(|h| h.str_of(a)));
425 } else if !matches!(a, Value::Undef) {
426 headers_arg = Some(a.clone());
427 }
428 }
429 if let Some(a) = args.get(2) {
430 if !matches!(a, Value::Undef) {
431 headers_arg = Some(a.clone());
432 }
433 }
434 let header_pairs = headers_arg.map(|h| object_pairs(&h)).unwrap_or_default();
435 RESPONSES.with(|r| {
436 if let Some(st) = r.borrow_mut().get_mut(&resid) {
437 st.status = status;
438 st.message = message;
439 for (k, v) in header_pairs {
440 upsert_header(&mut st.headers, &k, v);
441 }
442 }
443 });
444 set_prop(res, "statusCode", Value::Float(status as f64));
445 Ok(res.clone())
446 }
447 "setHeader" => {
448 let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)));
449 let v = with_host(|h| h.str_of(&args.get(1).cloned().unwrap_or(Value::Undef)));
450 RESPONSES.with(|r| {
451 if let Some(st) = r.borrow_mut().get_mut(&resid) {
452 upsert_header(&mut st.headers, &k, v);
453 }
454 });
455 Ok(Value::Undef)
456 }
457 "getHeader" => {
458 let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)))
459 .to_ascii_lowercase();
460 let val = RESPONSES.with(|r| {
461 r.borrow().get(&resid).and_then(|st| {
462 st.headers
463 .iter()
464 .find(|(hk, _)| hk.eq_ignore_ascii_case(&k))
465 .map(|(_, v)| v.clone())
466 })
467 });
468 Ok(val
469 .map(|v| with_host(|h| h.new_str(v)))
470 .unwrap_or(Value::Undef))
471 }
472 "removeHeader" => {
473 let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)));
474 RESPONSES.with(|r| {
475 if let Some(st) = r.borrow_mut().get_mut(&resid) {
476 st.headers.retain(|(hk, _)| !hk.eq_ignore_ascii_case(&k));
477 }
478 });
479 Ok(Value::Undef)
480 }
481 "flushHeaders" => Ok(Value::Undef),
482 "write" => {
483 let bytes = value_bytes(args.first());
484 RESPONSES.with(|r| {
485 if let Some(st) = r.borrow_mut().get_mut(&resid) {
486 st.body.extend_from_slice(&bytes);
487 }
488 });
489 Ok(Value::Bool(true))
490 }
491 "end" => {
492 if let Some(chunk) = args.first().filter(|v| !matches!(v, Value::Undef)) {
493 let bytes = value_bytes(Some(chunk));
494 RESPONSES.with(|r| {
495 if let Some(st) = r.borrow_mut().get_mut(&resid) {
496 st.body.extend_from_slice(&bytes);
497 }
498 });
499 }
500 finish_response(res, resid)?;
501 Ok(res.clone())
502 }
503 _ => Err(crate::host::type_error(&format!(
504 "res.{method} is not a function"
505 ))),
506 }
507}
508
509fn finish_response(res: &Value, resid: u64) -> Result<(), String> {
510 let js_status = u64_prop(res, "statusCode").map(|n| n as u16);
511 let st = RESPONSES.with(|r| r.borrow_mut().remove(&resid));
512 let Some(mut st) = st else { return Ok(()) };
513 if let Some(s) = js_status {
514 st.status = s;
515 }
516 let payload = serialize_response(&mut st);
517 super::tls::socket_write(st.sock_id, &payload);
518 super::tls::socket_end(st.sock_id);
519 super::events::instance_call(res, "emit", vec![with_host(|h| h.new_str("finish"))])?;
520 Ok(())
521}
522
523fn serialize_response(st: &mut ResState) -> Vec<u8> {
527 let reason = st
528 .message
529 .clone()
530 .unwrap_or_else(|| status_text(st.status).to_string());
531 let mut out = format!("HTTP/1.1 {} {}\r\n", st.status, reason).into_bytes();
532 let has = |name: &str| st.headers.iter().any(|(k, _)| k.eq_ignore_ascii_case(name));
533 let chunked = st.headers.iter().any(|(k, v)| {
534 k.eq_ignore_ascii_case("transfer-encoding") && v.to_ascii_lowercase().contains("chunked")
535 });
536 for (k, v) in &st.headers {
537 out.extend_from_slice(format!("{k}: {v}\r\n").as_bytes());
538 }
539 if !chunked && !has("content-length") {
540 out.extend_from_slice(format!("Content-Length: {}\r\n", st.body.len()).as_bytes());
541 }
542 if !has("connection") {
543 out.extend_from_slice(b"Connection: close\r\n");
544 }
545 out.extend_from_slice(b"\r\n");
546 out.extend_from_slice(&st.body);
547 out
548}
549
550fn upsert_header(headers: &mut Vec<(String, String)>, name: &str, value: String) {
551 if let Some(slot) = headers
552 .iter_mut()
553 .find(|(k, _)| k.eq_ignore_ascii_case(name))
554 {
555 slot.1 = value;
556 } else {
557 headers.push((name.to_string(), value));
558 }
559}
560
561fn object_pairs(obj: &Value) -> Vec<(String, String)> {
562 with_host(|h| match h.get(obj) {
563 Some(JsObj::Object(p)) => p
564 .iter()
565 .filter(|(k, _)| !k.starts_with("@@") && !k.starts_with('#'))
566 .map(|(k, v)| (k.clone(), h.str_of(v)))
567 .collect(),
568 _ => Vec::new(),
569 })
570}
571
572fn status_text(code: u16) -> &'static str {
573 for &(c, msg) in super::http::status_table() {
574 if c == code {
575 return msg;
576 }
577 }
578 "OK"
579}
580
581struct ClientReq {
586 host: String,
587 port: u16,
588 servername: String,
589 reject_unauthorized: bool,
590 method: String,
591 path: String,
592 headers: Vec<(String, String)>,
593 body: Vec<u8>,
594 request: Value,
596 sent: bool,
597}
598
599pub fn request(args: &[Value], is_get: bool) -> Result<Value, String> {
602 let mut host = "localhost".to_string();
603 let mut port: u16 = 443;
604 let mut path = "/".to_string();
605 let mut method = "GET".to_string();
606 let mut servername: Option<String> = None;
607 let mut reject_unauthorized = true;
608 let mut headers: Vec<(String, String)> = Vec::new();
609 let mut cb: Option<Value> = None;
610
611 for a in args {
612 if with_host(|h| crate::host::is_callable(h, a)) {
613 cb = Some(a.clone());
614 } else if with_host(|h| h.as_str(a)).is_some() {
615 let url = with_host(|h| h.str_of(a));
617 parse_url(&url, &mut host, &mut port, &mut path);
618 } else if matches!(a, Value::Obj(_)) {
619 for key in ["hostname", "host"] {
620 if let Some(v) = get_prop(a, key).filter(|v| with_host(|h| h.as_str(v)).is_some()) {
621 host = with_host(|h| h.str_of(&v));
622 }
623 }
624 if let Some(v) = get_prop(a, "port") {
625 let n = with_host(|h| h.to_number(&v));
626 if !n.is_nan() {
627 port = n as u16;
628 }
629 }
630 if let Some(v) = get_prop(a, "path").filter(|v| with_host(|h| h.as_str(v)).is_some()) {
631 path = with_host(|h| h.str_of(&v));
632 }
633 if let Some(v) = get_prop(a, "method").filter(|v| with_host(|h| h.as_str(v)).is_some())
634 {
635 method = with_host(|h| h.str_of(&v));
636 }
637 if let Some(v) =
638 get_prop(a, "servername").filter(|v| with_host(|h| h.as_str(v)).is_some())
639 {
640 servername = Some(with_host(|h| h.str_of(&v)));
641 }
642 if let Some(v) = get_prop(a, "rejectUnauthorized") {
643 reject_unauthorized = with_host(|h| h.truthy(&v));
644 }
645 if let Some(hv) = get_prop(a, "headers").filter(|v| matches!(v, Value::Obj(_))) {
646 for (k, val) in object_pairs(&hv) {
647 headers.push((k, val));
648 }
649 }
650 }
651 }
652 if is_get {
653 method = "GET".to_string();
654 }
655 let servername = servername.unwrap_or_else(|| host.clone());
656
657 let reqid = next_reqid();
658 let mut extra = IndexMap::new();
659 extra.insert("@@reqid".into(), Value::Float(reqid as f64));
660 extra.insert("method".into(), with_host(|h| h.new_str(method.clone())));
661 extra.insert("path".into(), with_host(|h| h.new_str(path.clone())));
662 let request = super::tls::new_emitter_object("HTTPSClientRequest", extra);
663 if let Some(cb) = cb {
666 super::events::instance_call(
667 &request,
668 "on",
669 vec![with_host(|h| h.new_str("response")), cb],
670 )?;
671 }
672 CLIENT_REQS.with(|c| {
673 c.borrow_mut().insert(
674 reqid,
675 ClientReq {
676 host,
677 port,
678 servername,
679 reject_unauthorized,
680 method,
681 path,
682 headers,
683 body: Vec::new(),
684 request: request.clone(),
685 sent: false,
686 },
687 );
688 });
689 if is_get {
691 dispatch_request(reqid)?;
692 }
693 Ok(request)
694}
695
696fn parse_url(url: &str, host: &mut String, port: &mut u16, path: &mut String) {
697 let rest = url.strip_prefix("https://").unwrap_or(url);
698 let (authority, p) = match rest.find('/') {
699 Some(i) => (&rest[..i], &rest[i..]),
700 None => (rest, "/"),
701 };
702 *path = if p.is_empty() {
703 "/".to_string()
704 } else {
705 p.to_string()
706 };
707 if let Some((h, port_str)) = authority.rsplit_once(':') {
708 *host = h.to_string();
709 if let Ok(n) = port_str.parse::<u16>() {
710 *port = n;
711 }
712 } else {
713 *host = authority.to_string();
714 *port = 443;
715 }
716}
717
718fn client_request_call(req: &Value, method: &str, args: Vec<Value>) -> Result<Value, String> {
719 let reqid = u64_prop(req, "@@reqid");
720 match method {
721 "write" => {
722 if let Some(id) = reqid {
723 let bytes = value_bytes(args.first());
724 CLIENT_REQS.with(|c| {
725 if let Some(r) = c.borrow_mut().get_mut(&id) {
726 r.body.extend_from_slice(&bytes);
727 }
728 });
729 }
730 Ok(Value::Bool(true))
731 }
732 "end" => {
733 if let Some(id) = reqid {
734 if let Some(chunk) = args.first().filter(|v| !matches!(v, Value::Undef)) {
735 let bytes = value_bytes(Some(chunk));
736 CLIENT_REQS.with(|c| {
737 if let Some(r) = c.borrow_mut().get_mut(&id) {
738 r.body.extend_from_slice(&bytes);
739 }
740 });
741 }
742 dispatch_request(id)?;
743 }
744 Ok(req.clone())
745 }
746 "setHeader" => {
747 if let Some(id) = reqid {
748 let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)));
749 let v = with_host(|h| h.str_of(&args.get(1).cloned().unwrap_or(Value::Undef)));
750 CLIENT_REQS.with(|c| {
751 if let Some(r) = c.borrow_mut().get_mut(&id) {
752 upsert_header(&mut r.headers, &k, v);
753 }
754 });
755 }
756 Ok(Value::Undef)
757 }
758 "getHeader" => {
759 let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)))
760 .to_ascii_lowercase();
761 let val = reqid.and_then(|id| {
762 CLIENT_REQS.with(|c| {
763 c.borrow().get(&id).and_then(|r| {
764 r.headers
765 .iter()
766 .find(|(hk, _)| hk.eq_ignore_ascii_case(&k))
767 .map(|(_, v)| v.clone())
768 })
769 })
770 });
771 Ok(val
772 .map(|v| with_host(|h| h.new_str(v)))
773 .unwrap_or(Value::Undef))
774 }
775 "removeHeader" => {
776 if let Some(id) = reqid {
777 let k = with_host(|h| h.str_of(&args.first().cloned().unwrap_or(Value::Undef)));
778 CLIENT_REQS.with(|c| {
779 if let Some(r) = c.borrow_mut().get_mut(&id) {
780 r.headers.retain(|(hk, _)| !hk.eq_ignore_ascii_case(&k));
781 }
782 });
783 }
784 Ok(Value::Undef)
785 }
786 "abort" | "destroy" | "setTimeout" => Ok(req.clone()),
787 _ => Err(crate::host::type_error(&format!(
788 "req.{method} is not a function"
789 ))),
790 }
791}
792
793fn dispatch_request(reqid: u64) -> Result<(), String> {
796 let sent = CLIENT_REQS.with(|c| c.borrow().get(&reqid).map(|r| r.sent).unwrap_or(true));
798 if sent {
799 return Ok(());
800 }
801 CLIENT_REQS.with(|c| {
802 if let Some(r) = c.borrow_mut().get_mut(&reqid) {
803 r.sent = true;
804 }
805 });
806
807 let (host, port, servername, reject, method, path, headers, body) = CLIENT_REQS.with(|c| {
808 let b = c.borrow();
809 let r = b.get(&reqid).unwrap();
810 (
811 r.host.clone(),
812 r.port,
813 r.servername.clone(),
814 r.reject_unauthorized,
815 r.method.clone(),
816 r.path.clone(),
817 r.headers.clone(),
818 r.body.clone(),
819 )
820 });
821
822 let config = super::tls::client_config(reject);
823 let io_tx = with_host(|h| h.io_sender());
824 with_host(|h| h.incr_handle());
825
826 let mut has_host = false;
828 let mut has_len = false;
829 let mut header_block = String::new();
830 for (k, v) in &headers {
831 if k.eq_ignore_ascii_case("host") {
832 has_host = true;
833 }
834 if k.eq_ignore_ascii_case("content-length") {
835 has_len = true;
836 }
837 if k.eq_ignore_ascii_case("connection") {
838 continue; }
840 header_block.push_str(&format!("{k}: {v}\r\n"));
841 }
842 let host_header = if port == 443 {
843 host.clone()
844 } else {
845 format!("{host}:{port}")
846 };
847 let mut request_bytes = format!("{method} {path} HTTP/1.1\r\n");
848 if !has_host {
849 request_bytes.push_str(&format!("Host: {host_header}\r\n"));
850 }
851 request_bytes.push_str(&header_block);
852 if !has_len && !body.is_empty() {
853 request_bytes.push_str(&format!("Content-Length: {}\r\n", body.len()));
854 }
855 request_bytes.push_str("Connection: close\r\n\r\n");
856 let mut wire = request_bytes.into_bytes();
857 wire.extend_from_slice(&body);
858
859 std::thread::spawn(move || {
860 let result = do_client_exchange(&host, port, &servername, config, &wire);
861 match result {
862 Ok(raw) => {
863 let _ = io_tx.send(Box::new(move || deliver_response(reqid, raw)));
864 }
865 Err(msg) => {
866 let _ = io_tx.send(Box::new(move || deliver_error(reqid, msg)));
867 }
868 }
869 });
870 Ok(())
871}
872
873fn do_client_exchange(
876 host: &str,
877 port: u16,
878 servername: &str,
879 config: std::sync::Arc<rustls::ClientConfig>,
880 request: &[u8],
881) -> Result<Vec<u8>, String> {
882 let server_name = ServerName::try_from(servername.to_string())
883 .map_err(|_| format!("Error: tls: invalid servername '{servername}'"))?;
884 let sock = TcpStream::connect((host, port))
885 .map_err(|e| format!("Error: connect ECONNREFUSED {host}:{port}: {e}"))?;
886 let conn =
887 ClientConnection::new(config, server_name).map_err(|e| format!("Error: tls: {e}"))?;
888 let mut stream = StreamOwned::new(conn, sock);
889 stream
890 .write_all(request)
891 .map_err(|e| format!("Error: https write: {e}"))?;
892 stream
893 .flush()
894 .map_err(|e| format!("Error: https flush: {e}"))?;
895 let mut raw = Vec::new();
896 let mut buf = [0u8; 16384];
900 loop {
901 match stream.read(&mut buf) {
902 Ok(0) => break,
903 Ok(n) => raw.extend_from_slice(&buf[..n]),
904 Err(ref e) if e.kind() == std::io::ErrorKind::UnexpectedEof => break,
905 Err(e) => {
906 if raw.is_empty() {
907 return Err(format!("Error: https read: {e}"));
908 }
909 break;
910 }
911 }
912 }
913 Ok(raw)
914}
915
916fn deliver_response(reqid: u64, raw: Vec<u8>) -> Result<(), String> {
919 let entry = CLIENT_REQS.with(|c| c.borrow_mut().remove(&reqid));
920 with_host(|h| h.decr_handle());
921 let _ = with_host(|h| h.io_sender()).send(Box::new(|| Ok(())));
922 let Some(entry) = entry else { return Ok(()) };
923
924 let (status, message, http_version, headers, body) = parse_response(&raw);
925
926 let headers_obj = with_host(|h| {
927 let mut m = IndexMap::new();
928 for (k, v) in &headers {
929 m.insert(k.clone(), h.new_str(v.clone()));
930 }
931 h.new_object(m)
932 });
933 let mut extra = IndexMap::new();
934 extra.insert("statusCode".into(), Value::Float(status as f64));
935 extra.insert("statusMessage".into(), with_host(|h| h.new_str(message)));
936 extra.insert("httpVersion".into(), with_host(|h| h.new_str(http_version)));
937 extra.insert("headers".into(), headers_obj);
938 let res = super::tls::new_emitter_object("IncomingMessage", extra);
939
940 super::events::instance_call(
942 &entry.request,
943 "emit",
944 vec![with_host(|h| h.new_str("response")), res.clone()],
945 )?;
946 if !body.is_empty() {
947 let chunk = super::buffer::from_bytes(&body);
948 super::events::instance_call(&res, "emit", vec![with_host(|h| h.new_str("data")), chunk])?;
949 }
950 super::events::instance_call(&res, "emit", vec![with_host(|h| h.new_str("end"))])?;
951 Ok(())
952}
953
954fn deliver_error(reqid: u64, msg: String) -> Result<(), String> {
955 let entry = CLIENT_REQS.with(|c| c.borrow_mut().remove(&reqid));
956 with_host(|h| h.decr_handle());
957 let _ = with_host(|h| h.io_sender()).send(Box::new(|| Ok(())));
958 if let Some(entry) = entry {
959 let err = with_host(|h| {
960 let mut m = IndexMap::new();
961 m.insert("message".into(), h.new_str(msg.clone()));
962 h.new_object(m)
963 });
964 super::events::instance_call(
965 &entry.request,
966 "emit",
967 vec![with_host(|h| h.new_str("error")), err],
968 )?;
969 }
970 Ok(())
971}
972
973fn parse_response(raw: &[u8]) -> (u16, String, String, Vec<(String, String)>, Vec<u8>) {
976 let head_end = find_subslice(raw, b"\r\n\r\n").unwrap_or(raw.len());
977 let head = String::from_utf8_lossy(&raw[..head_end]);
978 let body_start = (head_end + 4).min(raw.len());
979 let mut lines = head.split("\r\n");
980 let status_line = lines.next().unwrap_or("");
981 let mut sp = status_line.splitn(3, ' ');
982 let version = sp
983 .next()
984 .unwrap_or("HTTP/1.1")
985 .strip_prefix("HTTP/")
986 .unwrap_or("1.1")
987 .to_string();
988 let status = sp.next().and_then(|s| s.parse::<u16>().ok()).unwrap_or(0);
989 let message = sp.next().unwrap_or("").to_string();
990
991 let mut headers: Vec<(String, String)> = Vec::new();
992 let mut chunked = false;
993 for line in lines {
994 if line.is_empty() {
995 continue;
996 }
997 if let Some((k, v)) = line.split_once(':') {
998 let name = k.trim().to_ascii_lowercase();
999 let value = v.trim().to_string();
1000 if name == "transfer-encoding" && value.to_ascii_lowercase().contains("chunked") {
1001 chunked = true;
1002 }
1003 headers.push((name, value));
1004 }
1005 }
1006 let raw_body = &raw[body_start..];
1007 let body = if chunked {
1008 decode_chunked(raw_body)
1009 } else {
1010 raw_body.to_vec()
1011 };
1012 (status, message, version, headers, body)
1013}
1014
1015fn decode_chunked(mut data: &[u8]) -> Vec<u8> {
1018 let mut out = Vec::new();
1019 while let Some(nl) = find_subslice(data, b"\r\n") {
1020 let size_line = String::from_utf8_lossy(&data[..nl]);
1021 let size_hex = size_line.split(';').next().unwrap_or("").trim();
1022 let size = usize::from_str_radix(size_hex, 16).unwrap_or(0);
1023 if size == 0 {
1024 break;
1025 }
1026 let chunk_start = nl + 2;
1027 let chunk_end = (chunk_start + size).min(data.len());
1028 out.extend_from_slice(&data[chunk_start..chunk_end]);
1029 let next = chunk_end + 2;
1031 if next >= data.len() {
1032 break;
1033 }
1034 data = &data[next..];
1035 }
1036 out
1037}