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