Federated relay tier: peer flooding, per-member isolation, registry admission and relay discovery

This commit is contained in:
George Coles
2026-09-15 06:49:43 -04:00
parent d30e612be8
commit 2c97fd523f
11 changed files with 759 additions and 177 deletions
+13
View File
@@ -277,6 +277,7 @@ pub fn registry_init(dir: &Path) -> Result<()> {
issued_at: now_ts(),
ma_key: String::new(),
members: Vec::new(),
relays: Vec::new(),
};
let signed = registry::sign_registry(doc, &ma);
registry::save_registry(&registry_path, &signed)?;
@@ -390,11 +391,23 @@ pub fn registry_list(dir: &Path) -> Result<()> {
Ok(())
}
pub fn registry_set_relays(dir: &Path, relays: &[String]) -> Result<()> {
mutate_registry(dir, |doc| {
doc.relays = relays.to_vec();
Ok(())
})?;
println!("relays: {}", relays.join(" "));
Ok(())
}
pub fn registry_show(dir: &Path) -> Result<()> {
let (_, signed) = open_registry(dir)?;
println!("ma_key {}", signed.doc.ma_key);
println!("version {}", signed.doc.version);
println!("issued_at {}", signed.doc.issued_at);
for relay in &signed.doc.relays {
println!("relay {relay}");
}
Ok(())
}
+33 -2
View File
@@ -70,6 +70,16 @@ enum Command {
listen: String,
#[arg(long, default_value_t = relay::DEFAULT_CAPACITY)]
capacity: usize,
#[arg(long = "peer")]
peers: Vec<String>,
#[arg(long)]
url: Option<String>,
#[arg(long)]
registry: Option<String>,
#[arg(long)]
ma_key: Option<String>,
#[arg(long, default_value_t = 3)]
max_hops: usize,
},
Status,
Member {
@@ -147,6 +157,10 @@ enum RegistryCommand {
id: String,
},
List,
SetRelays {
#[arg(required = true)]
relays: Vec<String>,
},
Show,
Serve {
#[arg(long, default_value = "127.0.0.1:7800")]
@@ -231,10 +245,26 @@ async fn main() -> Result<()> {
println!("control API POST http://{}/v1/local/query", handle.addr);
tokio::signal::ctrl_c().await?;
}
Command::Relay { listen, capacity } => {
Command::Relay {
listen,
capacity,
peers,
url,
registry,
ma_key,
max_hops,
} => {
let options = relay::RelayOptions {
capacity,
peers,
url,
registry,
ma_key,
max_hops,
};
let (listener, addr) = relay::bind(&listen).await?;
println!("relay listening on http://{addr} (capacity {capacity})");
relay::run(listener, capacity).await?;
relay::run(listener, options).await?;
}
Command::Status => commands::status(&cli.config).await?,
Command::Member { command } => match command {
@@ -271,6 +301,7 @@ async fn main() -> Result<()> {
}
RegistryCommand::Remove { id } => commands::registry_remove(&dir, &id)?,
RegistryCommand::List => commands::registry_list(&dir)?,
RegistryCommand::SetRelays { relays } => commands::registry_set_relays(&dir, &relays)?,
RegistryCommand::Show => commands::registry_show(&dir)?,
RegistryCommand::Serve { listen } => commands::registry_serve(&dir, &listen).await?,
},
+41 -115
View File
@@ -23,7 +23,7 @@ use crate::message::{
AggregateBody, Envelope, QueryBody, ResponseBody, ResponseItem, TYPE_AGGREGATE, TYPE_QUERY,
TYPE_RESPONSE, build_response, timestamp_is_fresh,
};
use crate::registry::{SignedRegistry, authorized_keys, load_registry, verify_registry};
use crate::registry::Watcher;
#[derive(Default)]
struct Aggregates {
@@ -37,8 +37,7 @@ pub struct Node {
index: LocalIndex,
members: RwLock<Vec<Member>>,
members_mtime: Mutex<Option<SystemTime>>,
registry: Mutex<Option<SignedRegistry>>,
registry_mtime: Mutex<Option<SystemTime>>,
registry: Option<Arc<Watcher>>,
aggregates: Mutex<Aggregates>,
seen: Mutex<HashSet<String>>,
pending: Mutex<HashMap<String, Vec<(String, ResponseBody)>>>,
@@ -122,14 +121,28 @@ impl Node {
.timeout(Duration::from_secs(15))
.build()
.context("building http client")?;
let node = Arc::new(Self {
let registry = match (
config.node.registry.as_deref(),
config.node.ma_key.as_deref(),
) {
(Some(source), Some(ma_key)) => Some(Watcher::new(
source,
ma_key,
Some(config.registry_cache_path()),
)),
(Some(_), None) => return Err(anyhow!("registry configured without ma_key")),
_ => None,
};
if let Some(watcher) = &registry {
watcher.load_initial();
}
Ok(Arc::new(Self {
config,
key,
index,
members: RwLock::new(members),
members_mtime: Mutex::new(members_mtime),
registry: Mutex::new(None),
registry_mtime: Mutex::new(None),
registry,
aggregates: Mutex::new(Aggregates::default()),
seen: Mutex::new(HashSet::new()),
pending: Mutex::new(HashMap::new()),
@@ -137,90 +150,7 @@ impl Node {
client,
sent: AtomicU64::new(0),
received: AtomicU64::new(0),
});
node.load_initial_registry();
Ok(node)
}
fn load_initial_registry(&self) {
let Some(source) = &self.config.node.registry else {
return;
};
let result = if source.starts_with("http://") || source.starts_with("https://") {
load_registry(&self.config.registry_cache_path())
} else {
let path = std::path::PathBuf::from(source);
let loaded = load_registry(&path);
if loaded.is_ok() {
let mtime = fs::metadata(&path)
.and_then(|metadata| metadata.modified())
.ok();
*self.registry_mtime.lock().expect("registry mtime lock") = mtime;
}
loaded
};
if let Ok(signed) = result {
if let Err(error) = self.verify_and_apply_registry(signed) {
eprintln!("cached registry rejected: {error}");
}
}
}
fn verify_and_apply_registry(&self, signed: SignedRegistry) -> Result<()> {
let ma_key = self
.config
.node
.ma_key
.as_deref()
.ok_or_else(|| anyhow!("registry configured without ma_key"))?;
verify_registry(&signed, ma_key)?;
{
let current = self.registry.lock().expect("registry lock");
if let Some(existing) = current.as_ref() {
if signed.doc.version <= existing.doc.version {
return Err(anyhow!("registry version rollback rejected"));
}
}
}
*self.registry.lock().expect("registry lock") = Some(signed);
Ok(())
}
fn refresh_registry(&self) {
let Some(source) = &self.config.node.registry else {
return;
};
if source.starts_with("http://") || source.starts_with("https://") {
return;
}
let path = std::path::PathBuf::from(source);
let mtime = fs::metadata(&path)
.and_then(|metadata| metadata.modified())
.ok();
{
let last = self.registry_mtime.lock().expect("registry mtime lock");
if *last == mtime {
return;
}
}
let Ok(signed) = load_registry(&path) else {
return;
};
*self.registry_mtime.lock().expect("registry mtime lock") = mtime;
if let Err(error) = self.verify_and_apply_registry(signed) {
eprintln!("registry update rejected: {error}");
}
}
async fn fetch_registry(&self) -> Result<()> {
let Some(url) = &self.config.node.registry else {
return Ok(());
};
let response = self.client.get(url).send().await?;
let signed: SignedRegistry = response.json().await?;
self.verify_and_apply_registry(signed.clone())?;
crate::registry::save_registry(&self.config.registry_cache_path(), &signed)?;
Ok(())
}))
}
fn refresh_members(&self) {
@@ -311,23 +241,31 @@ impl Node {
.with_context(|| format!("binding {}", node.config.node.listen))?;
let addr = listener.local_addr()?;
let mut tasks = Vec::new();
if let Some(source) = &node.config.node.registry {
if source.starts_with("http://") || source.starts_with("https://") {
if let Err(error) = node.fetch_registry().await {
if let Some(watcher) = node.registry.clone() {
if watcher.is_url() {
if let Err(error) = watcher.fetch(&node.client).await {
eprintln!("initial registry fetch failed: {error}");
}
let node_for_registry = node.clone();
let client = node.client.clone();
tasks.push(tokio::spawn(async move {
loop {
if let Err(error) = node_for_registry.fetch_registry().await {
tokio::time::sleep(Duration::from_secs(60)).await;
if let Err(error) = watcher.fetch(&client).await {
eprintln!("registry fetch failed: {error}");
}
tokio::time::sleep(Duration::from_secs(60)).await;
}
}));
}
}
for relay in node.config.node.relays.clone() {
let relays = if node.config.node.relays.is_empty() {
node.registry
.as_ref()
.map(|watcher| watcher.relays())
.unwrap_or_default()
} else {
node.config.node.relays.clone()
};
for relay in relays {
tasks.push(tokio::spawn(poll_relay(node.clone(), relay)));
}
let app = router(node.clone());
@@ -456,18 +394,11 @@ impl Node {
return;
}
self.refresh_members();
self.refresh_registry();
let (class, listed) = if self.config.node.registry.is_some() {
let registry = self.registry.lock().expect("registry lock");
let authorized = registry
.as_ref()
.map(|signed| authorized_keys(signed, now_ts()));
match authorized {
Some(map) => match map.get(&envelope.key) {
Some((id, class)) if id == &envelope.from => (Some(class.clone()), true),
_ => (None, false),
},
None => (None, false),
let (class, listed) = if let Some(watcher) = &self.registry {
watcher.refresh_if_changed();
match watcher.authorized(&envelope.key, now_ts()) {
Some((id, class)) if id == envelope.from => (Some(class), true),
_ => (None, false),
}
} else {
let members = self.members.read().expect("members lock");
@@ -761,12 +692,7 @@ async fn local_query(
async fn local_status(State(node): State<Arc<Node>>) -> Response {
let aggregates = node.aggregate_for(&current_period(), None);
let registry_version = node
.registry
.lock()
.expect("registry lock")
.as_ref()
.map(|signed| signed.doc.version);
let registry_version = node.registry.as_ref().and_then(|watcher| watcher.version());
Json(json!({
"name": node.config.node.name,
"id": node.config.node.id,
+127 -1
View File
@@ -1,6 +1,8 @@
use std::collections::HashMap;
use std::fs;
use std::path::Path;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use std::time::SystemTime;
use anyhow::{Context, Result, anyhow};
use serde::{Deserialize, Serialize};
@@ -35,6 +37,8 @@ pub struct RegistryDoc {
pub ma_key: String,
#[serde(default)]
pub members: Vec<RegistryMember>,
#[serde(default)]
pub relays: Vec<String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
@@ -75,6 +79,127 @@ pub fn save_registry(path: &Path, signed: &SignedRegistry) -> Result<()> {
Ok(())
}
pub struct Watcher {
source: String,
ma_key: String,
cache_path: Option<PathBuf>,
current: Mutex<Option<SignedRegistry>>,
mtime: Mutex<Option<SystemTime>>,
}
impl Watcher {
pub fn new(source: &str, ma_key: &str, cache_path: Option<PathBuf>) -> Arc<Self> {
Arc::new(Self {
source: source.to_string(),
ma_key: ma_key.to_string(),
cache_path,
current: Mutex::new(None),
mtime: Mutex::new(None),
})
}
pub fn is_url(&self) -> bool {
self.source.starts_with("http://") || self.source.starts_with("https://")
}
pub fn load_initial(&self) {
let loaded = if self.is_url() {
match &self.cache_path {
Some(cache) => load_registry(cache),
None => Err(anyhow!("no cache path for URL registry")),
}
} else {
let path = PathBuf::from(&self.source);
let loaded = load_registry(&path);
if loaded.is_ok() {
*self.mtime.lock().expect("registry mtime lock") = fs::metadata(&path)
.and_then(|metadata| metadata.modified())
.ok();
}
loaded
};
if let Ok(signed) = loaded {
if let Err(error) = self.apply(signed) {
eprintln!("cached registry rejected: {error}");
}
}
}
pub fn refresh_if_changed(&self) {
if self.is_url() {
return;
}
let path = PathBuf::from(&self.source);
let mtime = fs::metadata(&path)
.and_then(|metadata| metadata.modified())
.ok();
{
let last = self.mtime.lock().expect("registry mtime lock");
if *last == mtime {
return;
}
}
let Ok(signed) = load_registry(&path) else {
return;
};
*self.mtime.lock().expect("registry mtime lock") = mtime;
if let Err(error) = self.apply(signed) {
eprintln!("registry update rejected: {error}");
}
}
pub async fn fetch(&self, client: &reqwest::Client) -> Result<()> {
if !self.is_url() {
return Ok(());
}
let response = client.get(&self.source).send().await?;
let signed: SignedRegistry = response.json().await?;
self.apply(signed.clone())?;
if let Some(cache) = &self.cache_path {
save_registry(cache, &signed)?;
}
Ok(())
}
fn apply(&self, signed: SignedRegistry) -> Result<()> {
verify_registry(&signed, &self.ma_key)?;
{
let current = self.current.lock().expect("registry lock");
if let Some(existing) = current.as_ref() {
if signed.doc.version <= existing.doc.version {
return Err(anyhow!("registry version rollback rejected"));
}
}
}
*self.current.lock().expect("registry lock") = Some(signed);
Ok(())
}
pub fn authorized(&self, key: &str, now: u64) -> Option<(String, String)> {
let current = self.current.lock().expect("registry lock");
current
.as_ref()
.and_then(|signed| authorized_keys(signed, now).remove(key))
}
pub fn version(&self) -> Option<u64> {
self.current
.lock()
.expect("registry lock")
.as_ref()
.map(|signed| signed.doc.version)
}
pub fn relays(&self) -> Vec<String> {
self.current
.lock()
.expect("registry lock")
.as_ref()
.map(|signed| signed.doc.relays.clone())
.unwrap_or_default()
}
}
pub fn authorized_keys(signed: &SignedRegistry, now: u64) -> HashMap<String, (String, String)> {
let mut authorized = HashMap::new();
for member in &signed.doc.members {
@@ -117,6 +242,7 @@ mod tests {
issued_at: now_ts(),
ma_key: String::new(),
members,
relays: Vec::new(),
},
key,
)
+229 -38
View File
@@ -3,7 +3,7 @@ use std::net::SocketAddr;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use anyhow::Result;
use anyhow::{Result, anyhow};
use axum::extract::{Query, State};
use axum::http::{StatusCode, header};
use axum::response::{IntoResponse, Response};
@@ -15,13 +15,40 @@ use tokio::net::TcpListener;
use crate::crypto::{now_ts, poll_signing_bytes, random_nonce, verify_signature};
use crate::message::{Envelope, TYPE_AGGREGATE, TYPE_QUERY, TYPE_RESPONSE, timestamp_is_fresh};
use crate::registry::Watcher;
pub const DEFAULT_CAPACITY: usize = 256;
const CHALLENGE_TTL: Duration = Duration::from_secs(120);
const SEEN_TTL: Duration = Duration::from_secs(30);
#[derive(Clone, Debug)]
pub struct RelayOptions {
pub capacity: usize,
pub peers: Vec<String>,
pub url: Option<String>,
pub registry: Option<String>,
pub ma_key: Option<String>,
pub max_hops: usize,
}
impl Default for RelayOptions {
fn default() -> Self {
Self {
capacity: DEFAULT_CAPACITY,
peers: Vec::new(),
url: None,
registry: None,
ma_key: None,
max_hops: 3,
}
}
}
#[derive(Default)]
struct MemberQueue {
items: VecDeque<(u64, Envelope)>,
lagging: bool,
missed: u64,
}
struct Challenge {
@@ -32,57 +59,136 @@ struct Challenge {
struct Inner {
members: HashMap<String, MemberQueue>,
challenges: HashMap<String, Challenge>,
seen: HashMap<String, Instant>,
seq: u64,
}
pub struct Relay {
inner: Mutex<Inner>,
capacity: usize,
options: RelayOptions,
registry: Option<Arc<Watcher>>,
client: reqwest::Client,
}
impl Relay {
pub fn new(capacity: usize) -> Arc<Self> {
Arc::new(Self {
pub fn new(options: RelayOptions) -> Result<Arc<Self>> {
let registry = match (options.registry.as_deref(), options.ma_key.as_deref()) {
(Some(source), Some(ma_key)) => Some(Watcher::new(source, ma_key, None)),
(Some(_), None) => return Err(anyhow!("registry configured without ma_key")),
_ => None,
};
if !options.peers.is_empty() && options.url.is_none() {
return Err(anyhow!("--url is required when --peer is set"));
}
let relay = Arc::new(Self {
inner: Mutex::new(Inner {
members: HashMap::new(),
challenges: HashMap::new(),
seen: HashMap::new(),
seq: 0,
}),
capacity,
})
options,
registry,
client: reqwest::Client::builder()
.timeout(Duration::from_secs(5))
.build()?,
});
if let Some(watcher) = &relay.registry {
watcher.load_initial();
}
Ok(relay)
}
fn push(&self, targets: Option<&str>, envelope: Envelope) -> Result<usize, StatusCode> {
let mut inner = self.inner.lock().expect("relay lock");
if inner
.members
.values()
.any(|q| q.items.len() >= self.capacity)
{
return Err(StatusCode::TOO_MANY_REQUESTS);
fn admit(&self, envelope: &Envelope) -> bool {
match &self.registry {
Some(watcher) => {
watcher.refresh_if_changed();
watcher.authorized(&envelope.key, now_ts()).is_some()
}
None => true,
}
}
fn mark_seen(&self, signature: &str) -> bool {
let mut inner = self.inner.lock().expect("relay lock");
let now = Instant::now();
inner
.seen
.retain(|_, at| now.duration_since(*at) < SEEN_TTL);
if inner.seen.contains_key(signature) {
return false;
}
inner.seen.insert(signature.to_string(), now);
true
}
fn push_local(&self, target: Option<&str>, envelope: Envelope) -> Result<usize, StatusCode> {
let mut inner = self.inner.lock().expect("relay lock");
inner.seq += 1;
let seq = inner.seq;
match targets {
match target {
Some(member) => {
let Some(queue) = inner.members.get_mut(member) else {
return Err(StatusCode::NOT_FOUND);
};
if queue.items.len() >= self.options.capacity {
queue.lagging = true;
queue.missed += 1;
return Err(StatusCode::TOO_MANY_REQUESTS);
}
queue.items.push_back((seq, envelope));
Ok(1)
}
None => {
let from = envelope.key.clone();
inner.members.entry(from).or_default();
let publisher = envelope.key.clone();
inner.members.entry(publisher).or_default();
let mut delivered = 0;
for queue in inner.members.values_mut() {
queue.items.push_back((seq, envelope.clone()));
delivered += 1;
if queue.items.len() >= self.options.capacity {
queue.lagging = true;
queue.missed += 1;
} else {
queue.items.push_back((seq, envelope.clone()));
delivered += 1;
}
}
Ok(delivered)
}
}
}
fn forward(&self, envelope: Envelope, origin: Option<&str>, hops: usize) {
if self.options.peers.is_empty() {
return;
}
let Some(url) = self.options.url.clone() else {
return;
};
let peers: Vec<String> = self
.options
.peers
.iter()
.filter(|peer| Some(peer.as_str()) != origin)
.cloned()
.collect();
if peers.is_empty() {
return;
}
let client = self.client.clone();
tokio::spawn(async move {
for peer in peers {
let target = format!(
"{}/v1/federation?origin={}&hops={}",
peer.trim_end_matches('/'),
url,
hops
);
if let Err(error) = client.post(&target).json(&envelope).send().await {
eprintln!("federation forward to {peer} failed: {error}");
}
}
});
}
}
pub fn router(relay: Arc<Relay>) -> Router {
@@ -90,13 +196,30 @@ pub fn router(relay: Arc<Relay>) -> Router {
.route("/health", get(health))
.route("/v1/challenge", get(challenge))
.route("/v1/publish", post(publish))
.route("/v1/federation", post(federation))
.route("/v1/unicast", post(unicast))
.route("/v1/poll", get(poll))
.with_state(relay)
}
pub async fn run(listener: TcpListener, capacity: usize) -> Result<()> {
let relay = Relay::new(capacity);
pub async fn run(listener: TcpListener, options: RelayOptions) -> Result<()> {
let relay = Relay::new(options)?;
if let Some(watcher) = relay.registry.clone() {
if watcher.is_url() {
let client = relay.client.clone();
if let Err(error) = watcher.fetch(&client).await {
eprintln!("initial registry fetch failed: {error}");
}
tokio::spawn(async move {
loop {
tokio::time::sleep(Duration::from_secs(60)).await;
if let Err(error) = watcher.fetch(&client).await {
eprintln!("registry fetch failed: {error}");
}
}
});
}
}
axum::serve(listener, router(relay)).await?;
Ok(())
}
@@ -152,14 +275,65 @@ async fn publish(State(relay): State<Arc<Relay>>, Json(envelope): Json<Envelope>
if envelope.msg_type != TYPE_QUERY {
return bad_request("relay carries broadcast queries only");
}
match relay.push(None, envelope) {
Ok(delivered) => (
StatusCode::ACCEPTED,
Json(json!({ "delivered": delivered })),
)
.into_response(),
Err(status) => backpressure(status),
if !relay.admit(&envelope) {
return bad_request("sender not admitted");
}
relay.mark_seen(&envelope.sig);
let delivered = relay.push_local(None, envelope.clone()).unwrap_or(0);
relay.forward(envelope, None, 1);
(
StatusCode::ACCEPTED,
Json(json!({ "delivered": delivered })),
)
.into_response()
}
#[derive(Deserialize)]
struct FederationParams {
origin: Option<String>,
hops: Option<usize>,
}
async fn federation(
State(relay): State<Arc<Relay>>,
Query(params): Query<FederationParams>,
Json(envelope): Json<Envelope>,
) -> Response {
let Some(origin) = params.origin.clone() else {
return bad_request("origin required");
};
if !relay.options.peers.iter().any(|peer| peer == &origin) {
return bad_request("unknown peer origin");
}
if envelope.verify().is_err() {
return bad_request("invalid signature");
}
if !timestamp_is_fresh(envelope.ts, now_ts()) {
return bad_request("stale timestamp");
}
if envelope.msg_type != TYPE_QUERY {
return bad_request("relay carries broadcast queries only");
}
if !relay.admit(&envelope) {
return bad_request("sender not admitted");
}
if !relay.mark_seen(&envelope.sig) {
return (
StatusCode::ACCEPTED,
Json(json!({ "delivered": 0, "duplicate": true })),
)
.into_response();
}
let delivered = relay.push_local(None, envelope.clone()).unwrap_or(0);
let hops = params.hops.unwrap_or(1);
if hops < relay.options.max_hops {
relay.forward(envelope, Some(&origin), hops + 1);
}
(
StatusCode::ACCEPTED,
Json(json!({ "delivered": delivered })),
)
.into_response()
}
#[derive(Deserialize)]
@@ -181,14 +355,17 @@ async fn unicast(
if envelope.msg_type != TYPE_RESPONSE && envelope.msg_type != TYPE_AGGREGATE {
return bad_request("unicast carries responses and aggregates only");
}
match relay.push(Some(&params.to), envelope) {
if !relay.admit(&envelope) {
return bad_request("sender not admitted");
}
match relay.push_local(Some(&params.to), envelope) {
Ok(_) => (StatusCode::OK, Json(json!({ "delivered": true }))).into_response(),
Err(StatusCode::NOT_FOUND) => (
StatusCode::NOT_FOUND,
Json(json!({ "error": "member not connected" })),
)
.into_response(),
Err(status) => backpressure(status),
Err(status) => backpressure(status, 0),
}
}
@@ -225,15 +402,29 @@ async fn poll(State(relay): State<Arc<Relay>>, Query(params): Query<PollParams>)
let timeout = Duration::from_millis(params.timeout_ms.unwrap_or(25_000).min(60_000));
let deadline = Instant::now() + timeout;
loop {
let batch: Vec<Envelope> = {
let (batch, lagged): (Vec<Envelope>, Option<u64>) = {
let mut inner = relay.inner.lock().expect("relay lock");
let queue = inner.members.entry(params.member.clone()).or_default();
queue
.items
.drain(..)
.map(|(_, envelope)| envelope)
.collect()
if queue.lagging {
let missed = queue.missed;
queue.lagging = false;
queue.missed = 0;
queue.items.clear();
(Vec::new(), Some(missed))
} else {
(
queue
.items
.drain(..)
.map(|(_, envelope)| envelope)
.collect(),
None,
)
}
};
if let Some(missed) = lagged {
return backpressure(StatusCode::TOO_MANY_REQUESTS, missed);
}
if !batch.is_empty() {
return (StatusCode::OK, Json(json!({ "messages": batch }))).into_response();
}
@@ -252,11 +443,11 @@ fn unauthorized(message: &str) -> Response {
(StatusCode::UNAUTHORIZED, Json(json!({ "error": message }))).into_response()
}
fn backpressure(status: StatusCode) -> Response {
fn backpressure(status: StatusCode, missed: u64) -> Response {
(
status,
[(header::RETRY_AFTER, "1")],
Json(json!({ "error": "transport backpressure" })),
Json(json!({ "error": "lagging", "missed": missed })),
)
.into_response()
}