use crate::error::NylonRingHostError;
use crate::{PluginCallGuard, context, context::HostContext};
use dashmap::DashMap;
use nylon_ring::NrStatus;
use rustc_hash::FxBuildHasher;
use std::collections::HashMap;
use std::sync::Arc;
use tokio::sync::{mpsc, oneshot};
pub type Result<T> = std::result::Result<T, NylonRingHostError>;
#[derive(Debug)]
pub(crate) enum Pending {
Unary(oneshot::Sender<(NrStatus, Vec<u8>)>),
Stream(mpsc::Sender<StreamFrame>),
}
#[derive(Debug)]
pub struct StreamFrame {
pub status: NrStatus,
pub data: Vec<u8>,
}
pub struct StreamReceiver {
inner: mpsc::Receiver<StreamFrame>,
host_ctx: Arc<HostContext>,
sid: u64,
_call_guard: Option<PluginCallGuard>,
}
impl StreamReceiver {
pub(crate) fn new(
inner: mpsc::Receiver<StreamFrame>,
host_ctx: Arc<HostContext>,
sid: u64,
call_guard: Option<PluginCallGuard>,
) -> Self {
Self {
inner,
host_ctx,
sid,
_call_guard: call_guard,
}
}
pub async fn recv(&mut self) -> Option<StreamFrame> {
self.inner.recv().await
}
}
impl Drop for StreamReceiver {
fn drop(&mut self) {
context::cleanup_sid(&self.host_ctx, self.sid);
}
}
pub(crate) type FastPendingMap = DashMap<u64, Pending, FxBuildHasher>;
pub(crate) type FastStateMap = DashMap<u64, HashMap<String, Vec<u8>>, FxBuildHasher>;
pub(crate) struct UnaryResultSlot {
pub(crate) sid: u64,
pub(crate) result: Option<(NrStatus, Vec<u8>)>,
}