Reqwest : comment diffuser les réponses HTTP en streaming
Plongée technique dans reqwest, le client HTTP Rust utilisé par Parsec pour le streaming des réponses SSE en contexte navigateur web.

Thierry Leblond
CEO & Co-founder
4 mins
- Technologie
Nous travaillons à l’implémentation des Server-Sent Events (SSE) sur toutes les plateformes prises en charge par Parsec, en utilisant reqwest, un client HTTP écrit en Rust. Parmi ces plateformes, il doit prendre en charge un contexte navigateur web.
Nous voulions savoir comment l’implémentation était réalisée pour le contexte navigateur web, afin de vérifier qu’elle ne causait pas de problème de performance ou de fuite.
Versions
Pour l’étude, nous utiliserons les versions suivantes :
Contexte
Sur les plateformes natives, l’OS fournit des API pour ouvrir des sockets TCP/UDP. L’application utilise ensuite ces sockets comme bon lui semble — par exemple en ajoutant des couches SSL/TLS et HTTP pour implémenter le protocole HTTPS, ou en implémentant tout autre protocole comme BitTorrent ou Tor.
Sur le web, le navigateur fournit des API de bien plus haut niveau, comme HTTP ou WebSocket. C’est très pratique et cela permet d’implémenter facilement des applications web dans le navigateur (qui joue le rôle d’un OS en contexte web). Cependant, selon les détails d’implémentation de l’API, les résultats peuvent être moins efficaces qu’une implémentation native.
Le serveur lent
Implémentons un serveur qui met du temps à répondre. Cela devrait nous aider à reproduire ce qui se passerait lors de l’envoi de réponses SSE.
L’exemple suivant utilise Python pour implémenter un serveur envoyant une réponse un octet à la fois.
import socketimport asyncio
async def handle_client(loop: asyncio.AbstractEventLoop, client: socket.socket, _address): peer_name = client.getpeername() resp = ( b"HTTP/1.1 418 I'm a teapot\r\n" b"Content-Length: 5\r\n" b"Access-Control-Allow-Origin: *\r\n" b"Access-Control-Expose-Headers: The-Answer-You-Are-Looking-For\r\n" b"The-Answer-You-Are-Looking-For: 42\r\n" b"\r\n" ) await loop.sock_sendall(client, resp) # Simulate a slow body being sent. try: for c in b"12345": await asyncio.sleep(1) print(f"Sending byte {c} to {peer_name}") await loop.sock_sendall(client, c.to_bytes()) except (ConnectionResetError, BrokenPipeError): pass client.close()
async def run_server(server: socket.socket): loop = asyncio.get_running_loop() while True: client, address = await loop.sock_accept(server) await asyncio.create_task(handle_client(loop, client, address))
server = socket.socket(socket.AF_INET6, socket.SOCK_STREAM)server.bind(('localhost', 8418))server.listen(4)asyncio.run(run_server(server))Le client web : une approche naïve
L’exemple JavaScript suivant utilise la Fetch API pour requêter une ressource sur notre serveur lent.
async function client() { console.time("Fetch request"); let response = await fetch("http://localhost:8418"); console.timeLog("Fetch request", "Received response") console.assert(response.headers.get('The-Answer-You-Are-Looking-For') === '42'); console.timeEnd("Fetch request");}
client().then(() => console.log("Finished"))Le fetch() mettra plusieurs secondes (~5s dans mon cas, avec Firefox) avant de recevoir la réponse complète de notre serveur lent.
Remarque : ce code se comporte différemment sous
node— le fetch en lui-même prend ~20ms, mais le script met ~5s à se terminer.
Le client natif
Plutôt que de bloquer pour lire la réponse dans son ensemble, nous préférerions lire les données partielles dès qu’elles sont disponibles. En Rust, on ferait quelque chose comme ceci :
use reqwest::{self, header::HeaderValue};use std::time::{Duration, Instant};
#[tokio::main(flavor = "current_thread")]async fn main() -> anyhow::Result<()> { let timer = Instant::now(); let response = reqwest::get("http://localhost:8418").await?; assert_eq!( response.headers().get("TheAnswerYouAreLookingFor"), Some(&HeaderValue::from_static("42")) ); let elapsed = timer.elapsed();
assert!(elapsed < Duration::from_secs(5)); println!("Request took {}ms", elapsed.as_millis()); Ok(())}L’appel reqwest::get prend ~10ms.
Pouvoir éviter de traiter la requête dans son intégralité est utile. Par exemple, certains serveurs web peuvent interrompre une requête tôt si elle est invalide.
C’est pourquoi nous ne voulons pas nous reposer sur la Fetch API lorsque nous travaillons avec le SSE sur le web.
Le client doit pouvoir lire les données dès que possible, sans attendre l’arrivée de la réponse HTTP complète.
Le SSE dans Parsec
Le serveur Parsec doit envoyer des événements au client via le protocole SSE.
Le protocole SSE est construit sur HTTP et sérialise les événements dans la réponse HTTP. En théorie, la réponse HTTP ne se termine jamais, puisque la connexion doit rester ouverte pour continuer à envoyer des événements.
C’est comme une réponse HTTP très lente.
Le fait que la réponse HTTP ne se termine jamais ne fait pas bon ménage avec la Fetch API, qui attendrait sa fin.
Heureusement, sur le web nous pouvons utiliser la SSE API.
Dans Parsec, le SSE est implémenté via reqwest-eventsource, qui s’appuie sur reqwest pour le HTTP proprement dit. Voici la méthode qui traite la réponse HTTP de reqwest, dans reqwest-eventsource/src/event_source.rs:145-150 :
fn handle_response(&mut self, res: Response) { self.last_retry.take(); let mut stream = res.bytes_stream().eventsource(); stream.set_last_event_id(self.last_event_id.clone()); self.cur_stream.replace(Box::pin(stream));}La ligne importante ici est :
let mut stream = res.bytes_stream().eventsource();Où res est un reqwest::Response. La méthode .eventsource() est fournie par le trait eventsource-stream::Eventsource ; elle consomme le flux renvoyé par .bytes_stream(). .bytes_stream() consomme la réponse et renvoie un flux qui produit des octets.
Alors, que se passe-t-il sous le capot ?
L’objectif est de comprendre comment reqwest implémente Response::bytes_stream pour la cible wasm32-unknown-unknown.
Rétro-ingénierie
Reqwest
reqwest dispose d’une feature stream qui active Response::bytes_stream dans Cargo.toml :
[package]name = "reqwest"version = "0.11.18" # remember to update html_root_urldescription = "higher level HTTP client library"
# ...
[features]# ...
stream = ["tokio/fs", "tokio-util", "wasm-streams"]
# ...
[target.'cfg(target_arch = "wasm32")'.dependencies]js-sys = "0.3.45"serde_json = "1.0"wasm-bindgen = "0.2.68"wasm-bindgen-futures = "0.4.18"wasm-streams = { version = "0.2", optional = true }
[target.'cfg(target_arch = "wasm32")'.dependencies.web-sys]version = "0.3.25"features = [ "AbortController", "AbortSignal", "Headers", "Request", "RequestInit", "RequestMode", "Response", "Window", "FormData", "Blob", "BlobPropertyBag", "ServiceWorkerGlobalScope", "RequestCredentials", "File", "ReadableStream"]
# ...Lorsque la feature stream est activée, la dépendance wasm-streams est ajoutée. Cette dépendance n’est utilisée qu’une seule fois dans le code de reqwest, dans reqwest/src/wasm/response.rs:135-157. Ces lignes correspondent à l’implémentation de Response::bytes_stream 😄
impl Response { /// Convert the response into a `Stream` of `Bytes` from the body. #[cfg(feature = "stream")] pub fn bytes_stream(self) -> impl futures_core::Stream<Item = crate::Result<Bytes>> { let web_response = self.http.into_body(); let abort = self._abort; let body = web_response .body() .expect("could not create wasm byte stream"); let body = wasm_streams::ReadableStream::from_raw(body.unchecked_into()); Box::pin(body.into_stream().map(move |buf_js| { // Keep the abort guard alive as long as this stream is. let _abort = &abort; let buffer = Uint8Array::new( &buf_js .map_err(crate::error::wasm) .map_err(crate::error::decode)?, ); let mut bytes = vec![0; buffer.length() as usize]; buffer.copy_to(&mut bytes); Ok(bytes.into()) })) }}Essayons de comprendre ce que fait ce code :
-
self.http.into_body()renvoieweb_sys::Response, comme le définitbodydans reqwest/src/wasm/response.rs:19-26 :/// A Response to a submitted `Request`.pub struct Response {http: http::Response<web_sys::Response>,_abort: AbortGuard,// Boxed to save space (11 words to 1 word), and it's not accessed// frequently internally.url: Box<Url>,}http::Responseest un wrapper autour du type interneweb_sys::Response.into_bodyrenvoie ceweb_sys::Responseinterne. -
body.unchecked_into()convertit leweb_sys::Responseenwasm_streams::readable::sys::ReadableStream. -
wasm_streams::ReadableStream::from_raw()consomme cesys::ReadableStream. -
Enfin, il convertit ce
ReadableStreamenStream.
Wasm-streams
wasm-streams est une crate dont le but principal est de communiquer avec l’API de flux du navigateur :
Working with the Web Streams API in Rust.
This crate provides wrappers around
ReadableStream,WritableStreamandTransformStream. It also supports converting to and fromStreams andSinks from thefuturescrate.
Deux éléments fournis par cette crate nous intéressent :
wasm_streams::readable::sys::ReadableStream, un wrapper autour de l’API navigateurReadableStreamwasm_streams::readable::ReadableStream, une structure permettant de convertir un flux JS/Rust en flux Rust/JS. Dans le cas dereqwest, elle est utilisée pour convertir un flux JS en flux Rust.
Conclusion
L’implémentation de bytes_stream par reqwest pour WebAssembly s’appuie sur la Streams API du navigateur. Cela devrait permettre d’implémenter le SSE dans Parsec avec une utilisation efficace des ressources, en évitant les fuites en dehors de l’implémentation interne du navigateur.
Article rédigé par :
- Florian Bennetot
- Marcos Medrano
- Emmanuel Leblond
Commencez à sécuriser vos données sensibles dès aujourd’hui
Profitez d'un essai gratuit de 15 jours. Vous pouvez résilier à tout moment.