pub mod codec;
pub mod compress;
pub mod health;
pub mod status;
use crate::courierust_body::Body;
use crate::courierust_bytes::Bytes;
use crate::courierust_client::Client;
use crate::courierust_error::{Error, Result};
use crate::courierust_h2::priority::Priority;
use crate::courierust_http::header::{HeaderMap, HeaderName, HeaderValue};
use crate::courierust_http::method::Method;
use crate::courierust_http::request::Request;
use crate::courierust_http::response::Response;
use crate::courierust_http::uri::PathAndQuery;
use crate::courierust_server::{Handler, Server, ServerConfig};
use std::net::ToSocketAddrs;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Duration;
pub const CONTENT_TYPE: &str = "application/grpc";
pub const TE_VALUE: &str = "trailers";
pub const DEFAULT_MAX_MESSAGE_SIZE: usize = 4 * 1024 * 1024;
pub fn grpc_timeout(d: Duration) -> String {
if d.is_zero() {
return "1n".to_string();
}
let nanos = d.as_nanos();
let hours = d.as_secs() / 3600;
if hours > 0 {
return format!("{hours}H");
}
let mins = d.as_secs() / 60;
if mins > 0 {
return format!("{mins}M");
}
let secs = d.as_secs();
if secs > 0 {
return format!("{secs}S");
}
let millis = nanos / 1_000_000;
if millis > 0 {
return format!("{millis}m");
}
let micros = nanos / 1_000;
if micros > 0 {
return format!("{micros}u");
}
format!("{nanos}n")
}
pub trait Interceptor: Send + Sync + 'static {
fn intercept(&self, method: &str, headers: &mut HeaderMap);
}
impl<F> Interceptor for F
where
F: Fn(&str, &mut HeaderMap) + Send + Sync + 'static,
{
fn intercept(&self, method: &str, headers: &mut HeaderMap) {
self(method, headers)
}
}
pub struct GrpcClientConfig {
pub base: String,
pub max_message_size: usize,
pub interceptor: Option<Arc<dyn Interceptor>>,
pub timeout: Option<Duration>,
pub compress: bool,
pub http_client: Client,
}
#[derive(Clone)]
pub struct GrpcClient {
pub client: Client,
base: String,
addrs: Arc<Vec<String>>,
cursor: Arc<AtomicUsize>,
max_message_size: usize,
compress: bool,
interceptor: Option<Arc<dyn Interceptor>>,
timeout: Option<Duration>,
}
impl GrpcClient {
pub fn new(base: &str) -> Result<Self> {
Self::with_config(GrpcClientConfig {
base: base.to_string(),
max_message_size: DEFAULT_MAX_MESSAGE_SIZE,
interceptor: None,
timeout: None,
compress: false,
http_client: Client::with_config(crate::courierust_client::ClientConfig {
http2: true,
user_agent: None,
..Default::default()
}),
})
}
pub fn new_with_client(base: &str, client: Client) -> Result<Self> {
Self::with_config(GrpcClientConfig {
base: base.to_string(),
max_message_size: DEFAULT_MAX_MESSAGE_SIZE,
interceptor: None,
timeout: None,
compress: false,
http_client: client,
})
}
pub fn with_config(config: GrpcClientConfig) -> Result<Self> {
let base = config.base.trim_end_matches('/').to_string();
let (base, addrs) = if let Some(a) = resolve_dns(&base)? {
(String::new(), Arc::new(a))
} else {
(base, Arc::new(Vec::new()))
};
Ok(Self {
client: config.http_client,
base,
addrs,
cursor: Arc::new(AtomicUsize::new(0)),
max_message_size: config.max_message_size,
compress: config.compress,
interceptor: config.interceptor,
timeout: config.timeout,
})
}
pub fn max_message_size(&self) -> usize {
self.max_message_size
}
fn effective_base(&self) -> String {
if self.addrs.is_empty() {
self.base.clone()
} else {
let i = self.cursor.fetch_add(1, Ordering::Relaxed) % self.addrs.len();
format!("http://{}", self.addrs[i])
}
}
pub fn call(&self, method: &str, req: Bytes) -> Result<Bytes> {
let mut stream = self.call_with_metadata(method, req, &HeaderMap::new())?;
match stream.next_message()? {
Some(msg) => Ok(msg),
None => Err(Error::grpc(
status::INTERNAL,
"empty response from unary call",
)),
}
}
pub fn call_unary<Req: codec::EncodeMessage, Resp: codec::DecodeMessage>(
&self,
method: &str,
req: &Req,
) -> Result<Resp> {
let body = req.encode_message()?;
let resp = self.call(method, Bytes::from(body))?;
Resp::decode_message(&resp)
}
pub fn call_with_metadata(
&self,
method: &str,
req: Bytes,
metadata: &HeaderMap,
) -> Result<MessageStream> {
let msg = frame_outbound(req, self.compress, self.max_message_size)?;
let request = self.build_request(method, metadata, Body::Bytes(Bytes::from(msg)))?;
let url = crate::courierust_http::uri::Url::parse(&format!(
"{}{}",
self.effective_base(),
method
))?;
let resp = self
.client
.execute_h2_raw(&url, request, Priority::default())?;
MessageStream::new(resp, self.max_message_size)
}
pub fn call_stream(&self, method: &str, req: Bytes) -> Result<MessageStream> {
self.call_with_metadata(method, req, &HeaderMap::new())
}
pub fn client_stream(
&self,
method: &str,
reqs: std::sync::mpsc::Receiver<Result<Bytes>>,
) -> Result<Bytes> {
let mut stream = self.bidi_stream(method, reqs)?;
match stream.next_message()? {
Some(msg) => Ok(msg),
None => Err(Error::grpc(
status::INTERNAL,
"empty response from client-streaming call",
)),
}
}
pub fn bidi_stream(
&self,
method: &str,
reqs: std::sync::mpsc::Receiver<Result<Bytes>>,
) -> Result<MessageStream> {
let body = frame_stream(reqs, self.max_message_size, self.compress)?;
let request = self.build_request(method, &HeaderMap::new(), body)?;
let url = crate::courierust_http::uri::Url::parse(&format!(
"{}{}",
self.effective_base(),
method
))?;
let resp = self
.client
.execute_h2_stream(&url, request, Priority::default())?;
MessageStream::new(resp, self.max_message_size)
}
fn build_request(
&self,
method: &str,
metadata: &HeaderMap,
body: Body,
) -> Result<Request<Body>> {
let uri = PathAndQuery::from_bytes(method.as_bytes())?;
let mut request = Request::new(Method::POST, uri);
request.headers.insert(
HeaderName::from_lowercase("content-type"),
HeaderValue::from_static(CONTENT_TYPE),
);
request.headers.insert(
HeaderName::from_lowercase("te"),
HeaderValue::from_static(TE_VALUE),
);
request.headers.insert(
HeaderName::from_lowercase("grpc-encoding"),
if self.compress {
HeaderValue::from_static("gzip")
} else {
HeaderValue::from_static("identity")
},
);
request.headers.insert(
HeaderName::from_lowercase("grpc-accept-encoding"),
HeaderValue::from_static("gzip, identity"),
);
if let Some(t) = self.timeout {
let v = grpc_timeout(t);
request.headers.insert(
HeaderName::from_lowercase("grpc-timeout"),
HeaderValue::from_bytes(v.as_bytes())?,
);
}
for (n, v) in metadata.iter() {
request.headers.append(n.clone(), v.clone());
}
if let Some(interceptor) = &self.interceptor {
interceptor.intercept(method, &mut request.headers);
}
request.body = body;
Ok(request)
}
}
fn resolve_dns(base: &str) -> Result<Option<Vec<String>>> {
let rest = base
.strip_prefix("dns:///")
.or_else(|| base.strip_prefix("dns://"))
.map(|s| s.trim_start_matches('/'));
let Some(rest) = rest else {
return Ok(None);
};
let (host, port) = match rest.rsplit_once(':') {
Some((h, p)) => (
h.to_string(),
p.parse::<u16>()
.map_err(|_| Error::grpc(status::INVALID_ARGUMENT, "invalid dns:/// port"))?,
),
None => (rest.to_string(), 80u16),
};
let mut seen = std::collections::HashSet::new();
let mut addrs = Vec::new();
for a in (host.as_str(), port)
.to_socket_addrs()
.map_err(|e| Error::grpc(status::UNAVAILABLE, format!("dns resolve: {e}")))?
{
if seen.insert(a) {
addrs.push(a.to_string());
}
}
if addrs.is_empty() {
return Err(Error::grpc(
status::UNAVAILABLE,
format!("no addresses for {rest}"),
));
}
Ok(Some(addrs))
}
pub struct MessageStream {
body: Body,
buf: Vec<u8>,
done: bool,
trailers: Arc<Mutex<Option<HeaderMap>>>,
head_headers: HeaderMap,
max_message_size: usize,
}
impl MessageStream {
fn new(
resp: crate::courierust_client::h2::H2Response,
max_message_size: usize,
) -> Result<Self> {
let ct = resp
.head
.headers
.get("content-type")
.map(|v| v.to_str().unwrap_or(""))
.unwrap_or("");
if !ct.starts_with(CONTENT_TYPE) {
return Err(Error::grpc(
status::INTERNAL,
format!("unexpected content-type: {ct}"),
));
}
Ok(Self {
body: resp.body,
buf: Vec::new(),
done: false,
trailers: resp.trailers,
head_headers: resp.head.headers,
max_message_size,
})
}
pub fn response_headers(&self) -> &HeaderMap {
&self.head_headers
}
pub fn trailers(&self) -> Option<HeaderMap> {
self.trailers.lock().unwrap().clone()
}
pub fn next_message(&mut self) -> Result<Option<Bytes>> {
loop {
if self.buf.len() >= 5 {
let (compressed, len) = read_frame_header(&self.buf, self.max_message_size)?;
if self.buf.len() >= 5 + len {
let payload = self.buf[5..5 + len].to_vec();
self.buf.drain(..5 + len);
return self.decode_payload(compressed, payload);
}
}
if self.done {
if !self.buf.is_empty() {
return Err(Error::protocol("truncated gRPC message"));
}
self.finish()?;
return Ok(None);
}
match &self.body {
Body::Channel(rx) => match rx.recv() {
Ok(Ok(chunk)) => self.buf.extend_from_slice(&chunk),
Ok(Err(e)) => return Err(e),
Err(_) => self.done = true,
},
Body::Empty => self.done = true,
Body::Bytes(b) => {
self.buf.extend_from_slice(b);
self.body = Body::Empty;
self.done = true;
}
}
}
}
fn decode_payload(&self, compressed: bool, payload: Vec<u8>) -> Result<Option<Bytes>> {
if compressed {
let raw = compress::gunzip(&payload, self.max_message_size)?;
return Ok(Some(Bytes::from(raw)));
}
Ok(Some(Bytes::from(payload)))
}
fn finish(&self) -> Result<()> {
let code = self
.trailers
.lock()
.unwrap()
.as_ref()
.and_then(|t| t.get("grpc-status"))
.and_then(|v| v.to_str().ok())
.and_then(|s| s.parse::<u32>().ok())
.or_else(|| {
self.head_headers
.get("grpc-status")
.and_then(|v| v.to_str().ok())
.and_then(|s| s.parse::<u32>().ok())
})
.unwrap_or(status::OK);
if code != status::OK {
let msg = self
.trailers
.lock()
.unwrap()
.as_ref()
.and_then(|t| t.get("grpc-message"))
.or_else(|| self.head_headers.get("grpc-message"))
.map(|v| percent_decode(v.to_str().unwrap_or("")))
.unwrap_or_default();
return Err(Error::grpc(code, msg));
}
Ok(())
}
}
fn read_frame_header(buf: &[u8], max: usize) -> Result<(bool, usize)> {
if buf.len() < 5 {
return Err(Error::protocol("truncated gRPC frame header"));
}
let compressed = buf[0] != 0;
let len = u32::from_be_bytes([buf[1], buf[2], buf[3], buf[4]]) as usize;
if len > max {
return Err(Error::overflow("gRPC message too large"));
}
Ok((compressed, len))
}
pub fn frame_message(payload: &[u8], compressed: bool) -> Vec<u8> {
let mut out = Vec::with_capacity(payload.len() + 5);
out.push(if compressed { 1 } else { 0 });
out.extend_from_slice(&(payload.len() as u32).to_be_bytes());
out.extend_from_slice(payload);
out
}
fn frame_outbound(msg: Bytes, compress: bool, max_message_size: usize) -> Result<Vec<u8>> {
if msg.len() > max_message_size {
return Err(Error::overflow("gRPC message too large"));
}
if compress {
let gz = compress::gzip(&msg);
Ok(frame_message(&gz, true))
} else {
Ok(frame_message(&msg, false))
}
}
fn frame_stream(
raw: std::sync::mpsc::Receiver<Result<Bytes>>,
max_message_size: usize,
compress: bool,
) -> Result<Body> {
let (tx, rx) = std::sync::mpsc::channel();
std::thread::Builder::new()
.name("courierust-grpc-frame".into())
.spawn(move || {
while let Ok(msg) = raw.recv() {
let frame = match msg {
Ok(m) => frame_outbound(m, compress, max_message_size),
Err(e) => Err(e),
};
if tx.send(frame.map(Bytes::from)).is_err() {
break;
}
}
})
.map_err(|e| Error::io(e.to_string()))?;
Ok(Body::Channel(rx))
}
pub fn percent_decode(s: &str) -> String {
let b = s.as_bytes();
let mut out = Vec::with_capacity(b.len());
let mut i = 0;
while i < b.len() {
if b[i] == b'%' && i + 2 < b.len() {
let hex = &b[i + 1..i + 3];
if let Ok(Ok(v)) = std::str::from_utf8(hex).map(|h| u8::from_str_radix(h, 16)) {
out.push(v);
i += 3;
continue;
}
}
out.push(b[i]);
i += 1;
}
String::from_utf8_lossy(&out).into_owned()
}
pub trait Service: Send + Sync + 'static {
fn call(&self, method: &str, req: Bytes) -> Result<Bytes>;
}
impl<F> Service for F
where
F: Fn(&str, Bytes) -> Result<Bytes> + Send + Sync + 'static,
{
fn call(&self, method: &str, req: Bytes) -> Result<Bytes> {
self(method, req)
}
}
pub trait StreamingService: Send + Sync + 'static {
fn serve(
&self,
method: &str,
reqs: &mut dyn Iterator<Item = Result<Bytes>>,
tx: &crate::courierust_body::BodySender,
) -> Result<()>;
}
impl<F> StreamingService for F
where
F: Fn(
&str,
&mut dyn Iterator<Item = Result<Bytes>>,
&crate::courierust_body::BodySender,
) -> Result<()>
+ Send
+ Sync
+ 'static,
{
fn serve(
&self,
method: &str,
reqs: &mut dyn Iterator<Item = Result<Bytes>>,
tx: &crate::courierust_body::BodySender,
) -> Result<()> {
self(method, reqs, tx)
}
}
pub struct GrpcServer {
server: Server,
service: Arc<dyn StreamingService>,
max_message_size: usize,
}
impl GrpcServer {
pub fn bind(
addr: impl std::net::ToSocketAddrs,
service: impl Service,
) -> std::io::Result<Self> {
Self::bind_streaming(addr, UnaryAdapter(Arc::new(service)))
}
pub fn bind_streaming(
addr: impl std::net::ToSocketAddrs,
service: impl StreamingService,
) -> std::io::Result<Self> {
Self::bind_streaming_with_config(addr, service, ServerConfig::default())
}
pub fn bind_streaming_with_config(
addr: impl std::net::ToSocketAddrs,
service: impl StreamingService,
http_cfg: ServerConfig,
) -> std::io::Result<Self> {
let cfg = ServerConfig {
http2: true,
..http_cfg
};
let server = Server::bind_with_config(addr, cfg)?;
Ok(Self {
server,
service: Arc::new(service),
max_message_size: DEFAULT_MAX_MESSAGE_SIZE,
})
}
pub fn local_addr(&self) -> std::io::Result<std::net::SocketAddr> {
self.server.local_addr()
}
pub fn max_message_size(mut self, n: usize) -> Self {
self.max_message_size = n.max(1);
self
}
pub fn serve(self) -> std::io::Result<()> {
let service = self.service;
let server = self.server;
let max = self.max_message_size;
server.serve_with_config(GrpcHandler {
service,
max_message_size: max,
})
}
pub fn serve_background(self) -> std::io::Result<crate::courierust_server::ServerHandle> {
let service = self.service;
let server = self.server;
let max = self.max_message_size;
server.serve_background(GrpcHandler {
service,
max_message_size: max,
})
}
}
pub fn unary(service: impl Service) -> impl StreamingService {
UnaryAdapter(Arc::new(service))
}
struct UnaryAdapter(Arc<dyn Service>);
impl StreamingService for UnaryAdapter {
fn serve(
&self,
method: &str,
reqs: &mut dyn Iterator<Item = Result<Bytes>>,
tx: &crate::courierust_body::BodySender,
) -> Result<()> {
let req = reqs.next().transpose()?.unwrap_or_default();
let resp = self.0.call(method, req)?;
tx.send(resp)?;
Ok(())
}
}
fn decode_messages(raw: &[u8], max: usize) -> Result<Vec<Result<Bytes>>> {
let mut out = Vec::new();
let mut pos = 0usize;
while pos < raw.len() {
let (compressed, len) = read_frame_header(&raw[pos..], max)?;
let start = pos + 5;
if raw.len() < start + len {
return Err(Error::protocol("truncated gRPC request message"));
}
let payload = &raw[start..start + len];
if compressed {
let plain = compress::gunzip(payload, max)?;
out.push(Ok(Bytes::from(plain)));
} else {
out.push(Ok(Bytes::from(payload)));
}
pos = start + len;
}
Ok(out)
}
struct GrpcHandler {
service: Arc<dyn StreamingService>,
max_message_size: usize,
}
struct Deadline(Option<std::time::Instant>);
impl Deadline {
fn parse(headers: &HeaderMap) -> Result<Self> {
let Some(v) = headers.get("grpc-timeout") else {
return Ok(Self(None));
};
let s = v
.to_str()
.map_err(|_| Error::grpc(status::INVALID_ARGUMENT, "malformed grpc-timeout value"))?;
if s.len() < 2 {
return Err(Error::grpc(
status::INVALID_ARGUMENT,
"grpc-timeout must be a value plus a unit",
));
}
let (num, unit) = s.split_at(s.len() - 1);
let n: u64 = num
.parse()
.map_err(|_| Error::grpc(status::INVALID_ARGUMENT, "bad grpc-timeout value"))?;
if n > 99_999_999 {
return Err(Error::grpc(
status::INVALID_ARGUMENT,
"grpc-timeout value too large",
));
}
let dur = match unit {
"H" => Duration::from_secs(n),
"M" => Duration::from_secs(n.saturating_mul(60)),
"S" => Duration::from_secs(n),
"m" => Duration::from_millis(n),
"u" => Duration::from_micros(n),
"n" => Duration::from_nanos(n),
_ => {
return Err(Error::grpc(
status::INVALID_ARGUMENT,
"bad grpc-timeout unit",
));
}
};
Ok(Self(Some(std::time::Instant::now() + dur)))
}
fn expired(&self) -> bool {
match self.0 {
Some(d) => std::time::Instant::now() >= d,
None => false,
}
}
}
struct DeadlineIter<'a, I> {
inner: I,
deadline: &'a Deadline,
}
impl<'a, I: Iterator<Item = Result<Bytes>>> Iterator for DeadlineIter<'a, I> {
type Item = Result<Bytes>;
fn next(&mut self) -> Option<Self::Item> {
if self.deadline.expired() {
return Some(Err(Error::grpc(
status::DEADLINE_EXCEEDED,
"deadline exceeded",
)));
}
self.inner.next()
}
}
impl Handler for GrpcHandler {
fn handle(&self, req: Request<Body>) -> Response<Body> {
let method = req.uri.as_str().to_string();
let is_grpc = req
.headers
.get("content-type")
.map(|v| v.to_str().unwrap_or("").starts_with(CONTENT_TYPE))
.unwrap_or(false);
if !is_grpc {
let mut resp = Response::<Body>::with_status(415.into());
resp.headers.insert(
HeaderName::from_lowercase("content-type"),
HeaderValue::from_static("text/plain"),
);
resp.body = Body::Bytes(Bytes::from_static(b"content-type must be application/grpc"));
return resp;
}
let deadline = match Deadline::parse(&req.headers) {
Ok(d) => d,
Err(e) => {
return grpc_error_response(
e.grpc_code().unwrap_or(status::INTERNAL),
&e.to_string(),
);
}
};
if deadline.expired() {
return grpc_error_response(status::DEADLINE_EXCEEDED, "deadline exceeded");
}
let accept_gzip = req
.headers
.get("grpc-accept-encoding")
.map(|v| v.to_str().unwrap_or(""))
.map(|v| {
v.to_ascii_lowercase()
.split(',')
.any(|c| c.trim() == "gzip")
})
.unwrap_or(false);
let raw = match req.body.collect() {
Ok(b) => b.to_vec(),
Err(e) => return grpc_error_response(status::INTERNAL, &e.to_string()),
};
let messages = match decode_messages(&raw, self.max_message_size) {
Ok(m) => m,
Err(e) => {
return grpc_error_response(
e.grpc_code().unwrap_or(status::INTERNAL),
&e.to_string(),
)
}
};
let (msg_tx, msg_rx) = std::sync::mpsc::channel();
let msg_sender = crate::courierust_body::BodySender::from_sender(msg_tx);
let (body_tx, body) = crate::courierust_body::channel();
let max = self.max_message_size;
let (started_tx, started_rx) = std::sync::mpsc::channel::<()>();
let (result_tx, result_rx) =
std::sync::mpsc::channel::<std::result::Result<(), (u32, String)>>();
let started_tx2 = started_tx.clone();
let _ = std::thread::Builder::new()
.name("courierust-grpc-frame".into())
.spawn(move || {
let mut first = true;
while let Ok(m) = msg_rx.recv() {
if first {
let _ = started_tx2.send(());
first = false;
}
let frame = match m {
Ok(raw) => frame_outbound(raw, accept_gzip, max),
Err(e) => Err(e),
};
if body_tx.send_result(frame.map(Bytes::from)).is_err() {
break;
}
}
});
let service = self.service.clone();
let method2 = method.clone();
let started_tx3 = started_tx;
let deadline_instant = deadline.0;
let _ = std::thread::Builder::new()
.name("courierust-grpc-serve".into())
.spawn(move || {
let outcome = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let mut iter = DeadlineIter {
inner: messages.into_iter(),
deadline: &deadline,
};
let r = service.serve(&method2, &mut iter, &msg_sender);
r.map_err(|e| (e.grpc_code().unwrap_or(status::INTERNAL), e.to_string()))
}));
let mapped = match outcome {
Ok(r) => r,
Err(_) => Err((status::INTERNAL, "service handler panicked".to_string())),
};
let _ = result_tx.send(mapped);
drop(started_tx3);
});
let wait = match deadline_instant {
Some(d) => {
let now = std::time::Instant::now();
if d <= now {
return grpc_error_response(status::DEADLINE_EXCEEDED, "deadline exceeded");
}
match started_rx.recv_timeout(d - now) {
Ok(()) => WaitOutcome::Started,
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {
WaitOutcome::DeadlineExceeded
}
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => WaitOutcome::Finished,
}
}
None => {
if started_rx.recv().is_ok() {
WaitOutcome::Started
} else {
WaitOutcome::Finished
}
}
};
let streaming_response = |body: Body| {
let mut resp = Response::<Body>::with_status(200.into());
resp.headers.insert(
HeaderName::from_lowercase("content-type"),
HeaderValue::from_static(CONTENT_TYPE),
);
resp.headers.insert(
HeaderName::from_lowercase("grpc-encoding"),
if accept_gzip {
HeaderValue::from_static("gzip")
} else {
HeaderValue::from_static("identity")
},
);
let mut tr = HeaderMap::new();
tr.insert(
HeaderName::from_lowercase("grpc-status"),
HeaderValue::from_static("0"),
);
resp.trailers = Some(tr);
resp.body = body;
resp
};
match wait {
WaitOutcome::Started => streaming_response(body),
WaitOutcome::Finished => match result_rx.recv() {
Ok(Ok(())) => streaming_response(body),
Ok(Err((code, msg))) => grpc_error_response(code, &msg),
Err(_) => grpc_error_response(status::INTERNAL, "service thread exited"),
},
WaitOutcome::DeadlineExceeded => {
grpc_error_response(status::DEADLINE_EXCEEDED, "deadline exceeded")
}
}
}
}
#[derive(Debug)]
enum WaitOutcome {
Started,
Finished,
DeadlineExceeded,
}
pub fn grpc_error_response(code: u32, message: &str) -> Response<Body> {
let mut resp = Response::<Body>::with_status(200.into());
resp.headers.insert(
HeaderName::from_lowercase("content-type"),
HeaderValue::from_static(CONTENT_TYPE),
);
resp.headers.insert(
HeaderName::from_lowercase("grpc-status"),
HeaderValue::from_bytes(code.to_string().as_bytes())
.unwrap_or_else(|_| HeaderValue::from_static("13")),
);
resp.headers.insert(
HeaderName::from_lowercase("grpc-message"),
HeaderValue::from_bytes(message.as_bytes())
.unwrap_or_else(|_| HeaderValue::from_static("")),
);
resp.body = Body::Empty;
resp
}