use std::collections::BTreeMap;
use vyre_driver::{BackendError, BackendRegistration, Completion, DeviceIdentity};
use vyre_foundation::diagnostics::RetryClass;
use vyre_megakernel::{ArtifactValueId, Digest};
use crate::artifact_admission::{ArtifactSession, ArtifactSessionError, RetainedArtifactSession};
use crate::recovery::classify_backend_error;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ResidentQueueState {
pub control: Vec<u8>,
pub ring: Vec<u8>,
pub debug_log: Vec<u8>,
pub io_queue: Vec<u8>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ResidentQueueCompletion {
pub state: ResidentQueueState,
pub device_ns: Option<u64>,
}
#[derive(Clone, Copy)]
struct ResidentQueueAbi {
control: ArtifactValueId,
ring: ArtifactValueId,
debug_log: ArtifactValueId,
io_queue: ArtifactValueId,
}
impl ResidentQueueAbi {
fn resolve(session: &ArtifactSession) -> Result<Self, ArtifactSessionError> {
Ok(Self {
control: session.resource("control")?,
ring: session.resource("ring_buffer")?,
debug_log: session.resource("debug_log")?,
io_queue: session.resource("io_queue")?,
})
}
fn bind(self, state: ResidentQueueState) -> BTreeMap<ArtifactValueId, Vec<u8>> {
BTreeMap::from([
(self.control, state.control),
(self.ring, state.ring),
(self.debug_log, state.debug_log),
(self.io_queue, state.io_queue),
])
}
fn completion(
self,
completion: Completion,
) -> Result<ResidentQueueCompletion, ArtifactSessionError> {
let mut retained = completion.retained;
let mut take = |value: ArtifactValueId,
name: &str|
-> Result<Vec<u8>, ArtifactSessionError> {
retained.remove(&value).ok_or_else(|| {
BackendError::InvalidProgram {
fix: format!(
"Fix: persistent artifact completion must return retained queue resource `{name}`."
),
}
.into()
})
};
let state = ResidentQueueState {
control: take(self.control, "control")?,
ring: take(self.ring, "ring_buffer")?,
debug_log: take(self.debug_log, "debug_log")?,
io_queue: take(self.io_queue, "io_queue")?,
};
if !retained.is_empty() {
return Err(BackendError::InvalidProgram {
fix: "Fix: persistent queue artifacts must expose exactly control, ring_buffer, debug_log, and io_queue as retained values.".to_string(),
}
.into());
}
Ok(ResidentQueueCompletion {
state,
device_ns: completion.device_ns,
})
}
}
pub struct PersistentExecutor {
session: RetainedArtifactSession,
abi: ResidentQueueAbi,
}
impl PersistentExecutor {
pub fn from_bytes(
registration: &'static BackendRegistration,
envelope_bytes: &[u8],
initial: ResidentQueueState,
) -> Result<Self, ArtifactSessionError> {
let session = ArtifactSession::from_bytes(registration, envelope_bytes)?;
let abi = ResidentQueueAbi::resolve(&session)?;
let session = RetainedArtifactSession::new(session, abi.bind(initial))?;
Ok(Self { session, abi })
}
pub fn artifact(&self) -> Result<Digest, ArtifactSessionError> {
self.session.artifact()
}
pub fn device(&self) -> Result<DeviceIdentity, ArtifactSessionError> {
self.session.device()
}
pub fn submit_and_wait(
&self,
state: ResidentQueueState,
) -> Result<ResidentQueueCompletion, ArtifactSessionError> {
self.session.replace_retained(self.abi.bind(state))?;
let bindings = self.session.bindings()?;
let completion = self.session.submit_and_wait(bindings)?;
self.abi.completion(completion)
}
pub fn recover(&self, failure: BackendError) -> Result<DeviceIdentity, ArtifactSessionError> {
if classify_backend_error(&failure) != RetryClass::NewDevice {
return Err(failure.into());
}
self.session.rematerialize()
}
}