Skip to main content

onlyne_client/session/
accept.rs

1use 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}