use appcore_contracts::InstallationId;
use appcore_types::{ClusterId, CoreId, TenantId};
use axum::extract::ws::Message;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use tokio::sync::mpsc::Sender;
pub const CONNECTION_BUFFER_CAPACITY: usize = 128;
static CONNECTION_GENERATION: AtomicU64 = AtomicU64::new(1);
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
pub struct WorkerConnectionKey {
pub tenant_id: TenantId,
pub installation_id: InstallationId,
pub core_id: CoreId,
}
#[derive(Debug, Clone)]
pub struct WorkerConnection {
pub key: WorkerConnectionKey,
pub sender: Sender<Message>,
cluster_id: Option<ClusterId>,
generation: u64,
last_heartbeat_ms: Arc<AtomicU64>,
}
impl WorkerConnection {
pub fn new(key: WorkerConnectionKey, sender: Sender<Message>, now_ms: u64) -> Self {
Self::new_inner(key, None, sender, now_ms)
}
pub fn new_in_cluster(
key: WorkerConnectionKey,
cluster_id: ClusterId,
sender: Sender<Message>,
now_ms: u64,
) -> Self {
Self::new_inner(key, Some(cluster_id), sender, now_ms)
}
fn new_inner(
key: WorkerConnectionKey,
cluster_id: Option<ClusterId>,
sender: Sender<Message>,
now_ms: u64,
) -> Self {
Self {
key,
sender,
cluster_id,
generation: CONNECTION_GENERATION.fetch_add(1, Ordering::Relaxed),
last_heartbeat_ms: Arc::new(AtomicU64::new(now_ms)),
}
}
pub fn cluster_id(&self) -> Option<&ClusterId> {
self.cluster_id.as_ref()
}
pub(crate) fn generation(&self) -> u64 {
self.generation
}
pub fn update_heartbeat(&self, now_ms: u64) {
self.last_heartbeat_ms.store(now_ms, Ordering::SeqCst);
}
pub fn last_heartbeat(&self) -> u64 {
self.last_heartbeat_ms.load(Ordering::SeqCst)
}
pub fn send(&self, message: Message) -> Result<(), crate::error::GatewayError> {
self.sender.try_send(message).map_err(|_| {
crate::error::GatewayError::Transport("worker connection closed".to_string())
})
}
}
#[derive(Debug, Clone)]
pub struct ClientConnection {
pub connection_id: String,
pub tenant_id: TenantId,
pub session_id: String,
pub sender: Sender<Message>,
}
impl ClientConnection {
pub fn new(
connection_id: String,
tenant_id: TenantId,
session_id: String,
sender: Sender<Message>,
) -> Self {
Self {
connection_id,
tenant_id,
session_id,
sender,
}
}
pub fn send(&self, message: Message) -> Result<(), crate::error::GatewayError> {
self.sender.try_send(message).map_err(|_| {
crate::error::GatewayError::Transport("client connection closed".to_string())
})
}
}