liminal_server/server/participant/
dispatch.rs1use std::collections::BTreeSet;
8use std::sync::Arc;
9
10use liminal::durability::DurableStore;
11use liminal::protocol::Frame;
12use liminal_protocol::lifecycle::ConnectionConversationTracking;
13use liminal_protocol::wire::{
14 BindingEpoch, ClientRequest, CodecError, ConnectionIncarnation, ConversationId,
15 ObserverRecoveryHandshake, ParticipantId, ServerValue, ValidatedFrameLimit,
16};
17
18use super::transport::{
19 ParticipantIngress, ParticipantSession, encode_server_value, gate_generic_frame,
20 normalize_configured_frame_limit,
21};
22use super::{
23 ObserverPublicationTarget, ParticipantOfferedProgress, ParticipantPublication,
24 ParticipantPublicationInbox, ParticipantPublicationRegistry,
25};
26
27#[derive(Debug, Default)]
40pub struct ParticipantConnectionConversations {
41 tracked: BTreeSet<ConversationId>,
42}
43
44impl ParticipantConnectionConversations {
45 #[must_use]
47 pub fn tracking(&self, conversation_id: ConversationId) -> ConnectionConversationTracking {
48 if self.tracked.contains(&conversation_id) {
49 ConnectionConversationTracking::AlreadyTracked
50 } else {
51 ConnectionConversationTracking::Untracked
52 }
53 }
54
55 #[must_use]
57 pub fn occupied(&self) -> u64 {
58 u64::try_from(self.tracked.len()).unwrap_or(u64::MAX)
62 }
63
64 pub fn track(&mut self, conversation_id: ConversationId) {
66 self.tracked.insert(conversation_id);
67 }
68
69 #[must_use]
72 pub fn tracked_conversations(&self) -> Vec<ConversationId> {
73 self.tracked.iter().copied().collect()
74 }
75}
76
77#[derive(Clone, Copy, Debug, PartialEq, Eq)]
79pub struct ParticipantConnectionContext {
80 connection_incarnation: ConnectionIncarnation,
81}
82
83impl ParticipantConnectionContext {
84 #[must_use]
86 pub const fn new(connection_incarnation: ConnectionIncarnation) -> Self {
87 Self {
88 connection_incarnation,
89 }
90 }
91
92 #[must_use]
94 pub const fn connection_incarnation(self) -> ConnectionIncarnation {
95 self.connection_incarnation
96 }
97}
98
99#[derive(Clone, Debug, thiserror::Error, PartialEq, Eq)]
105pub enum ParticipantSemanticError {
106 #[error("participant semantic service is unavailable")]
108 Unavailable,
109 #[error("participant semantic service failed: {message}")]
111 Internal {
112 message: String,
114 },
115}
116
117pub trait ParticipantSemanticHandler: core::fmt::Debug + Send + Sync {
119 fn publication_conversation_limit(&self) -> u64 {
135 0
136 }
137
138 fn ready_connection_incarnations(
146 &self,
147 _conversation_id: ConversationId,
148 ) -> Result<Vec<ConnectionIncarnation>, ParticipantSemanticError> {
149 Ok(Vec::new())
150 }
151
152 fn next_publication(
159 &self,
160 _connection_incarnation: ConnectionIncarnation,
161 _conversation_id: ConversationId,
162 _offered: Option<ParticipantOfferedProgress>,
163 ) -> Result<Option<ParticipantPublication>, ParticipantSemanticError> {
164 Ok(None)
165 }
166
167 fn publication_binding_is_current(
174 &self,
175 _conversation_id: ConversationId,
176 _participant_id: ParticipantId,
177 _binding_epoch: BindingEpoch,
178 ) -> Result<bool, ParticipantSemanticError> {
179 Ok(false)
180 }
181
182 fn record_publication_offer(
189 &self,
190 _publication: &ParticipantPublication,
191 ) -> Result<(), ParticipantSemanticError> {
192 Ok(())
193 }
194
195 fn handle_observer_recovery(
204 &self,
205 context: ParticipantConnectionContext,
206 conversations: &mut ParticipantConnectionConversations,
207 request: ObserverRecoveryHandshake,
208 target: Option<ObserverPublicationTarget>,
209 ) -> Result<ServerValue, ParticipantSemanticError> {
210 drop(target);
211 self.handle(
212 context,
213 conversations,
214 ClientRequest::ObserverRecovery(request),
215 )
216 }
217
218 fn handle(
226 &self,
227 context: ParticipantConnectionContext,
228 conversations: &mut ParticipantConnectionConversations,
229 request: ClientRequest,
230 ) -> Result<ServerValue, ParticipantSemanticError>;
231}
232
233#[derive(Clone, Debug)]
249pub struct InstalledParticipantService {
250 handler: Arc<dyn ParticipantSemanticHandler>,
251 durable_store: Arc<dyn DurableStore>,
252 frame_limit: ValidatedFrameLimit,
253 publication_registry: Arc<ParticipantPublicationRegistry>,
254}
255
256impl InstalledParticipantService {
257 pub(crate) fn new(
269 handler: Arc<dyn ParticipantSemanticHandler>,
270 durable_store: Arc<dyn DurableStore>,
271 configured_wf: u64,
272 ) -> Result<Self, CodecError> {
273 Ok(Self {
274 handler,
275 durable_store,
276 frame_limit: normalize_configured_frame_limit(configured_wf)?,
277 publication_registry: Arc::new(ParticipantPublicationRegistry::default()),
278 })
279 }
280
281 #[must_use]
283 pub(crate) fn durable_store(&self) -> Arc<dyn DurableStore> {
284 Arc::clone(&self.durable_store)
285 }
286
287 #[must_use]
290 pub(crate) const fn frame_limit(&self) -> ValidatedFrameLimit {
291 self.frame_limit
292 }
293
294 #[must_use]
297 pub(crate) fn publication_conversation_limit(&self) -> u64 {
298 self.handler.publication_conversation_limit()
299 }
300
301 #[must_use]
303 pub(crate) fn new_publication_inbox(&self) -> ParticipantPublicationInbox {
304 ParticipantPublicationInbox::new(self.handler.publication_conversation_limit())
305 }
306
307 #[must_use]
309 pub(crate) fn publication_registry(&self) -> &ParticipantPublicationRegistry {
310 &self.publication_registry
311 }
312
313 pub(crate) fn next_publication(
316 &self,
317 connection_incarnation: ConnectionIncarnation,
318 conversation_id: ConversationId,
319 offered: Option<ParticipantOfferedProgress>,
320 ) -> Result<Option<ParticipantPublication>, ParticipantSemanticError> {
321 self.handler
322 .next_publication(connection_incarnation, conversation_id, offered)
323 }
324
325 pub(crate) fn publication_binding_is_current(
326 &self,
327 conversation_id: ConversationId,
328 participant_id: ParticipantId,
329 binding_epoch: BindingEpoch,
330 ) -> Result<bool, ParticipantSemanticError> {
331 self.handler
332 .publication_binding_is_current(conversation_id, participant_id, binding_epoch)
333 }
334
335 pub(crate) fn record_publication_offer(
336 &self,
337 publication: &ParticipantPublication,
338 ) -> Result<(), ParticipantSemanticError> {
339 self.handler.record_publication_offer(publication)
340 }
341
342 fn notify_ready(
343 &self,
344 conversation_id: ConversationId,
345 ) -> Result<(), ParticipantSemanticError> {
346 for incarnation in self
347 .handler
348 .ready_connection_incarnations(conversation_id)?
349 {
350 self.publication_registry
351 .notify(incarnation, conversation_id)
352 .map_err(|error| ParticipantSemanticError::Internal {
353 message: format!("participant publication wake failed: {error}"),
354 })?;
355 }
356 Ok(())
357 }
358}
359
360impl ParticipantSemanticHandler for InstalledParticipantService {
361 fn publication_conversation_limit(&self) -> u64 {
362 self.handler.publication_conversation_limit()
363 }
364
365 fn handle(
366 &self,
367 context: ParticipantConnectionContext,
368 conversations: &mut ParticipantConnectionConversations,
369 request: ClientRequest,
370 ) -> Result<ServerValue, ParticipantSemanticError> {
371 if let ClientRequest::ObserverRecovery(request) = request {
372 let target = self
373 .publication_registry
374 .observer_target(context.connection_incarnation())
375 .map_err(|error| ParticipantSemanticError::Internal {
376 message: format!("observer publication target failed: {error}"),
377 })?;
378 return self
379 .handler
380 .handle_observer_recovery(context, conversations, request, target);
381 }
382 let conversation_id = request_conversation_id(&request);
383 let value = self.handler.handle(context, conversations, request)?;
384 if let Some(conversation_id) = conversation_id {
385 self.notify_ready(conversation_id)?;
386 }
387 Ok(value)
388 }
389}
390
391const fn request_conversation_id(request: &ClientRequest) -> Option<ConversationId> {
392 match request {
393 ClientRequest::Enrollment(request) => Some(request.conversation_id),
394 ClientRequest::CredentialAttach(request) => Some(request.conversation_id),
395 ClientRequest::Detach(request) => Some(request.conversation_id),
396 ClientRequest::ParticipantAck(request) => Some(request.conversation_id),
397 ClientRequest::Leave(request) => Some(request.conversation_id),
398 ClientRequest::MarkerAck(request) => Some(request.conversation_id),
399 ClientRequest::RecordAdmission(request) => Some(request.conversation_id),
400 ClientRequest::ObserverRecovery(_) => None,
401 }
402}
403
404#[derive(Debug)]
406pub enum ParticipantDispatch {
407 NotParticipant,
409 Respond(Frame),
411 RespondThenClose(Frame),
413 Fatal(ParticipantDispatchError),
415}
416
417#[derive(Debug, thiserror::Error)]
419pub enum ParticipantDispatchError {
420 #[error("invalid generic participant frame")]
422 InvalidGenericFrame,
423 #[error(transparent)]
425 Semantic(#[from] ParticipantSemanticError),
426 #[error("failed to encode participant response: {0:?}")]
428 Encode(CodecError),
429}
430
431#[must_use]
436pub fn dispatch_generic_frame(
437 frame: &Frame,
438 authenticated: bool,
439 session: ParticipantSession,
440 context: ParticipantConnectionContext,
441 conversations: &mut ParticipantConnectionConversations,
442 handler: &dyn ParticipantSemanticHandler,
443) -> ParticipantDispatch {
444 let (value, close_after_response) = match gate_generic_frame(frame, authenticated, session) {
445 ParticipantIngress::NotParticipant => return ParticipantDispatch::NotParticipant,
446 ParticipantIngress::Rejected(rejection) => {
447 (ServerValue::ParticipantTransportRejected(rejection), true)
448 }
449 ParticipantIngress::InvalidGenericFrame => {
450 return ParticipantDispatch::Fatal(ParticipantDispatchError::InvalidGenericFrame);
451 }
452 ParticipantIngress::Request(request) => {
453 match handler.handle(context, conversations, request) {
454 Ok(value) => (value, false),
455 Err(error) => {
456 return ParticipantDispatch::Fatal(ParticipantDispatchError::Semantic(error));
457 }
458 }
459 }
460 };
461 match encode_server_value(value) {
462 Ok(frame) if close_after_response => ParticipantDispatch::RespondThenClose(frame),
463 Ok(frame) => ParticipantDispatch::Respond(frame),
464 Err(error) => ParticipantDispatch::Fatal(ParticipantDispatchError::Encode(error)),
465 }
466}