mod era;
mod inertia;
mod settle_ctx;
mod spawn;
mod spawn_state;
use crate::context::{Context, Root};
use crate::effect::{Cleanup, DisposableList, TaskRegistrationError};
use crate::registry::ResidencyClaim;
use parking_lot::Mutex;
use std::fmt;
use std::future::Future;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
#[derive(Clone)]
pub struct FiberId(Arc<u8>);
impl PartialEq for FiberId {
fn eq(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.0, &other.0)
}
}
impl Eq for FiberId {}
impl std::hash::Hash for FiberId {
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
std::ptr::hash(Arc::as_ptr(&self.0), state);
}
}
impl fmt::Debug for FiberId {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str("FiberId(..)")
}
}
use std::time::Duration;
pub(crate) use inertia::InertiaSlot;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct GenerationClosed;
use inertia::SemanticTarget;
pub use spawn::{PluginFailure, PluginFailureKind, SpawnError};
use spawn_state::SpawnState;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FiberRole {
Root,
Ordinary,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FiberState {
Loading,
Active,
Pending,
Unloading,
Failed,
Disposed,
}
impl FiberState {
fn publication_index(self) -> usize {
match self {
Self::Loading => 0,
Self::Active => 1,
Self::Pending => 2,
Self::Unloading => 3,
Self::Failed => 4,
Self::Disposed => 5,
}
}
}
struct StatePublicationInner {
current: FiberState,
revisions: [u64; 6],
}
struct StatePublication {
inner: Mutex<StatePublicationInner>,
changed: tokio::sync::Notify,
}
impl StatePublication {
fn new(initial: FiberState) -> Self {
let mut revisions = [0; 6];
revisions[initial.publication_index()] = 1;
Self {
inner: Mutex::new(StatePublicationInner {
current: initial,
revisions,
}),
changed: tokio::sync::Notify::new(),
}
}
fn current(&self) -> FiberState {
self.inner.lock().current
}
fn publish(&self, next: FiberState) -> Option<FiberState> {
let old = {
let mut inner = self.inner.lock();
let old = inner.current;
if old == next {
return None;
}
inner.current = next;
let revision = &mut inner.revisions[next.publication_index()];
*revision = revision
.checked_add(1)
.expect("state publication revision overflow");
old
};
Some(old)
}
fn notify(&self) {
self.changed.notify_waiters();
}
fn snapshot(&self, target: FiberState) -> (FiberState, u64) {
let inner = self.inner.lock();
(inner.current, inner.revisions[target.publication_index()])
}
async fn wait_for(
&self,
target: FiberState,
timeout: Duration,
) -> std::result::Result<(), WaitStateError> {
let (_, baseline_revision) = self.snapshot(target);
self.wait_for_after(target, baseline_revision, timeout)
.await
}
async fn wait_for_after(
&self,
target: FiberState,
baseline_revision: u64,
timeout: Duration,
) -> std::result::Result<(), WaitStateError> {
let deadline = crate::deadline::watchdog(timeout);
tokio::pin!(deadline);
loop {
let notified = self.changed.notified();
tokio::pin!(notified);
notified.as_mut().enable();
let (current, revision) = self.snapshot(target);
if current == target || revision != baseline_revision {
return Ok(());
}
tokio::select! {
_ = &mut deadline => return Err(WaitStateError::Elapsed),
_ = &mut notified => {}
}
}
}
}
#[cfg(test)]
mod state_publication_tests {
use super::{FiberState, StatePublication};
use std::time::Duration;
#[tokio::test]
async fn requested_publication_survives_state_notification_race() {
let state = StatePublication::new(FiberState::Pending);
let (_, before_loading) = state.snapshot(FiberState::Loading);
state.publish(FiberState::Loading);
state.notify();
state.publish(FiberState::Active);
state.notify();
state
.wait_for_after(
FiberState::Loading,
before_loading,
Duration::from_millis(20),
)
.await
.expect("the durable Loading revision survives a superseding Active publication");
assert_eq!(state.current(), FiberState::Active);
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum EraSwapFailure {
#[error("successor dependency `{service}` has a contract mismatch")]
SuccessorDependencyContractMismatch {
service: String,
},
#[error("successor dependency `{service}` configuration is invalid: {diagnostic}")]
SuccessorDependencyConfiguration {
service: String,
diagnostic: String,
},
#[error("successor apply failed: {0}")]
SuccessorApply(PluginFailure),
#[error("the fresh successor was lost before handoff")]
SuccessorLost,
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum EraSwapError {
#[error("prepared change targets a different plugin contract")]
PluginContractMismatch,
#[error("fiber is closed")]
Closed,
#[error("{0}")]
Recursion(LifecycleRecursion),
#[error("era replacement incomplete: {0}")]
Incomplete(EraSwapFailure),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum UpdateOutcome {
Committed(FiberState),
Vetoed,
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum UpdateError {
#[error("prepared change targets a different plugin contract")]
PluginContractMismatch,
#[error("fiber is closed")]
Closed,
#[error("{0}")]
Recursion(LifecycleRecursion),
#[error("update control failed: {0}")]
Control(crate::events::InvocationFailure),
#[error("fiber lost update admission before commit")]
AdmissionLost,
#[error("plugin apply failed: {0}")]
Apply(PluginFailure),
}
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub(crate) struct DependencyEdge {
pub(crate) service: String,
pub(crate) realm: crate::context::RealmKey,
}
impl DependencyEdge {
pub(crate) fn new(service: String, realm: crate::context::RealmKey) -> Self {
Self { service, realm }
}
}
pub(crate) struct Fiber {
id: FiberId,
alive: AtomicBool,
disposing: AtomicBool,
generation_replacing: AtomicBool,
dispose_complete: AtomicBool,
dispose_done: tokio::sync::Notify,
pub(crate) name: String,
dependency_edges: Box<[DependencyEdge]>,
state: StatePublication,
pub(crate) disposables: Mutex<DisposableList>,
pub(crate) slot: InertiaSlot,
pub(crate) error: Mutex<Option<Arc<PluginFailure>>>,
pub(crate) spawn_state: SpawnState,
pub(crate) residency: Mutex<Option<ResidencyClaim>>,
draining: AtomicBool,
pub(crate) applied: AtomicBool,
pub(crate) creation_pending: AtomicBool,
}
impl std::fmt::Debug for Fiber {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Fiber")
.field("id", &self.id)
.field("name", &self.name)
.field("state", &self.state())
.finish()
}
}
impl Fiber {
pub(crate) fn new(name: impl Into<String>) -> Arc<Self> {
Self::new_with_edges(name, Vec::new())
}
pub(crate) fn new_with_edges(
name: impl Into<String>,
dependency_edges: Vec<DependencyEdge>,
) -> Arc<Self> {
Arc::new(Self {
id: FiberId(Arc::new(0)),
alive: AtomicBool::new(true),
disposing: AtomicBool::new(false),
generation_replacing: AtomicBool::new(false),
dispose_complete: AtomicBool::new(false),
dispose_done: tokio::sync::Notify::new(),
name: name.into(),
dependency_edges: dependency_edges.into_boxed_slice(),
state: StatePublication::new(FiberState::Pending),
disposables: Mutex::new(DisposableList::new()),
slot: InertiaSlot::new(),
error: Mutex::new(None),
spawn_state: SpawnState::new(),
residency: Mutex::new(None),
draining: AtomicBool::new(false),
applied: AtomicBool::new(false),
creation_pending: AtomicBool::new(false),
})
}
pub(crate) fn root() -> Arc<Self> {
let fiber = Self::new("root");
fiber.set_state(FiberState::Active);
fiber
}
pub(crate) fn id(&self) -> &FiberId {
&self.id
}
pub(crate) fn dependency_edges(&self) -> &[DependencyEdge] {
&self.dependency_edges
}
pub(crate) fn state(&self) -> FiberState {
self.state.current()
}
pub(crate) fn set_state(&self, state: FiberState) {
self.state.publish(state);
}
pub(crate) async fn transition(self: &Arc<Self>, next: FiberState) {
let old = self.state();
if old == next {
return;
}
let root = self.spawn_state.root();
let (visibility, drift) = if (old == FiberState::Active) != (next == FiberState::Active) {
match root.as_ref() {
Some(root) => root.services.commit_fiber_transition(root, self, old, next),
None => {
self.set_state(next);
(Vec::new(), None)
}
}
} else {
self.set_state(next);
(Vec::new(), None)
};
if let (Some(root), Some(drift)) = (root.as_ref(), drift) {
drift.kick(root);
}
self.state.notify();
if let Some(root) = root {
for slot in visibility {
let id = crate::observation::ServicePublicationId(slot.occurrence);
root.observations.publish(
crate::observation::RuntimeObservation::ServiceVisibility {
service: slot.service,
realm: crate::ServiceRealm::new(root.realm_membership.clone(), slot.realm),
previous: (old == FiberState::Active).then(|| id.clone()),
current: (next == FiberState::Active).then_some(id),
},
);
}
root.observations
.publish(crate::observation::RuntimeObservation::FiberState {
fiber: self.id().clone(),
previous: old,
current: next,
});
}
}
pub(crate) fn fiber_ctx(self: &Arc<Self>) -> Option<Context> {
self.spawn_state.context_for(self)
}
pub(crate) fn is_alive(&self) -> bool {
self.alive.load(Ordering::Relaxed)
}
pub(crate) fn assert_alive(&self) -> std::result::Result<(), GenerationClosed> {
if self.is_alive() {
Ok(())
} else {
Err(GenerationClosed)
}
}
pub(crate) fn assert_can_register(&self) -> std::result::Result<(), GenerationClosed> {
self.assert_alive()?;
if self.disposing.load(Ordering::SeqCst) || self.generation_replacing.load(Ordering::SeqCst)
{
return Err(GenerationClosed);
}
match self.state() {
FiberState::Loading | FiberState::Active => Ok(()),
FiberState::Pending
| FiberState::Unloading
| FiberState::Failed
| FiberState::Disposed => Err(GenerationClosed),
}
}
fn store_error(&self, next: Option<Arc<PluginFailure>>) {
*self.error.lock() = next;
}
pub(crate) fn take_error(&self) -> Option<Arc<PluginFailure>> {
self.error.lock().take()
}
pub(crate) fn missing_from(&self, root: &Root) -> Vec<String> {
self.dependency_edges
.iter()
.filter(|edge| {
root.services
.occurrence_id(&edge.realm, &edge.service)
.is_none()
})
.map(|edge| edge.service.clone())
.collect()
}
pub(crate) async fn run_cleanup_contained(self: &Arc<Self>, cleanup: Cleanup) {
let logger = self.fiber_ctx().map(|ctx| ctx.logger());
let (tx, rx) = tokio::sync::oneshot::channel();
let fiber = self.clone();
let report_logger = logger.clone();
let work = async move {
let outcome = settle_ctx::bracket(fiber, async move {
crate::effect::execute_cleanup(cleanup).await
})
.await;
if let Err(outcome) = tx.send(outcome)
&& let Some(failure) = outcome
{
crate::effect::report_cleanup_failure(report_logger.as_ref(), &failure);
}
};
crate::effect::detach(work);
if let Ok(Some(failure)) = rx.await {
crate::effect::report_cleanup_failure(logger.as_ref(), &failure);
}
}
pub(crate) async fn run_tokens_contained(
self: &Arc<Self>,
tokens: Vec<crate::effect::DisposableToken>,
) {
for token in tokens.into_iter().rev() {
if let Some(cleanup) = self.remove_disposable(token) {
self.run_cleanup_contained(cleanup).await;
}
}
}
pub(crate) async fn drain_disposables(self: &Arc<Self>) {
let tokens: Vec<_> = self.disposables.lock().tokens();
self.draining.store(true, Ordering::SeqCst);
self.run_tokens_contained(tokens).await;
self.draining.store(false, Ordering::SeqCst);
}
fn compute_target(&self, root: &Root) -> SemanticTarget {
let input = self
.spawn_state
.effective_input()
.expect("settled plugin fibers have committed effective input");
let assignments = self
.dependency_edges
.iter()
.map(|edge| {
let publication = root.services.occurrence_id(&edge.realm, &edge.service);
(edge.clone(), publication)
})
.collect();
SemanticTarget::new(input, assignments)
}
pub(crate) fn kick_committed_recheck(self: &Arc<Self>, root: &Arc<Root>) {
self.slot.kick(self, root);
}
async fn settle_once(self: &Arc<Self>, target: &SemanticTarget, owns_first_apply: bool) {
settle_ctx::bracket(
self.clone(),
self.settle_once_inner(target, owns_first_apply),
)
.await
}
async fn settle_once_inner(self: &Arc<Self>, target: &SemanticTarget, owns_first_apply: bool) {
if !self.is_alive() {
return; }
match self.state() {
FiberState::Active | FiberState::Failed | FiberState::Loading => {
self.transition(FiberState::Unloading).await;
self.drain_disposables().await;
}
FiberState::Pending | FiberState::Unloading | FiberState::Disposed => {}
}
if target.has_missing() {
self.store_error(None);
self.transition(FiberState::Pending).await;
return;
}
if !owns_first_apply
&& self.creation_pending.load(Ordering::SeqCst)
&& !self.applied.load(Ordering::SeqCst)
{
return;
}
match self.spawn_state.apply_snapshot() {
Ok(state) => {
self.generation_replacing.store(false, Ordering::SeqCst);
self.transition(FiberState::Loading).await;
let ctx = state.scope.with_fiber(self.clone());
let outcome = crate::contained::catch_contained(
state.plugin.clone().apply_boxed(ctx.clone()),
)
.await;
self.applied.store(true, Ordering::SeqCst);
match outcome {
Ok(Ok(())) => {
self.store_error(None);
self.transition(FiberState::Active).await;
}
Ok(Err(failure)) => {
self.transition(FiberState::Unloading).await;
self.drain_disposables().await;
self.store_error(Some(Arc::new(failure)));
self.transition(FiberState::Failed).await;
}
Err(panic) => {
self.transition(FiberState::Unloading).await;
self.drain_disposables().await;
self.store_error(Some(Arc::new(PluginFailure::panicked(
crate::contained::payload_text(&panic),
))));
self.transition(FiberState::Failed).await;
}
}
}
Err(_) => {
unreachable!("settling a plugin Fiber requires installed spawn state")
}
}
}
pub(crate) async fn initial_settle(
self: &Arc<Self>,
root: &Arc<Root>,
) -> inertia::InitialOutcome {
self.slot.initial_spawn_pass(self, root).await
}
pub(crate) async fn dispose(self: &Arc<Self>) {
if self.dispose_complete.load(Ordering::SeqCst) {
return;
}
self.slot.claim().await;
if self.dispose_complete.load(Ordering::SeqCst) {
self.slot.abandon();
return;
}
if self.disposing.swap(true, Ordering::SeqCst) {
self.slot.abandon();
self.wait_dispose_complete().await;
return;
}
let fiber = self.clone();
crate::effect::detach(async move {
fiber.complete_claimed_dispose().await;
});
self.wait_dispose_complete().await;
}
pub(crate) async fn complete_claimed_dispose(self: &Arc<Self>) {
settle_ctx::bracket(self.clone(), async {
self.teardown_body().await;
self.slot.abandon();
})
.await;
self.release_residency();
self.publish_dispose_complete();
}
async fn wait_dispose_complete(&self) {
loop {
let notified = self.dispose_done.notified();
tokio::pin!(notified);
notified.as_mut().enable();
if self.dispose_complete.load(Ordering::SeqCst) {
return;
}
notified.await;
}
}
fn publish_dispose_complete(&self) {
self.dispose_complete.store(true, Ordering::SeqCst);
self.dispose_done.notify_waiters();
}
async fn teardown_body(self: &Arc<Self>) {
self.creation_pending.store(false, Ordering::SeqCst);
if self.state() == FiberState::Disposed {
return; }
self.transition(FiberState::Unloading).await;
self.slot.pin_inactive();
self.alive.store(false, Ordering::Relaxed);
let root = self.spawn_state.root();
if let Some(root) = root {
root.deps.unregister(self);
}
self.drain_disposables().await;
self.transition(FiberState::Disposed).await;
self.store_error(None);
}
pub(crate) async fn creation_teardown(self: &Arc<Self>) {
self.disposing.store(true, Ordering::SeqCst);
settle_ctx::bracket(self.clone(), self.teardown_body()).await;
self.force_unlink();
self.publish_dispose_complete();
}
pub(crate) async fn creation_interrupted_teardown(self: &Arc<Self>) {
if self.slot.take_creation_ownership() {
self.disposing.store(true, Ordering::SeqCst);
} else {
if self.disposing.swap(true, Ordering::SeqCst) {
self.wait_dispose_complete().await;
return;
}
self.slot.claim().await;
if !self.is_alive() {
self.slot.abandon();
return;
}
}
settle_ctx::bracket(self.clone(), self.teardown_body()).await;
self.slot.abandon();
self.force_unlink();
self.publish_dispose_complete();
}
pub(crate) fn force_unlink(self: &Arc<Self>) {
self.release_residency();
}
pub(crate) fn release_residency(&self) {
let claim = self.residency.lock().take();
if let Some(claim) = claim {
let root = self.spawn_state.root();
let snapshot = root
.as_ref()
.map(|root| crate::observation::fiber_snapshot(root, self));
claim.release();
if let (Some(root), Some(fiber)) = (root, snapshot) {
root.observations
.publish(crate::observation::RuntimeObservation::FiberResidency {
change: crate::observation::ResidencyChange::Removed,
fiber,
});
}
}
}
pub(crate) fn remove_disposable(
&self,
token: crate::effect::DisposableToken,
) -> Option<Cleanup> {
self.disposables.lock().remove(token)
}
}
pub(crate) fn refuse_group_removal_recursion(
fibers: &[Arc<Fiber>],
) -> std::result::Result<(), LifecycleRecursion> {
for fiber in fibers {
settle_ctx::refuse_recursion(fiber, LifecycleOperation::RemovePlugins)?;
}
Ok(())
}
impl Context {
pub fn run<F>(&self, task: F) -> std::result::Result<(), TaskRegistrationError>
where
F: Future + Send + 'static,
{
let fiber = self.fiber().clone();
fiber
.assert_can_register()
.map_err(|_| TaskRegistrationError::InactiveContext)?;
let handle = tokio::runtime::Handle::try_current()
.map_err(|_| TaskRegistrationError::ExecutorUnavailable)?;
let cell: Arc<Mutex<RunSlot>> = Arc::new(Mutex::new(RunSlot::Pending));
let parked = Arc::new(tokio::sync::Notify::new());
let (join_cell, join_parked) = (cell.clone(), parked.clone());
let report_logger = self.logger();
let cleanup: Cleanup = crate::effect::fut_cleanup(move || {
Box::pin(async move {
loop {
let notified = join_parked.notified();
tokio::pin!(notified);
notified.as_mut().enable();
let joined = {
let mut slot = join_cell.lock();
match std::mem::replace(&mut *slot, RunSlot::Done) {
RunSlot::Running(join) => Some(join),
RunSlot::Pending => {
*slot = RunSlot::Pending;
None
}
RunSlot::Done => None,
}
};
if let Some(join) = joined {
crate::contained::contain_join("ctx.run task", Some(&report_logger), join)
.await;
return;
}
notified.await;
}
}) as crate::effect::BoxFuture<()>
});
crate::gated::push_gated(&fiber, cleanup, &mut crate::gated::NoPublish)
.map_err(|_| TaskRegistrationError::InactiveContext)?;
let scoped = fiber.clone();
let join: tokio::task::JoinHandle<()> = handle.spawn(async move {
settle_ctx::run_task(scoped, task).await;
});
*cell.lock() = RunSlot::Running(join);
parked.notify_waiters();
Ok(())
}
}
enum RunSlot {
Pending,
Running(tokio::task::JoinHandle<()>),
Done,
}
#[derive(Clone)]
pub struct FiberHandle {
pub(crate) fiber: Arc<Fiber>,
}
impl fmt::Debug for FiberHandle {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("FiberHandle")
.field("id", self.fiber.id())
.field("name", &self.fiber.name)
.finish()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum LifecycleOperation {
Ready,
WaitState,
Restart,
Update,
EraSwap,
Dispose,
RemovePlugins,
}
#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
#[error("{operation:?} cannot wait on {fiber_id:?} from that Fiber's settle context")]
pub struct LifecycleRecursion {
operation: LifecycleOperation,
fiber_id: FiberId,
}
impl LifecycleRecursion {
pub(crate) fn new(operation: LifecycleOperation, fiber_id: FiberId) -> Self {
Self {
operation,
fiber_id,
}
}
pub fn operation(&self) -> LifecycleOperation {
self.operation
}
pub fn fiber_id(&self) -> &FiberId {
&self.fiber_id
}
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum ReadyError {
#[error("{0}")]
Recursion(LifecycleRecursion),
#[error("plugin apply failed: {0}")]
Apply(PluginFailure),
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum RestartError {
#[error("fiber is closed")]
Closed,
#[error("{0}")]
Recursion(LifecycleRecursion),
#[error("plugin apply failed: {0}")]
Apply(PluginFailure),
}
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum WaitStateError {
#[error("wait_state timed out before the fiber published the target state")]
Elapsed,
#[error("{0}")]
Recursion(LifecycleRecursion),
}
impl FiberHandle {
pub(crate) fn new(fiber: Arc<Fiber>) -> Self {
Self { fiber }
}
pub fn state(&self) -> FiberState {
self.fiber.state()
}
pub fn id(&self) -> FiberId {
self.fiber.id().clone()
}
pub fn name(&self) -> &str {
&self.fiber.name
}
pub fn pending_missing(&self) -> Vec<String> {
if !self.fiber.is_alive() || self.fiber.state() != FiberState::Pending {
return Vec::new();
}
match self.fiber.spawn_state.root() {
Some(root) => self.fiber.missing_from(&root),
None => Vec::new(),
}
}
pub async fn ready(&self) -> std::result::Result<FiberState, ReadyError> {
settle_ctx::refuse_recursion(&self.fiber, LifecycleOperation::Ready)
.map_err(ReadyError::Recursion)?;
let fiber = &self.fiber;
let state = loop {
if fiber.slot.has_committed_recheck()
&& let Some(root) = fiber.spawn_state.root()
{
fiber.slot.kick(fiber, &root);
}
fiber.slot.wait_idle().await;
let observed = fiber.state();
if fiber.slot.is_idle() && !fiber.slot.has_committed_recheck() {
break observed;
}
};
match state {
FiberState::Failed => {
let failure = fiber
.error
.lock()
.clone()
.expect("Failed is published only with a parked PluginFailure");
Err(ReadyError::Apply(spawn::clone_owned(&failure)))
}
FiberState::Active | FiberState::Pending | FiberState::Disposed => Ok(state),
FiberState::Loading | FiberState::Unloading => {
unreachable!("an idle Fiber cannot publish a transient lifecycle state")
}
}
}
pub async fn wait_state(
&self,
state: FiberState,
timeout: Duration,
) -> std::result::Result<(), WaitStateError> {
settle_ctx::refuse_recursion(&self.fiber, LifecycleOperation::WaitState)
.map_err(WaitStateError::Recursion)?;
self.fiber.state.wait_for(state, timeout).await
}
pub async fn dispose(&self) -> std::result::Result<(), LifecycleRecursion> {
settle_ctx::refuse_recursion(&self.fiber, LifecycleOperation::Dispose)?;
self.fiber.dispose().await;
Ok(())
}
pub async fn restart(&self) -> std::result::Result<(), RestartError> {
settle_ctx::refuse_recursion(&self.fiber, LifecycleOperation::Restart)
.map_err(RestartError::Recursion)?;
if !self.fiber.is_alive() {
return Err(RestartError::Closed);
}
self.fiber.slot.restart_pass(&self.fiber).await
}
pub async fn update(
&self,
change: crate::PreparedChange,
) -> std::result::Result<UpdateOutcome, UpdateError> {
settle_ctx::refuse_recursion(&self.fiber, LifecycleOperation::Update)
.map_err(UpdateError::Recursion)?;
if !self.fiber.is_alive() {
return Err(UpdateError::Closed);
}
if self.fiber.spawn_state.contract() != Some(change.contract()) {
return Err(UpdateError::PluginContractMismatch);
}
let ctx = self
.fiber
.spawn_state
.required_context_for(&self.fiber)
.map_err(|_| UpdateError::Closed)?;
let accepted = ctx
.root
.updates
.control(change.contract(), ctx.scope.clone(), change)
.await
.map_err(UpdateError::Control)?;
let Some(change) = accepted else {
return Ok(UpdateOutcome::Vetoed);
};
let state = self.fiber.slot.update_pass(&self.fiber, change).await?;
debug_assert!(matches!(state, FiberState::Active | FiberState::Pending));
Ok(UpdateOutcome::Committed(state))
}
}