(
state_rc: Rc<RefCell<OpState>>,
resp: reqwest::Response,
)
| 89 | } |
| 90 | |
| 91 | fn handle_response( |
| 92 | state_rc: Rc<RefCell<OpState>>, |
| 93 | resp: reqwest::Response, |
| 94 | ) -> Result<ClientHttpResponse, AnyError> { |
| 95 | let mut resp_headers = HashMap::<String, String>::new(); |
| 96 | for (k, v) in resp.headers() { |
| 97 | let header_value = String::from_utf8(v.as_bytes().to_owned())?; |
| 98 | resp_headers.insert(k.to_string(), header_value); |
| 99 | } |
| 100 | |
| 101 | let status_code = resp.status(); |
| 102 | |
| 103 | // response body resource |
| 104 | let stream: BytesStream = Box::pin( |
| 105 | resp.bytes_stream() |
| 106 | .map(|r| r.map_err(|err| std::io::Error::new(std::io::ErrorKind::Other, err))), |
| 107 | ); |
| 108 | let stream_reader = StreamReader::new(stream); |
| 109 | let rid = state_rc |
| 110 | .borrow_mut() |
| 111 | .resource_table |
| 112 | .add(RequestReponseBodyResource { |
| 113 | body: AsyncRefCell::new(stream_reader), |
| 114 | cancel: CancelHandle::default(), |
| 115 | }); |
| 116 | |
| 117 | deno_core::unsync::spawn(async move { |
| 118 | tokio::time::sleep(Duration::from_secs(30)).await; |
| 119 | let mut borrowed = state_rc.borrow_mut(); |
| 120 | if borrowed.resource_table.has(rid) { |
| 121 | info!(%rid, "closing resource"); |
| 122 | if let Ok(r) = borrowed |
| 123 | .resource_table |
| 124 | .take::<RequestReponseBodyResource>(rid) |
| 125 | { |
| 126 | r.close() |
| 127 | }; |
| 128 | } |
| 129 | }); |
| 130 | |
| 131 | Ok(ClientHttpResponse { |
| 132 | body_resource_id: rid, |
| 133 | headers: resp_headers, |
| 134 | status_code: status_code.as_u16() as i32, |
| 135 | }) |
| 136 | } |
| 137 | |
| 138 | struct RequestBodyReceiver { |
| 139 | rx: mpsc::Receiver<std::io::Result<Vec<u8>>>, |
no test coverage detected