use {
agave_scheduler_bindings::{
PackToWorkerMessage, ProgressMessage, TpuToPackMessage, WorkerToPackMessage,
},
rts_alloc::Allocator,
thiserror::Error,
};
pub(crate) type RtsAllocError = rts_alloc::error::Error;
pub(crate) type ShaqError = shaq::error::Error;
pub const MAX_WORKERS: usize = 64;
pub(crate) const VERSION: u64 = 4;
pub(crate) const LOGON_SUCCESS: u8 = 0x01;
pub(crate) const LOGON_FAILURE: u8 = 0x02;
pub(crate) const MAX_ALLOCATOR_HANDLES: usize = 128;
pub(crate) const GLOBAL_ALLOCATORS: usize = 1;
#[derive(Debug, Default, Clone, Copy)]
#[repr(C)]
pub struct ClientLogon {
pub worker_count: usize,
pub allocator_size: usize,
pub allocator_handles: usize,
pub tpu_to_pack_capacity: usize,
pub progress_tracker_capacity: usize,
pub pack_to_worker_capacity: usize,
pub worker_to_pack_capacity: usize,
pub flags: u16,
}
impl ClientLogon {
pub fn try_from_bytes(buffer: &[u8]) -> Option<Self> {
if buffer.len() != core::mem::size_of::<Self>() {
return None;
}
Some(unsafe { core::ptr::read_unaligned(buffer.as_ptr().cast()) })
}
}
pub mod logon_flags {}
pub struct ClientSession {
pub allocators: Vec<Allocator>,
pub tpu_to_pack: shaq::spsc::Consumer<TpuToPackMessage>,
pub progress_tracker: shaq::spsc::Consumer<ProgressMessage>,
pub workers: Vec<ClientWorkerSession>,
}
pub struct ClientWorkerSession {
pub pack_to_worker: shaq::spsc::Producer<PackToWorkerMessage>,
pub worker_to_pack: shaq::spsc::Consumer<WorkerToPackMessage>,
}
#[derive(Debug, Error)]
pub enum ClientHandshakeError {
#[error("Io; err={0}")]
Io(#[from] std::io::Error),
#[error("Timed out")]
TimedOut,
#[error("Protocol violation")]
ProtocolViolation,
#[error("Rejected; reason={0}")]
Rejected(String),
#[error("Rts alloc; err={0}")]
RtsAlloc(#[from] RtsAllocError),
#[error("Shaq; err={0}")]
Shaq(#[from] ShaqError),
}
pub struct AgaveSession {
pub flags: u16,
pub tpu_to_pack: AgaveTpuToPackSession,
pub progress_tracker: shaq::spsc::Producer<ProgressMessage>,
pub workers: Vec<AgaveWorkerSession>,
}
pub struct AgaveTpuToPackSession {
pub allocator: Allocator,
pub producer: shaq::spsc::Producer<TpuToPackMessage>,
}
pub struct AgaveWorkerSession {
pub allocator: Allocator,
pub pack_to_worker: shaq::spsc::Consumer<PackToWorkerMessage>,
pub worker_to_pack: shaq::spsc::Producer<WorkerToPackMessage>,
}
#[derive(Debug, Error)]
pub enum AgaveHandshakeError {
#[error("Io; err={0}")]
Io(#[from] std::io::Error),
#[error("Timeout")]
Timeout,
#[error("Close during handshake")]
EofDuringHandshake,
#[error("Version; server={server}; client={client}")]
Version { server: u64, client: u64 },
#[error("Worker count; count={0}")]
WorkerCount(usize),
#[error("Allocator handles; count={0}")]
AllocatorHandles(usize),
#[error("Rts alloc; err={0:?}")]
RtsAlloc(#[from] RtsAllocError),
#[error("Shaq; err={0:?}")]
Shaq(#[from] ShaqError),
}