Docs
Data loading and streaming
Dedupe server loads, seed the browser query cache, apply cache policy, and stream slow page sections with Suspense.
Data loading is simple when a route owns one query. It becomes an architectural problem when a page grows into a tree of independent components: several components may need the same record, parent components start plumbing data they do not use, and uncoordinated fetches create duplicate database work or request waterfalls.
Proa's opinionated model lets components declare data by stable identity while one request-scoped resolver shares repeated reads. Independent Suspense boundaries can load concurrently, and origin and cache policy stay behind one DataLoader interface.
Server Data Pipeline
The complete server path fits into four files. Use the file tabs to follow one product request from its stable key to the rendered page:
use std::future::Future;
use bytes::Bytes;
use proa_core::{DataLoader, LoadError, LoadRequest, LoadValue};
#[derive(Clone)]
pub struct AppLoader {
// db: sqlx::PgPool,
// http: reqwest::Client,
}
impl DataLoader for AppLoader {
fn load(
&self,
request: &LoadRequest,
) -> impl Future<Output = Result<LoadValue, LoadError>> + Send {
let key = request.key().to_string();
async move {
match key.as_str() {
"product:sku_1" => {
let json = r#"{"title":"Trail Shoe","price_cents":12900}"#;
Ok(LoadValue::bytes(Bytes::from_static(json.as_bytes())))
}
_ => Err(LoadError::custom(format!("unknown load key: {key}"))),
}
}
}
}
use proa_core::LoadKey;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Deserialize, Serialize)]
pub struct Product {
pub title: String,
pub price_cents: u32,
}
pub fn product_key(id: &str) -> LoadKey {
LoadKey::named("product", id)
}
pub fn recommendations_key(id: &str) -> LoadKey {
LoadKey::typed("recommendations", id)
}
use std::future::Future;
use proa_core::{
DataLoader, RenderOutcome, WebContext, WebRender, WriteError,
};
use proa_macros::html;
use crate::data::product::{product_key, Product};
pub struct ProductPage {
pub id: &'static str,
}
impl<L: DataLoader> WebRender<L> for ProductPage {
fn render(
self,
cx: &mut WebContext<L>,
) -> RenderOutcome<impl Future<Output = Result<(), WriteError>>> {
RenderOutcome::pending(async move {
let product: Product = cx.load_json(product_key(self.id)).await?;
html! {
<main>
<h1>{product.title.as_str()}</h1>
<p>{proa_macros::text!("{} cents", product.price_cents)}</p>
</main>
}
.render(cx)
.resolve()
.await
})
}
}
use std::sync::Arc;
use proa_framework_axum::RouteResponse;
use proa_stream::DataResolver;
use crate::{
components::product_page::ProductPage,
data::loader::AppLoader,
};
pub async fn product_handler() -> impl axum::response::IntoResponse {
let app_loader = Arc::new(AppLoader {});
let loader = DataResolver::request_scope(app_loader);
RouteResponse::ssr_async_with_loader(ProductPage { id: "sku_1" }, loader)
}
Why Keyed Loading
Prefer keyed loading through WebContext. It keeps data requirements beside the components that own them while preventing repeated database or API work. Loading in the handler remains useful when one fetch has one clear owner and rendering cannot begin without it.
| Pattern | Use it when |
|---|---|
Keyed load through WebContext (preferred) | Data may be used by multiple components, should dedupe, or can stream independently. |
| Load in the Axum handler | One fetch has one clear owner and rendering cannot begin without it. |
The model draws from Haxl's request-scoped deduplication and TanStack Query's keyed browser cache. Proa uses explicit concurrency and does not automatically batch distinct keys.
How It Works
LoadKey identifies a resource, DataLoader fetches it and owns cache policy, and DataResolver shares repeated same-key loads across every component in one render. Use LoadKey::named for string ids or LoadKey::typed for values that implement Hash; LoadValue can carry JSON, raw bytes, or a typed Rust value.
A later request loads the resource again unless the DataLoader adds cross-request caching with StandardDataLoader or a custom shared cache. LoadHints can direct cache policies supported by that loader. Components request values through WebContext from inside WebRender. See Data caching.
Continue The Data Into The Browser
Server and browser caching are separate layers with an explicit handoff. LoadKey and DataResolver coordinate work while producing the response. If an interactive island needs to keep observing that data, pass the loaded value to rsjs::query_with_initial(...). RSJS serializes it once into the response's query seed table, then installs it in the page-scoped browser cache under a structural QueryKey.
| Phase | Identity and owner | What is shared |
|---|---|---|
| Server render | LoadKey through one request-scoped DataResolver | Origin work across server components |
| HTML handoff | query_with_initial(...) | The server result becomes the browser's initial query value |
| Hydrated page | QueryKey through the RSJS query cache | Cached data and one in-flight fetch across islands |
| Browser mutation | Mutation invalidation keys | Fresh data is refetched for every observer of the affected query |
Keep the same logical namespace and id on both sides—for example, product plus sku_1—but remember that LoadKey and QueryKey are different types with different lifetimes. Proa does not expose the server loader to the browser. The query_with_initial(...) call is the deliberate bridge.
use std::future::Future;
use proa_core::{
DataLoader, RenderOutcome, WebContext, WebRender, WriteError,
};
use crate::data::product::{product_key, Product};
pub struct ProductPrice;
#[rsjs::rsjs(component, client)]
impl<L: DataLoader> WebRender<L> for ProductPrice {
fn render(
self,
cx: &mut WebContext<L>,
) -> RenderOutcome<impl Future<Output = Result<(), WriteError>>> {
RenderOutcome::pending(async move {
let initial: Product = cx.load_json(product_key("sku_1")).await?;
let product: rsjs::Query<Product> = rsjs::query_with_initial(
rsjs::QueryKey::new(("product", "sku_1")),
initial,
rsjs::FetchRequest {
url: "/api/products/sku_1",
response: Some(rsjs::ResponseKind::Json),
..Default::default()
},
rsjs::QueryOptions {
stale_time_ms: 30_000,
refetch_on_focus: true,
..Default::default()
},
);
proa_macros::html! {
<article aria-live="polite">
<h2>{product.data().title}</h2>
<p>{product.data().price_cents}" cents"</p>
<button
type="button"
disabled={product.is_refetching()}
on_click={|_| product.refetch()}
>
{if product.is_refetching() { "Refreshing..." } else { "Refresh" }}
</button>
</article>
}
.render(cx)
.resolve()
.await
})
}
}
use std::future::Future;
use proa_core::{
DataLoader, RenderOutcome, WebContext, WebRender, WriteError,
};
pub struct ProductEditor;
#[rsjs::rsjs(component, client)]
impl<L: DataLoader> WebRender<L> for ProductEditor {
fn render(
self,
cx: &mut WebContext<L>,
) -> RenderOutcome<impl Future<Output = Result<(), WriteError>>> {
RenderOutcome::pending(async move {
let restock = rsjs::mutation::<i64, ()>(
// Identity of this mutation and its reactive state.
rsjs::QueryKey::new(("product", "sku_1", "restock")),
rsjs::FetchRequest {
url: "/api/products/sku_1/stock",
method: Some(rsjs::HttpMethod::Patch),
credentials: Some(rsjs::FetchCredentials::SameOrigin),
response: Some(rsjs::ResponseKind::Empty),
..Default::default()
},
// Query cache entries to refetch after a successful restock.
&[rsjs::QueryKey::new(("product", "sku_1"))],
);
proa_macros::html! {
<button
type="button"
disabled={restock.is_pending()}
on_click={|_| restock.run(1_i64)}
>
{if restock.is_pending() { "Saving..." } else { "Add inventory" }}
</button>
}
.render(cx)
.resolve()
.await
})
}
}
The mutation key identifies the write operation; the invalidation list names the query keys whose cached data is stale after that write succeeds.
This gives the page several useful properties:
- Useful HTML immediately. The first paint contains real data and does not wait for a browser fetch or hydration.
- No throwaway bootstrap request. The value loaded for SSR seeds the query cache; the browser refetches according to staleness policy rather than starting from an empty cache.
- Deduplication in both environments. Server components share one origin load per
LoadKey; hydrated islands share one cache entry and one in-flight browser request perQueryKey. - Local ownership without prop plumbing. A component can declare the data it needs on the server and the freshness behavior it needs in the browser.
- Built-in freshness tools. Queries support staleness windows, retries, polling, focus refetching, manual refetching, and shared invalidation after mutations.
- Private origins stay private. The server loader can use a database or internal service directly, while browser refreshes go through an authenticated public HTTP endpoint.
- JavaScript stays scoped. Static consumers remain server-only; Proa includes the RSJS query runtime only when rendered browser work actually uses a query or mutation.
The browser endpoint must return the same serialized shape used for the initial value. Mutations should invalidate the structural query keys whose data they changed. See Queries and mutations for retries, polling, progress, streaming text, mutation state, and cache lifetime details.
JSON, Bytes, And Typed Values
Choose the WebContext load helper that matches the value shape:
| Helper | Loader value | Use it for |
|---|---|---|
cx.load_json::<T>(key) | LoadValue::bytes(...) | JSON data that should deserialize into T. |
cx.load_bytes(key) | LoadValue::bytes(...) | Raw bytes such as pre-rendered fragments or API payloads. |
cx.load_typed::<T>(key) | LoadValue::typed(value) | Already-materialized Rust values stored in a loader or cache. |
Typed values return Arc<T>, so multiple components can share the same materialized value without cloning a large structure.
Cache Hints
LoadHints carry advisory policy without changing the identity of the data:
use std::time::Duration;
use proa_core::{CacheHint, CacheMode, LoadHints, LoadKey};
let product = cx
.load_json_with::<Product>(
LoadKey::named("product", "sku_1"),
LoadHints::new().with_cache(
CacheHint::ttl(Duration::from_secs(30)).with_mode(CacheMode::Refresh),
),
)
.await?;
Your loader can honor the hints, ignore them, or map them onto its own cache rules. See Data caching for the full cache-boundary model.
Cross-Request Cache
For a simple in-process cache, wrap an upstream loader with StandardDataLoader:
use std::sync::Arc;
use std::time::Duration;
use proa_core::StandardDataLoader;
let upstream = AppLoader {};
let cached = Arc::new(StandardDataLoader::new(upstream, Duration::from_secs(30)));
StandardDataLoader caches by LoadKey, applies the request TTL hint when present, and supports manual invalidation with invalidate(&key) or clear().
Use a process-local cache only when that matches your deployment model. For horizontally scaled apps, put shared cache policy in Redis, your database, your CDN, or the origin service. See Data caching for loader invalidation and Page caching for HTTP header patterns.
Streaming Slow Sections
Use RouteResponse::ssr_stream_with_loader and Suspense when a page has a fast shell plus independent slow sections:
use proa_macros::{html, html_sync};
use proa_stream::Suspense;
pub async fn product_handler() -> impl axum::response::IntoResponse {
let loader = DataResolver::request_scope(Arc::new(AppLoader {}));
RouteResponse::ssr_stream_with_loader(
html! {
<main>
<h1>"Product"</h1>
{Suspense::new(
Recommendations { product_id: "sku_1" },
html_sync! { <section>"Loading recommendations"</section> },
)}
</main>
},
loader,
)
}
The fallback renders into the shell immediately. The child renders when its future completes. Streaming responses default to out-of-order chunked rendering.
Out-of-order boundaries can opt into scheduling metadata:
use std::time::Duration;
use proa_core::Priority;
let recommendations = Suspense {
child: Recommendations { product_id: "sku_1" },
fallback: html_sync! { <section>"Recommendations unavailable"</section> },
priority: Priority::Optional,
timeout: Some(Duration::from_millis(750)),
};
All boundary futures start concurrently. Priority is a deterministic tie-break when multiple boundaries are ready in the same driver turn: Critical, then Optional, then Defer, with boundary id breaking equal-priority ties. It does not delay lower-priority work.
The timeout starts when the driver admits the boundary after rendering the shell. If it expires, Proa cancels that future, emits <!--proa:timeout:N--> for observability, preserves the shell fallback, and continues closing the response. The default framework adapter uses Tokio time. Runtime-independent hosts can drive proa_stream::driver::out_of_order::OutOfOrderBoundaryDriver::next_deadline() and expire(now) with their own clock; the clockless convenience adapter expires configured timeouts at admission rather than silently allowing an unbounded wait.
Buffered/static HTML, in-order HTML, and Markdown preserve source order and intentionally ignore both options. Cancelling an in-order child after it flushes could otherwise leave partial markup on the wire.
Streaming Markdown Output
Suspense also works in Markdown render trees. Use this when a route or adapter should return generated text/markdown for agents, CLI clients, feeds, or high-density text views.
The low-level driver lives in proa_stream:
use std::{rc::Rc, sync::Arc};
use proa_core::{DataLoader, MdRender, StreamChunkFlusher, WriteError};
use proa_stream::render_md_in_order_to_flusher_with_loader;
pub async fn render_markdown<L: DataLoader, R: MdRender<L>>(
page: R,
loader: Arc<L>,
flusher: Rc<dyn StreamChunkFlusher>,
) -> Result<(), WriteError> {
render_md_in_order_to_flusher_with_loader(page, loader, flusher).await
}
The driver renders the root through WebContext<L> in in-order mode with the supplied request-scoped loader, flushes parent Markdown before each streamed Suspense boundary, renders the child into the same buffer, then continues the document in order. Markdown-local state such as blank-line tracking is carried across the child boundary so streamed output matches buffered output.
Markdown Suspense boundaries preserve document order. The fallback is an HTML shell concern and is not emitted by the Markdown stream driver.
Preload Known Keys
If you already know which keys a streamed route will need, start them before the component tree reaches each boundary.
On the typed path, wrap the loader explicitly:
use std::sync::Arc;
use proa_core::LoadKey;
use proa_stream::{DataResolver, PreparedLoader};
let app_loader = Arc::new(AppLoader {});
let resolver = DataResolver::request_scope(app_loader);
let loader = Arc::new(PreparedLoader::new_typed(
resolver,
[
LoadKey::named("product", "sku_1"),
LoadKey::named("recommendations", "sku_1"),
],
));
RouteResponse::ssr_stream_with_loader(ProductPage { id: "sku_1" }, loader)
On the dyn path, use the response helper:
RouteResponse::ssr_stream_with_loader_dyn(ProductPage { id: "sku_1" }, loader)
.with_preload_keys([LoadKey::named("product", "sku_1")])
Use with_preload_requests(...) when a preload needs cache hints.
Pipeline CPU-Bound Work
DataLoader can coordinate expensive computation as well as I/O. If independent components need cryptographic verification, compression, image processing, or similar CPU-heavy results, route that work through a bounded application pipeline instead of running it on Tokio's async worker threads.
The complete pipeline spans five files. Use the tabs to follow a verification request from the component, through DataLoader, into the bounded CPU worker pipeline:
use std::sync::Arc;
use tokio::sync::{mpsc, oneshot};
use tokio::task::JoinSet;
#[derive(Debug)]
pub struct Verification {
pub valid: bool,
}
type CryptoResult = Result<Verification, String>;
type CryptoFn = dyn Fn(Vec<u8>) -> CryptoResult + Send + Sync;
struct CryptoJob {
input: Vec<u8>,
reply: oneshot::Sender<CryptoResult>,
}
#[derive(Clone)]
pub struct CryptoPipeline {
tx: mpsc::Sender<CryptoJob>,
}
impl CryptoPipeline {
pub fn start<F>(parallelism: usize, work: F) -> Self
where
F: Fn(Vec<u8>) -> CryptoResult + Send + Sync + 'static,
{
assert!(parallelism > 0);
let (tx, mut rx) = mpsc::channel::<CryptoJob>(parallelism * 2);
let work: Arc<CryptoFn> = Arc::new(work);
tokio::spawn(async move {
let mut running = JoinSet::new();
loop {
tokio::select! {
Some(job) = rx.recv(), if running.len() < parallelism => {
let work = Arc::clone(&work);
running.spawn_blocking(move || {
let result = work(job.input);
let _ = job.reply.send(result);
});
}
Some(_) = running.join_next(), if !running.is_empty() => {}
else => break,
}
}
while running.join_next().await.is_some() {}
});
Self { tx }
}
pub async fn run(&self, input: Vec<u8>) -> CryptoResult {
let (reply, result) = oneshot::channel();
self.tx
.send(CryptoJob { input, reply })
.await
.map_err(|_| "crypto pipeline stopped".to_owned())?;
result
.await
.map_err(|_| "crypto worker stopped".to_owned())?
}
}
use std::collections::HashMap;
use std::future::Future;
use std::sync::Arc;
use proa_core::{DataLoader, LoadError, LoadKey, LoadRequest, LoadValue};
use super::crypto_pipeline::CryptoPipeline;
pub fn verification_key(document_id: &str) -> LoadKey {
LoadKey::named("verification", document_id)
}
#[derive(Clone)]
pub struct AppLoader {
documents: Arc<HashMap<String, Vec<u8>>>,
crypto: CryptoPipeline,
}
impl AppLoader {
pub fn new(documents: HashMap<String, Vec<u8>>, crypto: CryptoPipeline) -> Self {
Self {
documents: Arc::new(documents),
crypto,
}
}
}
impl DataLoader for AppLoader {
fn load(
&self,
request: &LoadRequest,
) -> impl Future<Output = Result<LoadValue, LoadError>> + Send {
let documents = Arc::clone(&self.documents);
let crypto = self.crypto.clone();
let key = request.key().to_string();
async move {
let id = key
.strip_prefix("verification:")
.ok_or_else(|| LoadError::custom(format!("unknown load key: {key}")))?;
let input = documents
.get(id)
.cloned()
.ok_or_else(|| LoadError::custom("document not found"))?;
let verification = crypto.run(input).await.map_err(LoadError::custom)?;
Ok(LoadValue::typed(verification))
}
}
}
use std::future::Future;
use std::sync::Arc;
use proa_core::{DataLoader, RenderOutcome, WebContext, WebRender, WriteError};
use proa_macros::html;
use crate::data::crypto_pipeline::Verification;
use crate::data::loader::verification_key;
pub struct VerificationPanel {
pub document_id: String,
}
impl<L: DataLoader> WebRender<L> for VerificationPanel {
fn render(
self,
cx: &mut WebContext<L>,
) -> RenderOutcome<impl Future<Output = Result<(), WriteError>>> {
RenderOutcome::pending(async move {
let result: Arc<Verification> =
cx.load_typed(verification_key(&self.document_id)).await?;
html! {
<section aria-label="Signature verification">
{if result.valid { "Valid signature" } else { "Invalid signature" }}
</section>
}
.render(cx)
.resolve()
.await
})
}
}
use std::sync::Arc;
use axum::extract::{Path, State};
use axum::response::IntoResponse;
use proa_framework_axum::RouteResponse;
use proa_stream::DataResolver;
use crate::components::verification_panel::VerificationPanel;
use crate::data::loader::AppLoader;
pub async fn document(
Path(id): Path<String>,
State(app_loader): State<Arc<AppLoader>>,
) -> impl IntoResponse {
let loader = DataResolver::request_scope(app_loader);
RouteResponse::ssr_async_with_loader(
VerificationPanel { document_id: id },
loader,
)
}
use std::collections::HashMap;
use std::sync::Arc;
use axum::routing::get;
use axum::Router;
use crate::data::crypto_pipeline::{CryptoPipeline, Verification};
use crate::data::loader::AppLoader;
use crate::routes::documents::document;
pub fn router(documents: HashMap<String, Vec<u8>>) -> Router {
let parallelism =
std::thread::available_parallelism().map_or(1, |threads| threads.get());
let crypto = CryptoPipeline::start(parallelism, verify_signature);
let loader = Arc::new(AppLoader::new(documents, crypto));
Router::new()
.route("/documents/:id", get(document))
.with_state(loader)
}
fn verify_signature(input: Vec<u8>) -> Result<Verification, String> {
// Replace with the CPU-bound cryptography operation.
Ok(Verification {
valid: !input.is_empty(),
})
}
src/app.rs creates the pipeline once, while src/routes/documents.rs creates a fresh request-local DataResolver for each response. DataResolver deduplicates identical verification keys within that render; the shared pipeline controls concurrency across distinct jobs and requests. Replace the in-memory document map and verify_signature body with the application's storage and cryptography implementation.
Start known jobs with PreparedLoader, or place independent consumers in separate Suspense boundaries, so computation overlaps with other data loading and rendering instead of beginning only when the page must await its result.
Moving one indivisible calculation to another thread keeps the async runtime responsive, but does not make that calculation finish sooner. To reduce its own latency, split it into independent jobs or use a crypto implementation or CPU pool that parallelizes internally. Tokio recommends spawn_blocking for bounded blocking work and a specialized executor such as Rayon for substantial CPU-bound computation. Blocking jobs cannot be aborted after they start, so reject expired work before admission and keep the queue bounded. See Tokio's CPU-bound task guidance.