use super::*;
use super::env::{AGENT, REQUEST_TIMEOUT};
use super::projection::with_cluster;
use super::state::{DispatchInner, DispatchState};
pub(super) fn store_ack(inner: &DispatchInner, mut ack: AckArgs) {
if ack.op_id.is_none() {
ack.op_id = Some(onlyne_proto::new_op_id());
}
let Some(op_id) = ack.op_id.clone() else {
return;
};
match serde_json::to_value(ClientOp::Ack(ack)) {
Ok(value) => {
if let Err(error) = inner.store.enqueue_intent(&op_id, &value) {
tracing::warn!(error = %error, "settled ack was not stored");
}
}
Err(error) => tracing::warn!(error = %error, "settled ack did not serialize"),
}
}
pub(super) fn queue_outbound_locked(
inner: &mut DispatchInner,
envelope: &Envelope,
) -> Result<String> {
let mut stamped = envelope.clone();
let op_id = stamp_op_id(&mut stamped);
stamped
.validate()
.map_err(|error| anyhow!(error.to_string()))?;
inner
.store
.enqueue_intent(&op_id, &serde_json::to_value(&stamped)?)?;
Ok(op_id)
}
pub(super) async fn transport_envelope(state: &DispatchState, envelope: &Envelope) -> Result<()> {
let op = ClientOp::Send(Box::new(envelope.clone()));
let outbox = { state.inner.lock().outbox.clone() };
match outbox {
Some(outbox) => match outbox.send(op).await {
Ok(()) => Ok(()),
Err(error) => {
tracing::warn!(error = %error, "completion fell back to the intent queue");
state.enqueue_outbound(envelope).map(|_| ())
}
},
None => state.enqueue_outbound(envelope).map(|_| ()),
}
}
#[derive(Clone)]
pub struct ClientLink {
handle: ClientConn,
welcome: Arc<Welcome>,
hello: HandshakeArgs,
}
impl ClientLink {
pub async fn connect(init: &ClientInit, live_tasks: Vec<String>) -> Result<Self, NetError> {
let keypair = KeyPair::load(&init.key_path)?;
let settings = ConnSettings {
agent: AGENT.to_string(),
version: env!("CARGO_PKG_VERSION").to_string(),
..ConnSettings::new(PROTOCOL_VERSION)
};
let handle: ClientConn =
dial(&init.server, &keypair, &init.cert_pin, &init.role, settings).await?;
let hello = HandshakeArgs {
protocol: PROTOCOL_VERSION,
role: init.role.clone(),
key: keypair.public_str(),
signature: String::new(),
agent: AGENT.to_string(),
version: env!("CARGO_PKG_VERSION").to_string(),
aggregate: false,
live_tasks: Vec::new(),
};
let body = handle
.request(
Frame::req(
String::new(),
ClientOp::Hello(hello_with_live_tasks(&hello, live_tasks)),
),
REQUEST_TIMEOUT,
)
.await?;
if !body.ok {
let error = body.error.clone().unwrap_or(onlyne_proto::ErrorPayload {
code: onlyne_proto::ErrorCode::Internal,
message: "hello refused".to_string(),
field: None,
});
return Err(NetError::Rejected {
code: wire_code(error.code),
message: error.message,
});
}
let data = body.data().cloned().ok_or(NetError::BadFrame)?;
let welcome: Welcome = serde_json::from_value(data).map_err(|_| NetError::BadFrame)?;
Ok(Self {
handle,
welcome: Arc::new(welcome),
hello,
})
}
pub async fn authenticate(&self, live_tasks: Vec<String>) -> Result<(), NetError> {
let body = self
.handle
.request(
Frame::req(
String::new(),
ClientOp::Hello(hello_with_live_tasks(&self.hello, live_tasks)),
),
REQUEST_TIMEOUT,
)
.await?;
if !body.ok {
let error = body.error.clone().unwrap_or(onlyne_proto::ErrorPayload {
code: onlyne_proto::ErrorCode::Internal,
message: "hello refused".to_string(),
field: None,
});
return Err(NetError::Rejected {
code: wire_code(error.code),
message: error.message,
});
}
Ok(())
}
pub fn welcome(&self) -> &Welcome {
&self.welcome
}
pub async fn request(&self, op: ClientOp) -> Result<ResBody, NetError> {
self.handle
.request(Frame::req(String::new(), op), REQUEST_TIMEOUT)
.await
}
pub fn events(&self) -> broadcast::Receiver<Frame<ClientOp>> {
self.handle.events()
}
pub async fn close(&self) -> Result<(), NetError> {
self.handle.close().await
}
pub fn readiness(&self) -> ConnReadiness {
self.handle.readiness()
}
pub async fn failure(&self) -> Option<NetError> {
self.handle.failure().await
}
}
pub(crate) fn hello_with_live_tasks(
hello: &HandshakeArgs,
live_tasks: Vec<String>,
) -> HandshakeArgs {
let mut hello = hello.clone();
hello.live_tasks = live_tasks;
hello
}
fn wire_code(code: onlyne_proto::ErrorCode) -> String {
serde_json::to_value(code)
.ok()
.and_then(|value| value.as_str().map(str::to_string))
.unwrap_or_else(|| "internal".to_string())
}
pub trait Outbox: Send + Sync {
fn send(&self, op: ClientOp)
-> Pin<Box<dyn Future<Output = Result<(), NetError>> + Send + '_>>;
fn request(
&self,
op: ClientOp,
) -> Pin<Box<dyn Future<Output = Result<ResBody, NetError>> + Send + '_>>;
}
impl Outbox for ClientLink {
fn send(
&self,
op: ClientOp,
) -> Pin<Box<dyn Future<Output = Result<(), NetError>> + Send + '_>> {
Box::pin(async move { self.request(op).await.map(|_| ()) })
}
fn request(
&self,
op: ClientOp,
) -> Pin<Box<dyn Future<Output = Result<ResBody, NetError>> + Send + '_>> {
Box::pin(async move { ClientLink::request(self, op).await })
}
}
pub async fn send_frame(state: &DispatchState, op: ClientOp) -> Result<()> {
let op = match op {
ClientOp::Report(report) => ClientOp::Report(with_cluster(state, report)),
other => other,
};
if let Some(outbox) = state.outbox() {
if outbox.send(op.clone()).await.is_ok() {
return Ok(());
}
}
state.enqueue_op(&op)?;
Ok(())
}
impl DispatchState {
pub fn enqueue_outbound(&self, envelope: &Envelope) -> Result<String> {
queue_outbound_locked(&mut self.inner.lock(), envelope)
}
pub fn accept_new(&self) -> Arc<AtomicBool> {
self.inner.lock().accept_new.clone()
}
pub fn link_up(&self) -> bool {
self.inner.lock().link_up.load(Ordering::SeqCst)
}
pub fn set_link_up(&self, up: bool) {
self.inner.lock().link_up.store(up, Ordering::SeqCst);
}
pub fn cluster_ref(&self) -> String {
self.inner.lock().cluster_ref.clone()
}
pub fn set_topology(&self, cluster: &str) {
self.inner.lock().topology = cluster.trim().to_string();
}
pub fn topology(&self) -> String {
self.inner.lock().topology.clone()
}
pub fn set_cluster_ref(&self, aggregate: impl Into<String>) {
self.inner.lock().cluster_ref = aggregate.into();
}
pub async fn request(&self, op: ClientOp) -> Result<ResBody, NetError> {
let outbox = { self.inner.lock().outbox.clone() };
let Some(outbox) = outbox else {
return Err(NetError::NotReady);
};
outbox.request(op).await
}
pub fn attach_outbox(&self, outbox: Arc<dyn Outbox>) {
self.inner.lock().outbox = Some(outbox);
}
pub fn detach_outbox(&self) {
self.inner.lock().outbox = None;
}
fn outbox(&self) -> Option<Arc<dyn Outbox>> {
self.inner.lock().outbox.clone()
}
pub fn enqueue_op(&self, op: &ClientOp) -> Result<String> {
let op_id = onlyne_proto::new_id();
self.inner
.lock()
.store
.enqueue_intent(&op_id, &serde_json::to_value(op)?)?;
Ok(op_id)
}
}