Skip to main content

link_cli/protocol/
remote_links.rs

1//! [`RemoteLinks`]: a links store that lives behind a [`LinksServer`](super::LinksServer).
2//!
3//! It implements the same interfaces as a local store — [`doublets::Links`],
4//! [`doublets::Doublets`] and [`NamedTypeLinks`] — so code written against
5//! them, the [`QueryProcessor`](crate::QueryProcessor) included, switches
6//! from a local file to a server by swapping one value:
7//!
8//! ```no_run
9//! use link_cli::protocol::{RemoteLinks, TextLinoProtocol};
10//! use link_cli::{NamedTypeLinks, NamedTypesDecorator, QueryProcessor};
11//!
12//! fn run(store: &mut impl NamedTypeLinks) -> anyhow::Result<()> {
13//!     QueryProcessor::new(false).process_query(store, "() ((1 1))")?;
14//!     Ok(())
15//! }
16//!
17//! run(&mut NamedTypesDecorator::new("db.links", false)?)?;
18//! run(&mut RemoteLinks::connect("127.0.0.1:8080", TextLinoProtocol::new())?)?;
19//! # Ok::<(), anyhow::Error>(())
20//! ```
21//!
22//! Every call is one [`LinksOperation`] round trip. An observer of a write
23//! sees the net change of every link it touched, not the steps the store
24//! behind the server took. The methods of the upstream and CLI traits that
25//! cannot return an error panic when the connection fails, the way a local
26//! store panics when its file does; the inherent methods return every failure
27//! as a [`ProtocolError`].
28
29use 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/// A links store served by a [`LinksServer`](super::LinksServer), usable
47/// wherever a local store is.
48#[derive(Debug)]
49pub struct RemoteLinks {
50    client: Mutex<LinksClient>,
51}
52
53impl RemoteLinks {
54    /// Connects to the server at `address`.
55    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    /// Uses an existing connection.
63    pub fn new(client: LinksClient) -> Self {
64        Self {
65            client: Mutex::new(client),
66        }
67    }
68
69    /// Runs one operation on the server and returns its reply.
70    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    /// Number of links matching `restriction`.
78    pub fn count_matching(&self, restriction: &[Part]) -> ProtocolResult<u32> {
79        parse_count(&self.execute(&LinksOperation::Count(restriction.to_vec()))?)
80    }
81
82    /// Links matching `restriction`, ordered by index.
83    pub fn links_matching(&self, restriction: &[Part]) -> ProtocolResult<Vec<Link>> {
84        parse_links(&self.execute(&LinksOperation::Each(restriction.to_vec()))?)
85    }
86
87    /// Runs a `create`, `update` or `delete` and returns its changes.
88    pub fn changes(&self, operation: &LinksOperation) -> ProtocolResult<Vec<Change>> {
89        parse_changes(&self.execute(operation)?)
90    }
91
92    /// The link stored at `index`.
93    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    /// Runs an `update` or `delete` of `index`, reports its changes to
105    /// `observer` and returns the state `index` was in before.
106    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
134/// Unwraps the result of a call whose interface has no way to report a
135/// failure.
136fn 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
148/// Feeds `changes` to `handler` until it asks to stop.
149fn 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    /// Creates an empty `(index: 0 0)` link, like every `doublets` store.
172    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    /// Asks for links until the server hands out `id`, then deletes the ones
235    /// handed out on the way, exactly like [`LinkStorage::ensure_created`](crate::LinkStorage::ensure_created).
236    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    /// The server saves after every change, so there is nothing to flush.
319    fn save(&mut self) -> anyhow::Result<()> {
320        Ok(())
321    }
322}