onlyne_client/session/
accept.rs1use crate::session::dispatch::{DispatchState, dispatch};
2use anyhow::{Context, Result, anyhow};
3use onlyne_proto::{Delivery, Envelope, Report};
4use onlyne_session::{SessionBackend, SessionRef};
5use onlyne_store::ClientStore;
6use std::sync::Arc;
7
8#[derive(Clone)]
9pub struct AcceptPath {
10 pub dispatch: DispatchState,
11 pub prose: String,
12}
13
14impl AcceptPath {
15 pub fn new(dispatch: DispatchState, prose: impl Into<String>) -> Self {
16 Self {
17 dispatch,
18 prose: prose.into(),
19 }
20 }
21
22 pub fn accept_new(&self, delivery: &Delivery, accept_new: bool) -> Result<Option<SessionRef>> {
23 if !accept_new {
24 return Ok(None);
25 }
26 delivery
27 .envelope
28 .validate()
29 .map_err(|e| anyhow!(e.to_string()))?;
30 let task_id = delivery
31 .envelope
32 .task_id()
33 .context("delivery missing causality.task")?;
34 let _ = task_id;
35 Ok(Some(dispatch(&self.dispatch, &delivery.envelope)?))
36 }
37
38 pub fn ready_report(
39 &self,
40 task_id: &str,
41 session_id: &str,
42 generation: u64,
43 seq: u64,
44 ) -> Report {
45 Report::Ready {
46 task_id: task_id.into(),
47 session_id: session_id.into(),
48 generation,
49 seq,
50 cluster_ref: None,
51 }
52 }
53}
54
55pub fn validate_delivery(delivery: &Delivery) -> Result<&Envelope> {
56 delivery
57 .envelope
58 .validate()
59 .map_err(|e| anyhow!(e.to_string()))?;
60 Ok(&delivery.envelope)
61}
62
63pub fn resolve_session(
64 path: &AcceptPath,
65 delivery: &Delivery,
66 accept_new: bool,
67) -> Result<Option<SessionRef>> {
68 validate_delivery(delivery)?;
69 path.accept_new(delivery, accept_new)
70}
71
72pub async fn acknowledge_completion(
73 report: &Report,
74 _store: &ClientStore,
75) -> Result<Option<(String, onlyne_proto::Outcome)>> {
76 match report {
77 Report::Complete {
78 task_id, outcome, ..
79 } => Ok(Some((task_id.clone(), *outcome))),
80 _ => Ok(None),
81 }
82}
83
84pub fn backend_for_accept(backend: Arc<dyn SessionBackend>) -> Arc<dyn SessionBackend> {
85 backend
86}