1use super::client::LinksClient;
30use super::error::{ProtocolError, ProtocolResult};
31use super::links_operations::{
32 parse_changes, parse_count, parse_link_reply, parse_links, parse_name, Change, LinksOperation,
33 Part,
34};
35use super::mapping::LinoDocument;
36use super::protocols::LinoProtocol;
37use crate::link::Link;
38use crate::link_storage::ChangeObserver;
39use crate::link_storage_doublets::link_storage_constants;
40use crate::named_type_links::NamedTypeLinks;
41use doublets::data::{Flow, LinksConstants, ReadHandler, WriteHandler};
42use doublets::{Doublets, Error, Link as DoubletsLink, Links};
43use std::net::ToSocketAddrs;
44use std::sync::{Mutex, PoisonError};
45
46#[derive(Debug)]
49pub struct RemoteLinks {
50 client: Mutex<LinksClient>,
51}
52
53impl RemoteLinks {
54 pub fn connect(
56 address: impl ToSocketAddrs,
57 protocol: impl LinoProtocol + 'static,
58 ) -> ProtocolResult<Self> {
59 Ok(Self::new(LinksClient::connect(address, protocol)?))
60 }
61
62 pub fn new(client: LinksClient) -> Self {
64 Self {
65 client: Mutex::new(client),
66 }
67 }
68
69 pub fn execute(&self, operation: &LinksOperation) -> ProtocolResult<LinoDocument> {
71 self.client
72 .lock()
73 .unwrap_or_else(PoisonError::into_inner)
74 .request(&operation.to_document())
75 }
76
77 pub fn count_matching(&self, restriction: &[Part]) -> ProtocolResult<u32> {
79 parse_count(&self.execute(&LinksOperation::Count(restriction.to_vec()))?)
80 }
81
82 pub fn links_matching(&self, restriction: &[Part]) -> ProtocolResult<Vec<Link>> {
84 parse_links(&self.execute(&LinksOperation::Each(restriction.to_vec()))?)
85 }
86
87 pub fn changes(&self, operation: &LinksOperation) -> ProtocolResult<Vec<Change>> {
89 parse_changes(&self.execute(operation)?)
90 }
91
92 pub fn link(&self, index: u32) -> ProtocolResult<Option<Link>> {
94 Ok(self.links_matching(&[Some(index)])?.pop())
95 }
96
97 fn create_remote(&self, source: u32, target: u32) -> ProtocolResult<Link> {
98 self.changes(&LinksOperation::Create { source, target })?
99 .into_iter()
100 .find_map(|(before, after)| (before.is_null() && !after.is_null()).then_some(after))
101 .ok_or_else(|| ProtocolError::malformed("the create reply holds no created link"))
102 }
103
104 fn observed(
107 &self,
108 operation: &LinksOperation,
109 index: u32,
110 observer: ChangeObserver<'_>,
111 ) -> ProtocolResult<Link> {
112 let changes = self.changes(operation)?;
113 for (before, after) in &changes {
114 observer(*before, *after);
115 }
116 changes
117 .iter()
118 .find(|(before, _)| before.index == index)
119 .map(|(before, _)| *before)
120 .ok_or_else(|| {
121 ProtocolError::malformed(format!("the reply holds no change of {index}"))
122 })
123 }
124
125 fn parts(&self, query: &[u32]) -> Vec<Part> {
126 let any = self.constants().any;
127 query
128 .iter()
129 .map(|&part| (part != any).then_some(part))
130 .collect()
131 }
132}
133
134fn connected<T>(result: ProtocolResult<T>) -> T {
137 result.unwrap_or_else(|error| panic!("remote links store failed: {error}"))
138}
139
140fn doublets_link(link: &Link) -> DoubletsLink<u32> {
141 DoubletsLink::new(link.index, link.source, link.target)
142}
143
144fn doublets_error(error: ProtocolError) -> Error<u32> {
145 Error::Other(Box::new(error))
146}
147
148fn replay(changes: &[Change], handler: WriteHandler<'_, u32>) -> Flow {
150 for (before, after) in changes {
151 if handler(doublets_link(before), doublets_link(after)) == Flow::Break {
152 return Flow::Break;
153 }
154 }
155 Flow::Continue
156}
157
158fn part(query: &[u32], index: usize) -> u32 {
159 query.get(index).copied().unwrap_or(0)
160}
161
162impl Links<u32> for RemoteLinks {
163 fn constants(&self) -> &LinksConstants<u32> {
164 link_storage_constants()
165 }
166
167 fn count_links(&self, query: &[u32]) -> u32 {
168 connected(self.count_matching(&self.parts(query)))
169 }
170
171 fn create_links(
173 &mut self,
174 _query: &[u32],
175 handler: WriteHandler<'_, u32>,
176 ) -> Result<Flow, Error<u32>> {
177 let changes = self
178 .changes(&LinksOperation::Create {
179 source: 0,
180 target: 0,
181 })
182 .map_err(doublets_error)?;
183 Ok(replay(&changes, handler))
184 }
185
186 fn each_links(&self, query: &[u32], handler: ReadHandler<'_, u32>) -> Flow {
187 for link in connected(self.links_matching(&self.parts(query))) {
188 if handler(doublets_link(&link)) == Flow::Break {
189 return Flow::Break;
190 }
191 }
192 Flow::Continue
193 }
194
195 fn update_links(
196 &mut self,
197 query: &[u32],
198 change: &[u32],
199 handler: WriteHandler<'_, u32>,
200 ) -> Result<Flow, Error<u32>> {
201 let changes = self
202 .changes(&LinksOperation::Update {
203 index: part(query, 0),
204 source: part(change, 1),
205 target: part(change, 2),
206 })
207 .map_err(doublets_error)?;
208 Ok(replay(&changes, handler))
209 }
210
211 fn delete_links(
212 &mut self,
213 query: &[u32],
214 handler: WriteHandler<'_, u32>,
215 ) -> Result<Flow, Error<u32>> {
216 let changes = self
217 .changes(&LinksOperation::Delete(part(query, 0)))
218 .map_err(doublets_error)?;
219 Ok(replay(&changes, handler))
220 }
221}
222
223impl Doublets<u32> for RemoteLinks {
224 fn get_link(&self, index: u32) -> Option<DoubletsLink<u32>> {
225 connected(self.link(index)).as_ref().map(doublets_link)
226 }
227}
228
229impl NamedTypeLinks for RemoteLinks {
230 fn create(&mut self, source: u32, target: u32) -> u32 {
231 connected(self.create_remote(source, target)).index
232 }
233
234 fn ensure_created(&mut self, id: u32) -> u32 {
237 if id == 0 || self.exists(id) {
238 return id;
239 }
240 let mut passed_over = Vec::new();
241 loop {
242 let created = NamedTypeLinks::create(self, 0, 0);
243 if created == id {
244 break;
245 }
246 passed_over.push(created);
247 }
248 for address in passed_over {
249 connected(self.changes(&LinksOperation::Delete(address)));
250 }
251 id
252 }
253
254 fn get_link(&mut self, id: u32) -> Option<Link> {
255 connected(self.link(id))
256 }
257
258 fn exists(&mut self, id: u32) -> bool {
259 connected(self.link(id)).is_some()
260 }
261
262 fn update_observed(
263 &mut self,
264 id: u32,
265 source: u32,
266 target: u32,
267 observer: ChangeObserver<'_>,
268 ) -> anyhow::Result<Link> {
269 let operation = LinksOperation::Update {
270 index: id,
271 source,
272 target,
273 };
274 Ok(self.observed(&operation, id, observer)?)
275 }
276
277 fn delete_observed(&mut self, id: u32, observer: ChangeObserver<'_>) -> anyhow::Result<Link> {
278 Ok(self.observed(&LinksOperation::Delete(id), id, observer)?)
279 }
280
281 fn all_links(&mut self) -> Vec<Link> {
282 connected(self.links_matching(&[]))
283 }
284
285 fn search(&mut self, source: u32, target: u32) -> Option<u32> {
286 connected(self.links_matching(&[None, Some(source), Some(target)]))
287 .first()
288 .map(|link| link.index)
289 }
290
291 fn get_or_create(&mut self, source: u32, target: u32) -> u32 {
292 match self.search(source, target) {
293 Some(index) => index,
294 None => NamedTypeLinks::create(self, source, target),
295 }
296 }
297
298 fn get_name(&mut self, id: u32) -> anyhow::Result<Option<String>> {
299 Ok(parse_name(&self.execute(&LinksOperation::GetName(id))?)?)
300 }
301
302 fn set_name(&mut self, id: u32, name: &str) -> anyhow::Result<u32> {
303 parse_link_reply(&self.execute(&LinksOperation::SetName(id, name.to_string()))?)?
304 .ok_or_else(|| anyhow::anyhow!("the set-name reply holds no link"))
305 }
306
307 fn get_by_name(&mut self, name: &str) -> anyhow::Result<Option<u32>> {
308 Ok(parse_link_reply(
309 &self.execute(&LinksOperation::GetByName(name.to_string()))?,
310 )?)
311 }
312
313 fn remove_name(&mut self, id: u32) -> anyhow::Result<()> {
314 self.execute(&LinksOperation::RemoveName(id))?;
315 Ok(())
316 }
317
318 fn save(&mut self) -> anyhow::Result<()> {
320 Ok(())
321 }
322}