Skip to main content

fetch_happen/
lib.rs

1use serde::{Deserialize, Serialize};
2use serde_json::Value;
3use std::collections::HashMap;
4use std::fmt;
5use wasm_bindgen::prelude::*;
6use wasm_bindgen_futures::JsFuture;
7use web_sys::{
8    AbortSignal, ReadableStream, ReadableStreamDefaultReader, Request as WebRequest, RequestInit,
9    Response as WebResponse,
10};
11
12pub use web_sys::{AbortController, RequestMode};
13
14pub type Result<T> = std::result::Result<T, Error>;
15
16/// Errors that can occur when making a request
17#[derive(Debug)]
18pub enum Error {
19    /// JavaScript error
20    JsError(JsValue),
21    /// HTTP error with status code
22    HttpError(u16, String),
23    /// JSON parsing error
24    JsonError(String),
25    /// Request was aborted
26    Aborted,
27}
28
29impl From<JsValue> for Error {
30    fn from(value: JsValue) -> Self {
31        Error::JsError(value)
32    }
33}
34
35impl From<serde_json::Error> for Error {
36    fn from(err: serde_json::Error) -> Self {
37        Error::JsonError(err.to_string())
38    }
39}
40
41impl fmt::Display for Error {
42    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
43        match self {
44            Error::JsError(e) => write!(f, "JavaScript error: {:?}", e),
45            Error::HttpError(status, msg) => write!(f, "HTTP error {}: {}", status, msg),
46            Error::JsonError(e) => write!(f, "JSON error: {}", e),
47            Error::Aborted => write!(f, "Request was aborted"),
48        }
49    }
50}
51
52impl std::error::Error for Error {}
53
54/// HTTP methods
55#[derive(Debug, Clone, Copy)]
56pub enum Method {
57    GET,
58    POST,
59    PUT,
60    DELETE,
61    PATCH,
62    HEAD,
63    OPTIONS,
64}
65
66impl Method {
67    fn as_str(&self) -> &'static str {
68        match self {
69            Method::GET => "GET",
70            Method::POST => "POST",
71            Method::PUT => "PUT",
72            Method::DELETE => "DELETE",
73            Method::PATCH => "PATCH",
74            Method::HEAD => "HEAD",
75            Method::OPTIONS => "OPTIONS",
76        }
77    }
78}
79
80/// A builder for HTTP requests
81pub struct RequestBuilder {
82    url: String,
83    method: Method,
84    headers: HashMap<String, String>,
85    body: Option<String>,
86    mode: RequestMode,
87    signal: Option<AbortSignal>,
88}
89
90impl RequestBuilder {
91    fn new(method: Method, url: impl Into<String>) -> Self {
92        Self {
93            url: url.into(),
94            method,
95            headers: HashMap::new(),
96            body: None,
97            mode: RequestMode::Cors,
98            signal: None,
99        }
100    }
101
102    /// Set a header
103    pub fn header(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
104        self.headers.insert(key.into(), value.into());
105        self
106    }
107
108    /// Set multiple headers
109    pub fn headers(mut self, headers: HashMap<String, String>) -> Self {
110        self.headers.extend(headers);
111        self
112    }
113
114    /// Set the request mode (Cors, NoCors, SameOrigin)
115    pub fn mode(mut self, mode: RequestMode) -> Self {
116        self.mode = mode;
117        self
118    }
119
120    /// Set the request body as a string
121    pub fn body(mut self, body: impl Into<String>) -> Self {
122        self.body = Some(body.into());
123        self
124    }
125
126    /// Set the request body as JSON
127    pub fn json<T: Serialize>(mut self, json: &T) -> Result<Self> {
128        let body = serde_json::to_string(json)?;
129        self.body = Some(body);
130        self.headers
131            .insert("Content-Type".to_string(), "application/json".to_string());
132        Ok(self)
133    }
134
135    /// Set an abort signal for the request
136    pub fn abort_signal(mut self, signal: AbortSignal) -> Self {
137        self.signal = Some(signal);
138        self
139    }
140
141    /// Send the request and get a Response
142    pub async fn send(self) -> Result<Response> {
143        let opts = RequestInit::new();
144        opts.set_method(self.method.as_str());
145        opts.set_mode(self.mode);
146
147        if let Some(body) = &self.body {
148            opts.set_body(&JsValue::from_str(body));
149        }
150
151        if let Some(signal) = &self.signal {
152            opts.set_signal(Some(signal));
153        }
154
155        let request = WebRequest::new_with_str_and_init(&self.url, &opts)?;
156        let headers = request.headers();
157
158        for (key, value) in &self.headers {
159            headers.set(key, value)?;
160        }
161
162        let window = web_sys::window()
163            .ok_or_else(|| Error::JsError(JsValue::from_str("Failed to get window")))?;
164
165        let resp_value = JsFuture::from(window.fetch_with_request(&request))
166            .await
167            .map_err(|e| {
168                // Check if this is an abort error
169                if let Some(error) = e.dyn_ref::<js_sys::Error>() {
170                    if error.name() == "AbortError" {
171                        return Error::Aborted;
172                    }
173                }
174                Error::JsError(e)
175            })?;
176        let web_response: WebResponse = resp_value
177            .dyn_into()
178            .map_err(|_| Error::JsError(JsValue::from_str("Response conversion failed")))?;
179
180        Ok(Response::from_web_response(web_response))
181    }
182}
183
184/// A response from a fetch request
185pub struct Response {
186    inner: WebResponse,
187}
188
189impl Response {
190    fn from_web_response(response: WebResponse) -> Self {
191        Self { inner: response }
192    }
193
194    /// Get the status code
195    pub fn status(&self) -> u16 {
196        self.inner.status()
197    }
198
199    /// Check if the response was successful (status 200-299)
200    pub fn ok(&self) -> bool {
201        self.inner.ok()
202    }
203
204    /// Get a header value
205    pub fn header(&self, name: &str) -> Result<Option<String>> {
206        Ok(self.inner.headers().get(name)?)
207    }
208
209    /// Get the response body as text
210    pub async fn text(&self) -> Result<String> {
211        let promise = self.inner.text().map_err(Error::JsError)?;
212        let text = JsFuture::from(promise).await?;
213
214        text.as_string()
215            .ok_or_else(|| Error::JsError(JsValue::from_str("Failed to convert to string")))
216    }
217
218    /// Get the response body as JSON
219    pub async fn json<T: for<'de> Deserialize<'de>>(&self) -> Result<T> {
220        let text = self.text().await?;
221        Ok(serde_json::from_str(&text)?)
222    }
223
224    /// Get the response body as a dynamic JSON value
225    pub async fn json_value(&self) -> Result<Value> {
226        self.json().await
227    }
228
229    /// Get the response body as bytes
230    pub async fn bytes(&self) -> Result<Vec<u8>> {
231        let promise = self.inner.array_buffer().map_err(Error::JsError)?;
232        let array_buffer = JsFuture::from(promise).await?;
233        let uint8_array = js_sys::Uint8Array::new(&array_buffer);
234        Ok(uint8_array.to_vec())
235    }
236
237    /// Ensure the response was successful, returning an error if not
238    pub fn error_for_status(self) -> Result<Self> {
239        if self.ok() {
240            Ok(self)
241        } else {
242            let status = self.status();
243            let text = format!("HTTP Error {}", status);
244            Err(Error::HttpError(status, text))
245        }
246    }
247
248    /// Get the response body as a readable stream
249    pub fn stream(&self) -> Result<ReadableStream> {
250        self.inner
251            .body()
252            .ok_or_else(|| Error::JsError(JsValue::from_str("No body in response")))
253    }
254
255    /// Get a stream reader for reading chunks from the response
256    pub fn stream_reader(&self) -> Result<StreamReader> {
257        let stream = self.stream()?;
258        let reader = stream
259            .get_reader()
260            .dyn_into::<ReadableStreamDefaultReader>()
261            .map_err(|_| Error::JsError(JsValue::from_str("Failed to get stream reader")))?;
262        Ok(StreamReader { reader })
263    }
264}
265
266/// A reader for streaming response bodies chunk by chunk
267pub struct StreamReader {
268    reader: ReadableStreamDefaultReader,
269}
270
271impl StreamReader {
272    /// Read the next chunk from the stream
273    /// Returns Ok(Some(bytes)) if a chunk is available
274    /// Returns Ok(None) if the stream is finished
275    pub async fn read_chunk(&self) -> Result<Option<Vec<u8>>> {
276        let result = JsFuture::from(self.reader.read()).await?;
277
278        let done = js_sys::Reflect::get(&result, &JsValue::from_str("done"))?
279            .as_bool()
280            .unwrap_or(false);
281
282        if done {
283            return Ok(None);
284        }
285
286        let value = js_sys::Reflect::get(&result, &JsValue::from_str("value"))?;
287        let uint8_array = js_sys::Uint8Array::new(&value);
288        Ok(Some(uint8_array.to_vec()))
289    }
290
291    /// Release the reader lock
292    pub fn cancel(self) -> Result<()> {
293        self.reader.release_lock();
294        Ok(())
295    }
296}
297
298/// Main client for making HTTP requests
299pub struct Client;
300
301impl Client {
302    /// Make a GET request
303    pub fn get(&self, url: impl Into<String>) -> RequestBuilder {
304        RequestBuilder::new(Method::GET, url)
305    }
306
307    /// Make a POST request
308    pub fn post(&self, url: impl Into<String>) -> RequestBuilder {
309        RequestBuilder::new(Method::POST, url)
310    }
311
312    /// Make a PUT request
313    pub fn put(&self, url: impl Into<String>) -> RequestBuilder {
314        RequestBuilder::new(Method::PUT, url)
315    }
316
317    /// Make a DELETE request
318    pub fn delete(&self, url: impl Into<String>) -> RequestBuilder {
319        RequestBuilder::new(Method::DELETE, url)
320    }
321
322    /// Make a PATCH request
323    pub fn patch(&self, url: impl Into<String>) -> RequestBuilder {
324        RequestBuilder::new(Method::PATCH, url)
325    }
326
327    /// Make a HEAD request
328    pub fn head(&self, url: impl Into<String>) -> RequestBuilder {
329        RequestBuilder::new(Method::HEAD, url)
330    }
331}
332
333/// Convenience function for making a GET request
334pub async fn get(url: impl Into<String>) -> Result<Response> {
335    Client.get(url).send().await
336}
337
338/// Convenience function for making a POST request with JSON body
339pub async fn post_json<T: Serialize>(url: impl Into<String>, json: &T) -> Result<Response> {
340    Client.post(url).json(json)?.send().await
341}
342
343#[cfg(all(feature = "examples", target_arch = "wasm32"))]
344pub mod examples {
345    use super::*;
346    use web_sys::console;
347    use wasm_bindgen::prelude::wasm_bindgen;
348
349    /// Example of streaming a large response body in chunks
350    #[wasm_bindgen]
351    pub async fn stream_large_file() {
352        let client = Client;
353        let url = "https://raw.githubusercontent.com/yaptown/yap/refs/heads/main/out/deu/frequency_lists/combined/frequencies.jsonl";
354
355        console::log_1(&"Starting streaming download...".into());
356
357        let response = match client.get(url).send().await {
358            Ok(r) => r,
359            Err(e) => {
360                console::error_1(&format!("Request failed: {}", e).into());
361                return;
362            }
363        };
364
365        let response = match response.error_for_status() {
366            Ok(r) => r,
367            Err(e) => {
368                console::error_1(&format!("HTTP error: {}", e).into());
369                return;
370            }
371        };
372
373        // Get a stream reader
374        let reader = match response.stream_reader() {
375            Ok(r) => r,
376            Err(e) => {
377                console::error_1(&format!("Failed to get stream reader: {}", e).into());
378                return;
379            }
380        };
381
382        let mut total_bytes = 0;
383        let mut chunk_count = 0;
384
385        // Read chunks until the stream is done
386        loop {
387            match reader.read_chunk().await {
388                Ok(Some(chunk)) => {
389                    total_bytes += chunk.len();
390                    chunk_count += 1;
391                    console::log_1(&format!("Received chunk {}: {} bytes", chunk_count, chunk.len()).into());
392                }
393                Ok(None) => break,
394                Err(e) => {
395                    console::error_1(&format!("Error reading chunk: {}", e).into());
396                    return;
397                }
398            }
399        }
400
401        console::log_1(&format!("✓ Total: {} bytes in {} chunks", total_bytes, chunk_count).into());
402    }
403
404    /// Example of streaming text content line by line
405    #[wasm_bindgen]
406    pub async fn stream_text_content() {
407        let client = Client;
408        let url = "https://raw.githubusercontent.com/yaptown/yap/refs/heads/main/out/deu/frequency_lists/combined/frequencies.jsonl";
409
410        console::log_1(&"Starting line-by-line streaming...".into());
411
412        let response = match client.get(url).send().await.and_then(|r| r.error_for_status()) {
413            Ok(r) => r,
414            Err(e) => {
415                console::error_1(&format!("Request failed: {}", e).into());
416                return;
417            }
418        };
419
420        let reader = match response.stream_reader() {
421            Ok(r) => r,
422            Err(e) => {
423                console::error_1(&format!("Failed to get stream reader: {}", e).into());
424                return;
425            }
426        };
427
428        let mut buffer = Vec::new();
429        let mut line_count = 0;
430
431        loop {
432            let chunk = match reader.read_chunk().await {
433                Ok(Some(c)) => c,
434                Ok(None) => break,
435                Err(e) => {
436                    console::error_1(&format!("Error reading chunk: {}", e).into());
437                    return;
438                }
439            };
440
441            buffer.extend_from_slice(&chunk);
442
443            // Process complete lines from the buffer
444            while let Some(newline_pos) = buffer.iter().position(|&b| b == b'\n') {
445                let line_bytes = buffer.drain(..=newline_pos).collect::<Vec<_>>();
446                let line = String::from_utf8_lossy(&line_bytes);
447                line_count += 1;
448
449                // Only log first few lines to avoid spam
450                if line_count <= 5 {
451                    console::log_1(&format!("Line {}: {}", line_count, line.trim()).into());
452                }
453            }
454        }
455
456        // Process any remaining data in the buffer
457        if !buffer.is_empty() {
458            let line = String::from_utf8_lossy(&buffer);
459            line_count += 1;
460            console::log_1(&format!("Last line: {}", line.trim()).into());
461        }
462
463        console::log_1(&format!("✓ Processed {} lines total", line_count).into());
464    }
465
466    /// Example of downloading with progress tracking
467    #[wasm_bindgen]
468    pub async fn download_with_progress() {
469        let client = Client;
470        let url = "https://raw.githubusercontent.com/yaptown/yap/refs/heads/main/out/deu/frequency_lists/combined/frequencies.jsonl";
471
472        console::log_1(&"Starting download with progress tracking...".into());
473
474        let response = match client.get(url).send().await.and_then(|r| r.error_for_status()) {
475            Ok(r) => r,
476            Err(e) => {
477                console::error_1(&format!("Request failed: {}", e).into());
478                return;
479            }
480        };
481
482        // Get content length if available
483        let content_length = response
484            .header("content-length")
485            .ok()
486            .flatten()
487            .and_then(|s| s.parse::<usize>().ok());
488
489        if let Some(total) = content_length {
490            console::log_1(&format!("Content-Length: {} bytes", total).into());
491        } else {
492            console::log_1(&"Content-Length not available".into());
493        }
494
495        let reader = match response.stream_reader() {
496            Ok(r) => r,
497            Err(e) => {
498                console::error_1(&format!("Failed to get stream reader: {}", e).into());
499                return;
500            }
501        };
502
503        let mut downloaded = Vec::new();
504        let mut last_logged_percent = 0;
505
506        loop {
507            let chunk = match reader.read_chunk().await {
508                Ok(Some(c)) => c,
509                Ok(None) => break,
510                Err(e) => {
511                    console::error_1(&format!("Error reading chunk: {}", e).into());
512                    return;
513                }
514            };
515
516            downloaded.extend_from_slice(&chunk);
517
518            if let Some(total) = content_length {
519                let progress = (downloaded.len() as f64 / total as f64) * 100.0;
520                let progress_int = progress as u32;
521
522                // Only log every 10%
523                if progress_int >= last_logged_percent + 10 {
524                    console::log_1(&format!("Progress: {:.1}% ({}/{})", progress, downloaded.len(), total).into());
525                    last_logged_percent = progress_int;
526                }
527            }
528        }
529
530        console::log_1(&format!("✓ Download complete: {} bytes", downloaded.len()).into());
531    }
532}