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}