Skip to main content

lxmf_runtime/
backend.rs

1use std::sync::{Arc, Mutex};
2
3use lxmf_sdk::{
4    Ack, CancelResult, ConfigPatch, DeliverySnapshot, DeliveryState, EventBatch, EventCursor,
5    MessageId, NegotiationRequest, NegotiationResponse, RouterStats, RouterStoragePolicy,
6    RouterStoragePolicyPatch, RuntimeSnapshot, SdkBackend, SdkError, SendRequest, Severity,
7    ShutdownMode,
8};
9use rns_transport::hash::AddressHash;
10use serde_json::{json, Value as JsonValue};
11
12use crate::config::InProcessBackendConfig;
13use crate::delivery::{
14    request_link_attempts, request_link_timeout, request_resource_timeout, InProcessSendReport,
15    SendContext,
16};
17use crate::state::{internal_error, BackendState};
18
19#[derive(Clone)]
20pub struct InProcessBackend {
21    config: Arc<InProcessBackendConfig>,
22    state: Arc<Mutex<BackendState>>,
23    propagation_relay: Arc<Mutex<Option<AddressHash>>>,
24    router_storage_policy: Arc<Mutex<RouterStoragePolicy>>,
25}
26
27impl InProcessBackend {
28    pub fn new(config: InProcessBackendConfig) -> Self {
29        let state = BackendState::new(config.runtime_id.clone(), config.limits);
30        let propagation_relay = Arc::new(Mutex::new(config.propagation_relay));
31        Self {
32            config: Arc::new(config),
33            state: Arc::new(Mutex::new(state)),
34            propagation_relay,
35            router_storage_policy: Arc::new(Mutex::new(RouterStoragePolicy::default())),
36        }
37    }
38
39    pub fn send_report(
40        &self,
41        message_id: &MessageId,
42    ) -> Result<Option<InProcessSendReport>, SdkError> {
43        Ok(self
44            .state
45            .lock()
46            .map_err(|_| internal_error("in-process backend state poisoned"))?
47            .send_report(&message_id.0))
48    }
49
50    pub fn set_propagation_relay(&self, relay: Option<AddressHash>) -> Result<(), SdkError> {
51        *self
52            .propagation_relay
53            .lock()
54            .map_err(|_| internal_error("in-process propagation relay state poisoned"))? = relay;
55        Ok(())
56    }
57
58    pub fn record_delivery(
59        &self,
60        message_id: &MessageId,
61        state: DeliveryState,
62        reason: Option<String>,
63    ) -> Result<(), SdkError> {
64        self.state
65            .lock()
66            .map_err(|_| internal_error("in-process backend state poisoned"))?
67            .record_delivery(&message_id.0, state, reason)
68    }
69
70    pub fn record_event(
71        &self,
72        event_type: &str,
73        severity: Severity,
74        payload: JsonValue,
75    ) -> Result<(), SdkError> {
76        self.state
77            .lock()
78            .map_err(|_| internal_error("in-process backend state poisoned"))?
79            .record_event(event_type, severity, payload)
80    }
81}
82
83impl SdkBackend for InProcessBackend {
84    fn negotiate(&self, _req: NegotiationRequest) -> Result<NegotiationResponse, SdkError> {
85        let runtime_id = self
86            .state
87            .lock()
88            .map_err(|_| internal_error("in-process backend state poisoned"))?
89            .runtime_id()
90            .to_owned();
91        serde_json::from_value(json!({
92            "runtime_id": runtime_id,
93            "active_contract_version": 2,
94            "effective_capabilities": [
95                "sdk.capability.event_stream",
96                "sdk.capability.cursor_replay",
97                "sdk.capability.receipt_terminality",
98                "sdk.capability.config_revision_cas",
99                "sdk.capability.idempotency_ttl",
100                "reticulum.capability.raw_bytes",
101                "reticulum.capability.msgpack_fields"
102            ],
103            "effective_limits": {
104                "max_poll_events": 128,
105                "max_event_bytes": 65536,
106                "max_batch_bytes": 1048576,
107                "max_extension_keys": 16,
108                "idempotency_ttl_ms": 43200000
109            },
110            "contract_release": lxmf_sdk::CONTRACT_RELEASE,
111            "schema_namespace": lxmf_sdk::SCHEMA_NAMESPACE,
112        }))
113        .map_err(|err| internal_error(format!("invalid negotiation response: {err}")))
114    }
115
116    fn send(&self, request: SendRequest) -> Result<MessageId, SdkError> {
117        let link_timeout = request_link_timeout(&request, self.config.link_connect_timeout);
118        let link_attempts = request_link_attempts(&request, self.config.link_connect_attempts);
119        let resource_timeout =
120            request_resource_timeout(&request, self.config.resource_transfer_timeout);
121        let context = SendContext {
122            transport: &self.config.transport,
123            identity: &self.config.identity,
124            source_destination: self.config.source_destination,
125            propagation_relay: *self
126                .propagation_relay
127                .lock()
128                .map_err(|_| internal_error("in-process propagation relay state poisoned"))?,
129            link_connect_timeout: link_timeout,
130            link_connect_attempts: link_attempts,
131            resource_transfer_timeout: resource_timeout,
132        };
133        let result = self.config.runtime_handle.block_on(crate::delivery::send(context, &request));
134        match result {
135            Ok(report) => {
136                let message_id = report.message_id.clone();
137                self.state
138                    .lock()
139                    .map_err(|_| internal_error("in-process backend state poisoned"))?
140                    .record_send(report)?;
141                Ok(message_id)
142            }
143            Err(error) => Err(error),
144        }
145    }
146
147    fn cancel(&self, _id: MessageId) -> Result<CancelResult, SdkError> {
148        Ok(CancelResult::Unsupported)
149    }
150
151    fn status(&self, id: MessageId) -> Result<Option<DeliverySnapshot>, SdkError> {
152        Ok(self
153            .state
154            .lock()
155            .map_err(|_| internal_error("in-process backend state poisoned"))?
156            .status(&id))
157    }
158
159    fn configure(&self, _expected_revision: u64, _patch: ConfigPatch) -> Result<Ack, SdkError> {
160        let revision = self
161            .state
162            .lock()
163            .map_err(|_| internal_error("in-process backend state poisoned"))?
164            .advance_config_revision();
165        make_ack(Some(revision))
166    }
167
168    fn poll_events(&self, cursor: Option<EventCursor>, max: usize) -> Result<EventBatch, SdkError> {
169        self.state
170            .lock()
171            .map_err(|_| internal_error("in-process backend state poisoned"))?
172            .poll(cursor.as_ref(), max)
173    }
174
175    fn snapshot(&self) -> Result<RuntimeSnapshot, SdkError> {
176        self.state
177            .lock()
178            .map_err(|_| internal_error("in-process backend state poisoned"))?
179            .snapshot()
180    }
181
182    fn shutdown(&self, _mode: ShutdownMode) -> Result<Ack, SdkError> {
183        make_ack(None)
184    }
185
186    fn router_stats(&self) -> Result<RouterStats, SdkError> {
187        let (messages, outbound_inflight) = self
188            .state
189            .lock()
190            .map_err(|_| internal_error("in-process backend state poisoned"))?
191            .router_counts();
192        let storage_policy = self
193            .router_storage_policy
194            .lock()
195            .map_err(|_| internal_error("in-process router policy state poisoned"))?
196            .clone();
197        Ok(RouterStats { messages, outbound_inflight, storage_policy, ..Default::default() })
198    }
199
200    fn router_storage_policy(&self) -> Result<RouterStoragePolicy, SdkError> {
201        self.router_storage_policy
202            .lock()
203            .map_err(|_| internal_error("in-process router policy state poisoned"))
204            .map(|policy| policy.clone())
205    }
206
207    fn set_router_storage_policy(
208        &self,
209        patch: RouterStoragePolicyPatch,
210    ) -> Result<RouterStoragePolicy, SdkError> {
211        let mut policy = self
212            .router_storage_policy
213            .lock()
214            .map_err(|_| internal_error("in-process router policy state poisoned"))?;
215        if let Some(limit) = patch.message_limit_bytes {
216            policy.message_limit_bytes = Some(limit);
217        }
218        if let Some(limit) = patch.information_limit_bytes {
219            policy.information_limit_bytes = Some(limit);
220        }
221        if let Some(retain) = patch.retain_node_lxms {
222            policy.retain_node_lxms = retain;
223        }
224        Ok(policy.clone())
225    }
226}
227
228fn make_ack(revision: Option<u64>) -> Result<Ack, SdkError> {
229    serde_json::from_value(json!({
230        "accepted": true,
231        "revision": revision,
232    }))
233    .map_err(|err| internal_error(format!("invalid acknowledgement: {err}")))
234}