use appcore_contracts::ServiceId;
use appcore_types::{
CapabilityDescriptor, ClusterId, CoreId, CoreIdentity, DistributedCoreManifest, PeerEndpoint,
RuntimeOperationalMode, TenantId, TraceContext,
};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use std::future::Future;
use std::pin::Pin;
pub const CONTROL_PLANE_PROTOCOL_VERSION: u16 = 1;
pub const CONTROL_REGISTER_PATH: &str = "/v1/control/register";
pub const CONTROL_HEARTBEAT_PATH: &str = "/v1/control/heartbeat";
pub const CONTROL_PEERS_PATH: &str = "/v1/control/peers";
pub const CONTROL_SERVICE_LEASE_PATH: &str = "/v1/control/service-lease";
pub const CONTROL_SERVICE_LEASE_RELEASE_PATH: &str = "/v1/control/service-lease/release";
pub type ControlPlaneResult<T> = Result<T, ControlPlaneError>;
pub type ControlPlaneFuture<'a, T> =
Pin<Box<dyn Future<Output = ControlPlaneResult<T>> + Send + 'a>>;
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
pub enum ControlPlaneError {
#[error("control plane is offline")]
Offline,
#[error("control plane request timed out")]
Timeout,
#[error("control plane rejected operation: {0}")]
Rejected(String),
#[error("control plane state conflict: {0}")]
Conflict(String),
#[error("invalid control plane response: {0}")]
InvalidResponse(String),
#[error("control plane transport failed: {0}")]
Transport(String),
#[error("control plane lease is unavailable")]
LeaseUnavailable,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CoreRegistration {
pub manifest: DistributedCoreManifest,
pub registered_at_ms: u64,
pub operation_mode: RuntimeOperationalMode,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CorePresence {
pub identity: CoreIdentity,
pub operation_mode: RuntimeOperationalMode,
pub healthy: bool,
pub last_seen_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct HeartbeatRequest {
pub identity: CoreIdentity,
pub operation_mode: RuntimeOperationalMode,
pub sent_at_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct HeartbeatResponse {
pub accepted: bool,
pub server_time_ms: u64,
pub operation_mode: RuntimeOperationalMode,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PeerDirectory {
pub tenant_id: TenantId,
pub cluster_id: Option<ClusterId>,
pub peers: Vec<PeerRecord>,
pub refreshed_at_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct PeerRecord {
pub identity: CoreIdentity,
pub endpoints: Vec<PeerEndpoint>,
pub capabilities: Vec<CapabilityDescriptor>,
pub healthy: bool,
pub last_seen_ms: u64,
pub metadata: BTreeMap<String, String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ServiceLeaderLease {
pub service_id: ServiceId,
pub tenant_id: TenantId,
pub cluster_id: ClusterId,
pub holder_core_id: CoreId,
pub epoch: u64,
pub acquired_at_ms: u64,
pub expires_at_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ServiceLeaseRequest {
pub identity: CoreIdentity,
pub service_id: ServiceId,
pub ttl_ms: u64,
pub now_ms: u64,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct EmptyResponse {}
pub trait ControlPlaneProvider: Send + Sync {
fn register<'a>(
&'a self,
registration: CoreRegistration,
) -> ControlPlaneFuture<'a, CorePresence>;
fn heartbeat<'a>(
&'a self,
request: HeartbeatRequest,
) -> ControlPlaneFuture<'a, HeartbeatResponse>;
fn discover_peers<'a>(
&'a self,
identity: &'a CoreIdentity,
) -> ControlPlaneFuture<'a, PeerDirectory>;
fn acquire_or_renew_service_lease<'a>(
&'a self,
identity: &'a CoreIdentity,
service_id: &'a ServiceId,
ttl_ms: u64,
now_ms: u64,
) -> ControlPlaneFuture<'a, ServiceLeaderLease>;
fn release_service_lease<'a>(&'a self, lease: ServiceLeaderLease)
-> ControlPlaneFuture<'a, ()>;
fn register_traced<'a>(
&'a self,
registration: CoreRegistration,
_trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, CorePresence> {
self.register(registration)
}
fn heartbeat_traced<'a>(
&'a self,
request: HeartbeatRequest,
_trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, HeartbeatResponse> {
self.heartbeat(request)
}
fn discover_peers_traced<'a>(
&'a self,
identity: &'a CoreIdentity,
_trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, PeerDirectory> {
self.discover_peers(identity)
}
fn acquire_or_renew_service_lease_traced<'a>(
&'a self,
identity: &'a CoreIdentity,
service_id: &'a ServiceId,
ttl_ms: u64,
now_ms: u64,
_trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, ServiceLeaderLease> {
self.acquire_or_renew_service_lease(identity, service_id, ttl_ms, now_ms)
}
fn release_service_lease_traced<'a>(
&'a self,
lease: ServiceLeaderLease,
_trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, ()> {
self.release_service_lease(lease)
}
}
pub trait DiscoveryProvider: Send + Sync {
fn discover<'a>(&'a self, identity: &'a CoreIdentity) -> ControlPlaneFuture<'a, PeerDirectory>;
fn discover_traced<'a>(
&'a self,
identity: &'a CoreIdentity,
trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, PeerDirectory>;
}
impl<T> DiscoveryProvider for T
where
T: ControlPlaneProvider + ?Sized,
{
fn discover<'a>(&'a self, identity: &'a CoreIdentity) -> ControlPlaneFuture<'a, PeerDirectory> {
self.discover_peers(identity)
}
fn discover_traced<'a>(
&'a self,
identity: &'a CoreIdentity,
trace: Option<&'a TraceContext>,
) -> ControlPlaneFuture<'a, PeerDirectory> {
self.discover_peers_traced(identity, trace)
}
}
pub trait ServiceLeadershipGuard: Send + Sync {
fn current_service_lease(&self, service_id: &ServiceId) -> Option<ServiceLeaderLease>;
fn check_service_write_permission(
&self,
service_id: &ServiceId,
tenant_id: &TenantId,
cluster_id: &ClusterId,
core_id: &CoreId,
min_epoch: Option<u64>,
now_ms: u64,
) -> LeadershipDecision;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LeadershipDecision {
Allowed,
NoLease,
Expired,
StaleEpoch,
WrongHolder,
}
#[cfg(test)]
#[path = "tests.rs"]
mod tests;