use std::{
collections::VecDeque,
sync::{
Arc, Condvar, Mutex,
atomic::{AtomicU64, Ordering},
mpsc::{SyncSender, sync_channel},
},
thread::{self, JoinHandle},
time::{Duration, Instant},
};
use anyhow::{Context, Result, anyhow, bail};
use skippy_server::frontend::{
GenerationAbort, GenerationCommit, GenerationReceipt, GenerationStart, LinearProposal,
LinearProposalDiscardReason, LinearProposalQuery, LinearProposalReceipt,
LinearProposalSourceOutcome, LinearProposalSourceTelemetry,
};
use crate::ActivePlugin;
const PLUGIN_COMMAND_CAPACITY: usize = 1_024;
const PLUGIN_TERMINAL_RESERVE: usize = 64;
const CLEAN_SHUTDOWN_TIMEOUT: Duration = Duration::from_secs(5);
pub(crate) enum PluginCommand {
Begin(GenerationStart),
Committed(GenerationCommit),
Abort(GenerationAbort),
Finish(GenerationReceipt),
Proposal(LinearProposalQuery, SyncSender<ProposalResponse>),
ReportHandoff(LinearProposalReceipt),
Report(LinearProposalReceipt, SyncSender<Result<()>>),
Discard(Vec<u8>, LinearProposalDiscardReason),
Fence(SyncSender<()>),
}
pub(crate) struct ProposalResponse {
pub(crate) proposal: std::result::Result<Option<LinearProposal>, String>,
pub(crate) telemetry: LinearProposalSourceTelemetry,
}
struct QueuedPluginCommand {
enqueued_at: Instant,
command: PluginCommand,
}
struct QueueState {
commands: VecDeque<QueuedPluginCommand>,
closed: bool,
}
pub(crate) struct PluginCommandQueue {
state: Mutex<QueueState>,
available: Condvar,
}
#[derive(Debug)]
pub(crate) enum PluginCommandQueueError {
Full,
Stopped,
Poisoned,
}
impl PluginCommandQueue {
pub(crate) fn new() -> Self {
Self {
state: Mutex::new(QueueState {
commands: VecDeque::with_capacity(PLUGIN_COMMAND_CAPACITY),
closed: false,
}),
available: Condvar::new(),
}
}
fn enqueue(&self, command: PluginCommand) -> Result<()> {
self.try_enqueue(command).map_err(|error| match error {
PluginCommandQueueError::Full => anyhow!("native serving plugin command queue is full"),
PluginCommandQueueError::Stopped => anyhow!("native serving plugin worker stopped"),
PluginCommandQueueError::Poisoned => {
anyhow!("native serving plugin command queue lock poisoned")
}
})
}
pub(crate) fn try_enqueue(
&self,
command: PluginCommand,
) -> std::result::Result<(), PluginCommandQueueError> {
self.enqueue_within(command, PLUGIN_COMMAND_CAPACITY)
}
fn try_enqueue_terminal(
&self,
command: PluginCommand,
) -> std::result::Result<(), PluginCommandQueueError> {
let capacity = if matches!(command, PluginCommand::Fence(_)) {
PLUGIN_COMMAND_CAPACITY + PLUGIN_TERMINAL_RESERVE
} else {
PLUGIN_COMMAND_CAPACITY + PLUGIN_TERMINAL_RESERVE - 1
};
self.enqueue_within(command, capacity)
}
fn enqueue_within(
&self,
command: PluginCommand,
capacity: usize,
) -> std::result::Result<(), PluginCommandQueueError> {
let mut state = self
.state
.lock()
.map_err(|_| PluginCommandQueueError::Poisoned)?;
if state.closed {
return Err(PluginCommandQueueError::Stopped);
}
if state.commands.len() >= capacity {
return Err(PluginCommandQueueError::Full);
}
state.commands.push_back(QueuedPluginCommand {
enqueued_at: Instant::now(),
command,
});
self.available.notify_one();
Ok(())
}
pub(crate) fn close(&self) {
if let Ok(mut state) = self.state.lock() {
state.closed = true;
}
self.available.notify_all();
}
fn next(&self) -> Option<QueuedPluginCommand> {
let mut state = self
.state
.lock()
.expect("native serving plugin command queue lock must not be poisoned");
loop {
if let Some(command) = state.commands.pop_front() {
return Some(command);
}
if state.closed {
return None;
}
state = self
.available
.wait(state)
.expect("native serving plugin command queue lock must not be poisoned");
}
}
}
struct WorkerExit {
exited: Mutex<bool>,
signal: Condvar,
}
impl WorkerExit {
fn new() -> Self {
Self {
exited: Mutex::new(false),
signal: Condvar::new(),
}
}
fn mark_exited(&self) {
if let Ok(mut exited) = self.exited.lock() {
*exited = true;
}
self.signal.notify_all();
}
fn wait_for_exit(&self, timeout: Duration) -> bool {
let Ok(exited) = self.exited.lock() else {
return false;
};
self.signal
.wait_timeout_while(exited, timeout, |exited| !*exited)
.is_ok_and(|(exited, _)| *exited)
}
}
struct WorkerStopGuard {
queue: Arc<PluginCommandQueue>,
exit: Arc<WorkerExit>,
}
impl Drop for WorkerStopGuard {
fn drop(&mut self) {
self.queue.close();
self.exit.mark_exited();
}
}
pub(crate) struct PluginDriver {
pub(crate) queue: Arc<PluginCommandQueue>,
passive_queue: Arc<PluginCommandQueue>,
active: Arc<ActivePlugin>,
pub(crate) fatal_error: Arc<Mutex<Option<String>>>,
lifecycle_delivery_failures: Arc<AtomicU64>,
report_delivery_failures: Arc<AtomicU64>,
worker: WorkerHandle,
passive_worker: WorkerHandle,
}
struct WorkerHandle {
exit: Arc<WorkerExit>,
handle: Mutex<Option<JoinHandle<()>>>,
}
impl PluginDriver {
pub(crate) fn spawn(active: ActivePlugin) -> Result<Self> {
let queue = Arc::new(PluginCommandQueue::new());
let passive_queue = Arc::new(PluginCommandQueue::new());
let active = Arc::new(active);
let fatal_error = Arc::new(Mutex::new(None));
let lifecycle_delivery_failures = Arc::new(AtomicU64::new(0));
let report_delivery_failures = Arc::new(AtomicU64::new(0));
let exit = Arc::new(WorkerExit::new());
let passive_exit = Arc::new(WorkerExit::new());
let worker_queue = Arc::clone(&queue);
let worker_passive_queue = Arc::clone(&passive_queue);
let worker_active = Arc::clone(&active);
let worker_failures = Arc::clone(&lifecycle_delivery_failures);
let worker_report_failures = Arc::clone(&report_delivery_failures);
let worker_exit = Arc::clone(&exit);
let handle = thread::Builder::new()
.name("mesh-native-serving-plugin".to_string())
.spawn(move || {
plugin_worker(
worker_active,
worker_queue,
worker_passive_queue,
worker_failures,
worker_report_failures,
worker_exit,
);
})
.context("spawn native serving plugin worker")?;
let passive_worker_queue = Arc::clone(&passive_queue);
let passive_worker_active = Arc::clone(&active);
let passive_worker_exit = Arc::clone(&passive_exit);
let passive_handle = thread::Builder::new()
.name("mesh-native-serving-plugin-passive".to_string())
.spawn(move || {
plugin_passive_worker(
passive_worker_active,
passive_worker_queue,
passive_worker_exit,
);
})
.context("spawn native serving plugin passive worker")?;
Ok(Self {
queue,
passive_queue,
active,
fatal_error,
lifecycle_delivery_failures,
report_delivery_failures,
worker: WorkerHandle {
exit,
handle: Mutex::new(Some(handle)),
},
passive_worker: WorkerHandle {
exit: passive_exit,
handle: Mutex::new(Some(passive_handle)),
},
})
}
pub(crate) fn ensure_healthy(&self) -> Result<()> {
let error = self
.fatal_error
.lock()
.map_err(|_| anyhow!("native serving plugin health lock poisoned"))?;
if let Some(error) = error.as_deref() {
bail!("native serving plugin worker failed: {error}");
}
Ok(())
}
pub(crate) fn enqueue(&self, command: PluginCommand) -> Result<()> {
self.ensure_healthy()?;
self.enqueue_recovery(command)
}
pub(crate) fn enqueue_recovery(&self, command: PluginCommand) -> Result<()> {
self.queue_for(&command).enqueue(command)
}
pub(crate) fn enqueue_terminal(&self, command: PluginCommand) -> Result<()> {
self.ensure_healthy()?;
self.queue_for(&command)
.try_enqueue_terminal(command)
.map_err(|error| match error {
PluginCommandQueueError::Full => {
anyhow!("native serving plugin terminal queue is full")
}
PluginCommandQueueError::Stopped => {
anyhow!("native serving plugin passive worker stopped")
}
PluginCommandQueueError::Poisoned => {
anyhow!("native serving plugin terminal queue lock poisoned")
}
})
}
fn queue_for(&self, command: &PluginCommand) -> &Arc<PluginCommandQueue> {
if matches!(
command,
PluginCommand::Report(_, _) | PluginCommand::Discard(_, _) | PluginCommand::Fence(_)
) {
&self.passive_queue
} else {
&self.queue
}
}
pub(crate) fn lifecycle_delivery_failures(&self) -> u64 {
self.lifecycle_delivery_failures.load(Ordering::Relaxed)
}
pub(crate) fn report_delivery_failures(&self) -> u64 {
self.report_delivery_failures.load(Ordering::Relaxed)
}
pub(crate) fn propose(&self, query: LinearProposalQuery) -> Result<ProposalResponse> {
self.ensure_healthy()?;
let deadline = query.deadline;
let submitted_at = Instant::now();
if submitted_at >= deadline {
return Ok(abstention(
0,
LinearProposalSourceOutcome::HostDeadlineExceeded,
));
}
let (reply, response) = sync_channel(0);
match self
.queue
.try_enqueue(PluginCommand::Proposal(query, reply))
{
Ok(()) => {}
Err(PluginCommandQueueError::Full) => {
return Ok(abstention(0, LinearProposalSourceOutcome::QueueFull));
}
Err(PluginCommandQueueError::Stopped) => {
bail!("native serving plugin worker stopped before accepting proposal")
}
Err(PluginCommandQueueError::Poisoned) => {
bail!("native serving plugin command queue lock poisoned")
}
}
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Ok(abstention(
elapsed_us(submitted_at),
LinearProposalSourceOutcome::HostDeadlineExceeded,
));
}
match response.recv_timeout(remaining) {
Ok(result) => Ok(result),
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => Ok(abstention(
elapsed_us(submitted_at),
LinearProposalSourceOutcome::HostDeadlineExceeded,
)),
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
bail!("native serving plugin worker stopped before replying")
}
}
}
}
fn abstention(queue_wait_us: u64, outcome: LinearProposalSourceOutcome) -> ProposalResponse {
ProposalResponse {
proposal: Ok(None),
telemetry: LinearProposalSourceTelemetry {
queue_wait_us,
callback_elapsed_us: 0,
outcome,
},
}
}
fn elapsed_us(started: Instant) -> u64 {
u64::try_from(started.elapsed().as_micros()).unwrap_or(u64::MAX)
}
impl Drop for PluginDriver {
fn drop(&mut self) {
stop_worker(&self.queue, &self.worker, "callback");
stop_worker(&self.passive_queue, &self.passive_worker, "passive");
if let Some(active) = Arc::get_mut(&mut self.active)
&& let Err(error) = active.shutdown()
{
eprintln!("native serving plugin shutdown failed: {error:#}");
}
}
}
fn stop_worker(queue: &Arc<PluginCommandQueue>, worker: &WorkerHandle, label: &str) {
queue.close();
if !worker.exit.wait_for_exit(CLEAN_SHUTDOWN_TIMEOUT) {
eprintln!(
"native serving plugin {label} worker did not stop within {CLEAN_SHUTDOWN_TIMEOUT:?}; \
deferring plugin shutdown to that thread"
);
return;
}
if let Ok(mut handle) = worker.handle.lock()
&& let Some(handle) = handle.take()
{
let _ = handle.join();
}
}
fn plugin_worker(
active: Arc<ActivePlugin>,
queue: Arc<PluginCommandQueue>,
passive_queue: Arc<PluginCommandQueue>,
lifecycle_delivery_failures: Arc<AtomicU64>,
report_delivery_failures: Arc<AtomicU64>,
exit: Arc<WorkerExit>,
) {
let _stop_guard = WorkerStopGuard {
queue: Arc::clone(&queue),
exit,
};
while let Some(QueuedPluginCommand {
enqueued_at,
command,
}) = queue.next()
{
let (result, lifecycle) = match command {
PluginCommand::Begin(event) => (active.begin(&event), true),
PluginCommand::Committed(event) => (active.committed(&event), true),
PluginCommand::Abort(event) => (active.abort(&event), true),
PluginCommand::Finish(event) => (
finish_after_passive_fence(&passive_queue, || active.finish(&event)),
true,
),
PluginCommand::ReportHandoff(event) => {
let result = report_after_passive_handoff(&passive_queue, event);
(result, false)
}
PluginCommand::Proposal(query, reply) => {
run_proposal(&active, &passive_queue, enqueued_at, query, &reply);
continue;
}
PluginCommand::Report(_, _)
| PluginCommand::Discard(_, _)
| PluginCommand::Fence(_) => {
unreachable!("passive plugin callbacks must use the passive worker queue")
}
};
if let Err(error) = result {
if lifecycle {
lifecycle_delivery_failures.fetch_add(1, Ordering::Relaxed);
} else {
report_delivery_failures.fetch_add(1, Ordering::Relaxed);
eprintln!("native serving plugin report handoff failed: {error:#}");
}
}
}
}
fn run_proposal(
active: &ActivePlugin,
passive_queue: &PluginCommandQueue,
enqueued_at: Instant,
query: LinearProposalQuery,
reply: &SyncSender<ProposalResponse>,
) {
let deadline = query.deadline;
let queue_wait_us = elapsed_us(enqueued_at);
if Instant::now() >= deadline {
send_proposal_response(
passive_queue,
reply,
abstention(
queue_wait_us,
LinearProposalSourceOutcome::DeadlineExceededBeforeDispatch,
),
);
return;
}
let remaining = deadline.saturating_duration_since(Instant::now());
let fence_timeout = remaining.min(CLEAN_SHUTDOWN_TIMEOUT);
let fence_consumes_query_deadline = remaining <= CLEAN_SHUTDOWN_TIMEOUT;
if let Err(error) = fence_passive(passive_queue, fence_timeout) {
if fence_consumes_query_deadline && matches!(&error, PassiveFenceError::Timeout(_)) {
send_proposal_response(
passive_queue,
reply,
abstention(
queue_wait_us,
LinearProposalSourceOutcome::HostDeadlineExceeded,
),
);
return;
}
eprintln!("native serving plugin proposal fence failed: {error}");
send_proposal_response(
passive_queue,
reply,
ProposalResponse {
proposal: Err(error.to_string()),
telemetry: LinearProposalSourceTelemetry {
queue_wait_us,
callback_elapsed_us: 0,
outcome: LinearProposalSourceOutcome::SourceError,
},
},
);
return;
}
if Instant::now() >= deadline {
send_proposal_response(
passive_queue,
reply,
abstention(
queue_wait_us,
LinearProposalSourceOutcome::DeadlineExceededBeforeDispatch,
),
);
return;
}
let callback_started = Instant::now();
let result = active.propose(query);
let callback_finished = Instant::now();
let callback_elapsed_us = elapsed_us(callback_started);
let deadline_missed = callback_finished >= deadline;
let (proposal, outcome) = match result {
Ok(Some(proposal)) if deadline_missed => {
discard_late_candidate(passive_queue, &proposal);
(
Ok(None),
LinearProposalSourceOutcome::CandidateReturnedTooLate,
)
}
Ok(Some(proposal)) => (Ok(Some(proposal)), LinearProposalSourceOutcome::Ready),
Ok(None) if deadline_missed => (
Ok(None),
LinearProposalSourceOutcome::DeadlineExceededInPlugin,
),
Ok(None) => (Ok(None), LinearProposalSourceOutcome::Abstained),
Err(error) => {
eprintln!("native serving plugin proposal failed: {error:#}");
(
Err(format!("{error:#}")),
LinearProposalSourceOutcome::SourceError,
)
}
};
send_proposal_response(
passive_queue,
reply,
ProposalResponse {
proposal,
telemetry: LinearProposalSourceTelemetry {
queue_wait_us,
callback_elapsed_us,
outcome,
},
},
);
}
fn send_proposal_response(
passive_queue: &PluginCommandQueue,
reply: &SyncSender<ProposalResponse>,
response: ProposalResponse,
) {
let candidate_decision_id = response
.proposal
.as_ref()
.ok()
.and_then(|proposal| proposal.as_ref())
.map(|proposal| proposal.decision_id.as_bytes().to_vec());
if reply.send(response).is_ok() {
return;
}
if let Some(decision_id) = candidate_decision_id
&& let Err(error) = passive_queue.try_enqueue_terminal(PluginCommand::Discard(
decision_id,
LinearProposalDiscardReason::DeadlineExceeded,
))
{
eprintln!(
"native serving plugin could not deliver the terminal discard for a detached proposal reply: {error:?}"
);
}
}
fn discard_late_candidate(passive_queue: &PluginCommandQueue, proposal: &LinearProposal) {
if let Err(error) = passive_queue.try_enqueue_terminal(PluginCommand::Discard(
proposal.decision_id.as_bytes().to_vec(),
LinearProposalDiscardReason::DeadlineExceeded,
)) {
eprintln!(
"native serving plugin could not deliver the terminal discard for a late proposal: \
{error:?}"
);
}
}
fn report_after_passive_handoff(
passive_queue: &PluginCommandQueue,
receipt: LinearProposalReceipt,
) -> Result<()> {
let (ack, result) = sync_channel(0);
passive_queue
.try_enqueue_terminal(PluginCommand::Report(receipt, ack))
.map_err(|error| match error {
PluginCommandQueueError::Full => {
anyhow!("native serving plugin terminal queue is full")
}
PluginCommandQueueError::Stopped => {
anyhow!("native serving plugin passive worker stopped")
}
PluginCommandQueueError::Poisoned => {
anyhow!("native serving plugin terminal queue lock poisoned")
}
})?;
match result.recv_timeout(CLEAN_SHUTDOWN_TIMEOUT) {
Ok(result) => result,
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => bail!(
"native serving plugin report callback did not complete within {CLEAN_SHUTDOWN_TIMEOUT:?}"
),
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => {
bail!("native serving plugin passive worker stopped before report ack")
}
}
}
#[derive(Debug)]
enum PassiveFenceError {
Timeout(Duration),
Failure(anyhow::Error),
}
impl std::fmt::Display for PassiveFenceError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Timeout(timeout) => write!(
formatter,
"native serving plugin passive callbacks did not complete within {timeout:?}"
),
Self::Failure(error) => write!(formatter, "{error:#}"),
}
}
}
fn fence_passive(
passive_queue: &PluginCommandQueue,
timeout: Duration,
) -> std::result::Result<(), PassiveFenceError> {
let (reply, response) = sync_channel(1);
passive_queue
.try_enqueue_terminal(PluginCommand::Fence(reply))
.map_err(|error| {
PassiveFenceError::Failure(match error {
PluginCommandQueueError::Full => {
anyhow!("native serving plugin terminal queue is full")
}
PluginCommandQueueError::Stopped => {
anyhow!("native serving plugin passive worker stopped")
}
PluginCommandQueueError::Poisoned => {
anyhow!("native serving plugin terminal queue lock poisoned")
}
})
})?;
match response.recv_timeout(timeout) {
Ok(()) => Ok(()),
Err(std::sync::mpsc::RecvTimeoutError::Timeout) => Err(PassiveFenceError::Timeout(timeout)),
Err(std::sync::mpsc::RecvTimeoutError::Disconnected) => Err(PassiveFenceError::Failure(
anyhow!("native serving plugin passive worker stopped before fence"),
)),
}
}
fn finish_after_passive_fence<T>(
passive_queue: &PluginCommandQueue,
finish: impl FnOnce() -> Result<T>,
) -> Result<T> {
fence_passive(passive_queue, CLEAN_SHUTDOWN_TIMEOUT).map_err(|error| anyhow!("{error}"))?;
finish()
}
fn plugin_passive_worker(
active: Arc<ActivePlugin>,
queue: Arc<PluginCommandQueue>,
exit: Arc<WorkerExit>,
) {
let _stop_guard = WorkerStopGuard {
queue: Arc::clone(&queue),
exit,
};
while let Some(queued) = queue.next() {
match queued.command {
PluginCommand::Report(event, ack) => {
let result = active.report(&event);
if ack.send(result).is_err() {
eprintln!(
"native serving plugin report callback acknowledgement receiver dropped"
);
}
}
PluginCommand::Discard(decision_id, reason) => {
let _ = active.discard(&decision_id, reason);
}
PluginCommand::Fence(reply) => {
let _ = reply.send(());
}
PluginCommand::Begin(_)
| PluginCommand::Committed(_)
| PluginCommand::Abort(_)
| PluginCommand::Finish(_)
| PluginCommand::ReportHandoff(_)
| PluginCommand::Proposal(_, _) => {
unreachable!("lifecycle and proposal callbacks must use the primary worker queue")
}
}
}
}
#[cfg(test)]
mod tests;