pub struct ObservingTaskRegistry { /* private fields */ }Expand description
A SessionTaskRegistry decorator that fans real task transitions out to
registered TaskTransitionObservers.
This is the reusable, storage-agnostic form of the fan-out the server’s
DbSessionTaskRegistry performs inline (EVE-729): it wraps any inner
registry (in-memory, SQLite, gRPC) so an embedder — e.g. everruns-host
with a crate::wake_queue::SessionWakeQueue — gets the same transition
notifications without depending on the control-plane server.
Transition detection mirrors DbSessionTaskRegistry exactly so mid-turn and
between-turn delivery agree on when a wake fires:
Terminal— fired once when an update moves a non-terminal task into a terminal state, gated on this update’s own intent (an update that does not set a terminal state never fires it, so a racing heartbeat cannot wake on another writer’s transition).AwaitingInput— fired once on the transition intoawaiting_input.Message— fired for each outbound message.
Observers are awaited in registration order (fast, in-process consumers); a failing observer is logged and never fails the underlying task op.
Implementations§
Source§impl ObservingTaskRegistry
impl ObservingTaskRegistry
pub fn new(inner: Arc<dyn SessionTaskRegistry>) -> Self
Sourcepub fn with_observer(self, observer: Arc<dyn TaskTransitionObserver>) -> Self
pub fn with_observer(self, observer: Arc<dyn TaskTransitionObserver>) -> Self
Register an observer to receive every real transition.
Sourcepub fn has_observers(&self) -> bool
pub fn has_observers(&self) -> bool
Whether any observer is registered (fan-out is otherwise a no-op).
Trait Implementations§
Source§impl SessionTaskRegistry for ObservingTaskRegistry
impl SessionTaskRegistry for ObservingTaskRegistry
Source§fn create<'life0, 'async_trait>(
&'life0 self,
input: CreateSessionTask,
) -> Pin<Box<dyn Future<Output = Result<SessionTask>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn create<'life0, 'async_trait>(
&'life0 self,
input: CreateSessionTask,
) -> Pin<Box<dyn Future<Output = Result<SessionTask>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Source§fn update<'life0, 'life1, 'async_trait>(
&'life0 self,
session_id: SessionId,
task_id: &'life1 str,
update: SessionTaskUpdate,
) -> Pin<Box<dyn Future<Output = Result<Option<SessionTask>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn update<'life0, 'life1, 'async_trait>(
&'life0 self,
session_id: SessionId,
task_id: &'life1 str,
update: SessionTaskUpdate,
) -> Pin<Box<dyn Future<Output = Result<Option<SessionTask>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
apply_task_update invariants.fn get<'life0, 'life1, 'async_trait>(
&'life0 self,
session_id: SessionId,
task_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<SessionTask>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn list<'life0, 'life1, 'async_trait>(
&'life0 self,
session_id: SessionId,
filter: Option<&'life1 SessionTaskFilter>,
) -> Pin<Box<dyn Future<Output = Result<Vec<SessionTask>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Source§fn request_cancel<'life0, 'life1, 'async_trait>(
&'life0 self,
session_id: SessionId,
task_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<SessionTask>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn request_cancel<'life0, 'life1, 'async_trait>(
&'life0 self,
session_id: SessionId,
task_id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<SessionTask>>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Source§fn record_message<'life0, 'life1, 'async_trait>(
&'life0 self,
session_id: SessionId,
task_id: &'life1 str,
message: NewTaskMessage,
) -> Pin<Box<dyn Future<Output = Result<TaskMessage>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn record_message<'life0, 'life1, 'async_trait>(
&'life0 self,
session_id: SessionId,
task_id: &'life1 str,
message: NewTaskMessage,
) -> Pin<Box<dyn Future<Output = Result<TaskMessage>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
in_reply_to set) clear a matching pending input request and return
the task to running.