Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/comm/response.rs

20.0 KiB, 121 runs

created by r1870400018:771, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1use crate::{
2 prelude::*,
3 base::id::{
4 self,
5 OzoneBotId,
6 Ticket,
7 },
8 comm::msg::OzoneMsg,
9 data::{
10 core::Value,
11 },
12};
13
14use oxedyne_fe2o3_core::{
15 alt::Override,
16 channels::{
17 Recv,
18 Simplex,
19 },
20};
21use oxedyne_fe2o3_crypto::enc::EncryptionScheme;
22use oxedyne_fe2o3_iop_crypto::enc::EncrypterDefAlt;
23use oxedyne_fe2o3_iop_db::api::Meta;
24use oxedyne_fe2o3_jdat::{
25 prelude::*,
26 id::NumIdDat,
27};
28
29use std::{
30 collections::HashSet,
31 time::{
32 Duration,
33 Instant,
34 },
35};
36
37#[derive(Clone, Debug)]
38pub struct Responder<
39 const UIDL: usize,
40 UID: NumIdDat<UIDL>,
41 ENC: Encrypter,
42 KH: Hasher,
43> {
44 ozid: Option<OzoneBotId>, // Source
45 tik: Ticket,
46 chan: Option<Simplex<OzoneMsg<UIDL, UID, ENC, KH>>>,
47}
48
49impl<
50 const UIDL: usize,
51 UID: NumIdDat<UIDL> + 'static,
52 ENC: Encrypter + 'static,
53 KH: Hasher + 'static,
54>
55 Responder<UIDL, UID, ENC, KH>
56{
57 pub fn new(ozid: Option<&OzoneBotId>) -> Self {
58 Self {
59 ozid: ozid.map(|id| id.clone()),
60 tik: Ticket::new(),
61 chan: Some(Simplex::default()),
62 }
63 }
64
65 pub fn make(
66 ozid: Option<&OzoneBotId>,
67 chan: Option<Simplex<OzoneMsg<UIDL, UID, ENC, KH>>>,
68 )
69 -> Self
70 {
71 Self {
72 ozid: ozid.map(|id| id.clone()),
73 tik: Ticket::new(),
74 chan: chan,
75 }
76 }
77
78 pub fn none(ozid: Option<&OzoneBotId>) -> Self {
79 Self {
80 ozid: ozid.map(|id| id.clone()),
81 tik: Ticket::new(),
82 chan: None,
83 }
84 }
85
86 pub fn is_none(&self) -> bool { self.chan.is_none() }
87 pub fn is_some(&self) -> bool { self.chan.is_some() }
88 pub fn ozid(&self) -> &Option<OzoneBotId> { &self.ozid }
89 pub fn ticket(&self) -> &Ticket { &self.tik }
90 pub fn channel(&self) -> Option<&Simplex<OzoneMsg<UIDL, UID, ENC, KH>>> { self.chan.as_ref() }
91
92 pub fn recv_block(&self) -> Outcome<OzoneMsg<UIDL, UID, ENC, KH>> {
93 match self.channel() {
94 None => Err(err!("This responder does not have a channel."; Channel, Missing)),
95 Some(ref simplex) => {
96 match simplex.recv() {
97 Err(e) => return Err(err!(e,
98 "Could not read from responder channel";
99 Channel, Read)),
100 Ok(OzoneMsg::Error(e)) => return Err(e),
101 Ok(msg) => Ok(msg),
102 }
103 },
104 }
105 }
106
107 pub fn send(&self, msg: OzoneMsg<UIDL, UID, ENC, KH>) -> Outcome<()> {
108 match self.channel() {
109 Some(chan) => chan.send(msg),
110 None => return Err(err!("This responder has no channel."; Channel, Missing)),
111 }
112 }
113
114 /// A receiver waiting for a complete `Dat` wrapped byte vector. Also returns whether
115 /// garbage collection has just been performed on the read file, during which time it is
116 /// possible the value may have been updated. This method does not assemble a `Dat`
117 /// wrapped byte vector from chunks, use `db::fetch_chunks` for that.
118 pub fn recv_daticle(
119 &self,
120 enc: &EncrypterDefAlt<EncryptionScheme, ENC>,
121 or: Option<&Override<EncryptionScheme, ENC>>,
122 )
123 -> Outcome<(Option<(Dat, Meta<UIDL, UID>)>, bool)>
124 {
125 match self.channel() {
126 None => Err(err!("This responder does not have a channel."; Channel, Missing)),
127 Some(ref simplex) => {
128 match simplex.recv_timeout(constant::USER_REQUEST_TIMEOUT) {
129 Recv::Empty => Err(err!(
130 "Failed to receive a message via responder within {:.2} [s].",
131 constant::USER_REQUEST_TIMEOUT.as_secs_f32();
132 Missing, Data)),
133 Recv::Result(Err(e)) => return Err(err!(e,
134 "Could not read from responder channel";
135 Channel, Read)),
136 Recv::Result(Ok(msg)) => match msg {
137 OzoneMsg::Error(e) => return Err(e),
138 // A single record is decoded the same whether it was read under a Complete
139 // key or a chunk key: the chunk arm is reached when a chunk-data key is
140 // fetched directly -- which an orphan sweep does to read a chunk record's
141 // tombstone, and which also arises when a deleted chunked value's chunk key
142 // is scanned (its newest record is a deleted-kind marker and scans as a
143 // Tup5u64 main key). Recognising the marker in both arms, not just the
144 // Complete one, is what lets such a key read back as absent rather than as
145 // an "unexpected" error.
146 OzoneMsg::Value(Value::Complete(Some((dat, meta)), postgc)) |
147 OzoneMsg::Value(Value::Chunk(Some((dat, meta)), _, postgc)) => {
148 // A deletion is stored as a tombstone value under the deleted key, and
149 // the reader finds it exactly as it finds any other value. A key whose
150 // newest value is a tombstone has no value, so say so here, once, for
151 // every caller: a deleted key reads as absent, indistinguishable from
152 // one that was never written. The tombstone is deliberately left
153 // unencrypted so that it can be recognised without a key.
154 if let Dat::Usr(kind, _) = &dat {
155 if *kind == id::usr_kind_id_deleted() {
156 return Ok((None, postgc));
157 }
158 }
159 let or_is_some = match or {
160 Some(or) => or.is_some(),
161 None => false,
162 };
163 if enc.is_none() && !or_is_some {
164 return Ok((Some((dat, meta)), postgc));
165 }
166 let val = try_extract_dat!(dat, BU8, BU16, BU32, BU64);
167 let plain = res!(enc.or_decrypt(&val, or));
168 match Dat::from_bytes(&plain) {
169 Err(e) => return Err(err!(e,
170 "Could not form a Dat from the value bytes, \
171 this could be due to the use of an encryption scheme \
172 differing from the one provided ({}).", enc.or_debug(or);
173 Decode, Bytes)),
174 Ok((dat, _)) => return Ok((Some((dat, meta)), postgc)),
175 }
176 },
177 OzoneMsg::Value(Value::Complete(None, _)) |
178 OzoneMsg::Value(Value::Chunk(None, ..)) => Ok((None, false)),
179 msg => return Err(err!(
180 "Expected a OzoneMsg::Value containing a Value::Complete \
181 wrapping a Dat::BU64 but received a {:?}.", msg;
182 Unexpected, Input)),
183 }
184 }
185 },
186 }
187 }
188
189 /// Collect one and only one reply within a given time, otherwise return an error.
190 pub fn recv_timeout(&self, timeout: Duration) -> Outcome<OzoneMsg<UIDL, UID, ENC, KH>> {
191 match self.channel() {
192 None => Err(err!("This responder does not have a channel."; Channel, Missing)),
193 Some(simplex) => {
194 match simplex.recv_timeout(timeout) {
195 Recv::Empty => Err(err!(
196 "Failed to receive a message via responder within {:.2} [s].",
197 timeout.as_secs_f32();
198 Missing, Data)),
199 Recv::Result(Err(e)) => Err(err!(e,
200 "Could not read from responder channel.";
201 Channel, Read)),
202 Recv::Result(Ok(msg)) => Ok(msg),
203 }
204 },
205 }
206 }
207
208 /// Waits for the answers to a store dispatched with this responder: the count of records it
209 /// was split into, then each record confirmed written and all of them confirmed durable.
210 /// Returns whether the key already held a value, and the number of records.
211 pub fn recv_store_ack(&self) -> Outcome<(bool, usize)> {
212 // The count is sent before anything is dispatched, so it is waiting already.
213 let n = match res!(self.recv_timeout(constant::USER_REQUEST_TIMEOUT)) {
214 OzoneMsg::Chunks(n) => n,
215 OzoneMsg::Error(e) => return Err(e),
216 msg => return Err(err!(
217 "Expected an OzoneMsg::Chunks counting the records of a store, received {:?}.", msg;
218 Channel, Unexpected)),
219 };
220 // Passed on as it is, not wrapped: its tags are the caller's only way to tell a write
221 // that landed unconfirmed (`Unconfirmed`) from one that did not, and a wrapper's tags are
222 // all that `Error::tags` reports.
223 let acks = ok!(self.recv_write_acks(
224 n,
225 constant::USER_REQUEST_TIMEOUT,
226 constant::DURABILITY_TIMEOUT,
227 ));
228 let mut exists = false;
229 for ack in acks {
230 match ack {
231 OzoneMsg::KeyExists(b) => exists = b,
232 OzoneMsg::KeyChunkExists(b, 0) => exists = b, // the bunch key
233 _ => (),
234 }
235 }
236 Ok((exists, n))
237 }
238
239 /// Waits for the answer to a delete dispatched with this responder, and returns whether the
240 /// key held a value.
241 pub fn recv_delete_ack(&self) -> Outcome<bool> {
242 // Passed on as it is, for its tags, as in `recv_store_ack`.
243 let acks = ok!(self.recv_write_acks(
244 1,
245 constant::USER_REQUEST_TIMEOUT,
246 constant::DURABILITY_TIMEOUT,
247 ));
248 match acks.first() {
249 Some(OzoneMsg::KeyExists(b)) => Ok(*b),
250 msg => Err(err!(
251 "Expected an OzoneMsg::KeyExists answering a delete, received {:?}.", msg;
252 Channel, Unexpected)),
253 }
254 }
255
256 /// Collects the final answers to `n` records written under this responder. Each record is
257 /// first confirmed written, which is the writer's own work and so is held to `liveness`, and
258 /// then confirmed durable and readable, which waits on the disk and is held to `durability`.
259 /// Both deadlines measure silence, counted from the last answer of any kind: a value of many
260 /// chunks is not failed for taking long while its writers keep answering, nor a busy disk
261 /// for being slow while its barriers keep completing.
262 ///
263 /// Expiry of the first says a writer did not answer, so whether its record lands is unknown.
264 /// Expiry of the second says every record is written but not all are confirmed durable, which
265 /// is not a failure: the records are in the files and become durable, and readable, when the
266 /// disk completes them. That, and every error a store sends about a record after writing it,
267 /// is tagged `Unconfirmed`, which nothing that fails a write before it lands is.
268 pub fn recv_write_acks(
269 &self,
270 n: usize,
271 liveness: Duration,
272 durability: Duration,
273 )
274 -> Outcome<Vec<OzoneMsg<UIDL, UID, ENC, KH>>>
275 {
276 let chan = match self.channel() {
277 Some(chan) => chan,
278 None => return Err(err!("This responder does not have a channel."; Channel, Missing)),
279 };
280 let mut heard = Instant::now(); // the last answer, or the call
281 let mut written = 0;
282 let mut acks = Vec::with_capacity(n);
283 while acks.len() < n {
284 let deadline = if written < n { liveness } else { durability };
285 let left = deadline.saturating_sub(heard.elapsed());
286 if left.is_zero() {
287 if written < n {
288 return Err(err!(
289 "{} of {} records of this write were confirmed written, and the writer \
290 holding the rest has said nothing for {:?}, so whether they land is not \
291 known.", written, n, liveness;
292 Channel, Timeout));
293 }
294 return Err(err!(
295 "All {} records of this write were written, but {} of them were not \
296 confirmed durable, the disk having completed nothing for {:?}. The write \
297 has not failed: its records are in the store's files and become durable, \
298 and readable, when the disk completes them, unless the machine stops first.",
299 n, n - acks.len(), durability;
300 Write, Timeout, Unconfirmed));
301 }
302 match chan.recv_timeout(left) {
303 Recv::Empty => (), // Out of time, which the next pass reports.
304 Recv::Result(Err(e)) => return Err(err!(e,
305 "Could not read from responder channel.";
306 Channel, Read)),
307 Recv::Result(Ok(msg)) => match msg {
308 OzoneMsg::Written => {
309 written += 1;
310 heard = Instant::now();
311 },
312 OzoneMsg::KeyExists(_) |
313 OzoneMsg::KeyChunkExists(..) => {
314 acks.push(msg);
315 heard = Instant::now();
316 },
317 OzoneMsg::Finish => (),
318 OzoneMsg::Error(e) => return Err(e),
319 msg => return Err(err!(
320 "Expected the answer to a write, received {:?}.", msg;
321 Channel, Unexpected)),
322 },
323 }
324 }
325 Ok(acks)
326 }
327
328 /// Collect replies within a given time.
329 pub fn recv_number(
330 &self,
331 n: usize,
332 wait: Wait,
333 )
334 -> Outcome<(Instant, Vec<OzoneMsg<UIDL, UID, ENC, KH>>)>
335 {
336 let mut msgs = Vec::new();
337 let start = Instant::now();
338 let mut count: usize = 0;
339 match self.channel() {
340 None => return Err(err!(
341 "This responder does not have a channel.";
342 Missing, Data)),
343 Some(chan) => {
344 loop {
345 match chan.recv_timeout(wait.check_interval) {
346 Recv::Empty => (),
347 Recv::Result(Err(e)) => return Err(err!(e,
348 "Could not read from responder channel.";
349 Channel, Read)),
350 // Neither is an answer: `Finish` trails one, and `Written` precedes a
351 // write's final answer, which is what a caller counting answers wants.
352 Recv::Result(Ok(OzoneMsg::Finish)) |
353 Recv::Result(Ok(OzoneMsg::Written)) => {
354 continue;
355 }
356 Recv::Result(Ok(msg)) => {
357 msgs.push(msg);
358 count += 1;
359 if count == n {
360 break;
361 }
362 }
363 }
364 if start.elapsed() > wait.max_wait {
365 if count < n {
366 return Err(err!(
367 "Expecting {} messages via responder, received {} when \
368 timed out after {:?}.", n, count, wait.max_wait;
369 Input, Mismatch, Size));
370 } else {
371 break;
372 }
373 }
374 if count > n {
375 return Err(err!(
376 "Expecting {} messages via responder, received {} after \
377 {:?}.", count, n, start.elapsed();
378 Input, Mismatch, Size));
379 }
380 }
381 },
382 }
383 Ok((start, msgs))
384 }
385
386 /// Collect pong messages within a given time, from a dedicated ping/pong `Responder`.
387 pub fn recv_pongs(
388 &self,
389 wait: Wait,
390 )
391 -> Outcome<(Instant, HashSet<OzoneBotId>)>
392 {
393 let mut ozids = HashSet::new();
394 let start = Instant::now();
395 match self.channel() {
396 None => return Err(err!(
397 "This responder does not have a channel.";
398 Missing, Data)),
399 Some(chan) => {
400 loop {
401 match chan.recv_timeout(wait.check_interval) {
402 Recv::Empty => {}
403 Recv::Result(Err(e)) => return Err(err!(e,
404 "Could not read from responder channel.";
405 Channel, Read)),
406 Recv::Result(Ok(OzoneMsg::Pong(ozid, _errs))) => {
407 ozids.insert(ozid);
408 }
409 Recv::Result(Ok(msg)) => {
410 error!(sync_log::stream(), err!(
411 "Expecting an OzoneMsg::Pong, received a {:?}.", msg;
412 Input, Mismatch, Unexpected));
413 }
414 }
415 if start.elapsed() > wait.max_wait {
416 break;
417 }
418 }
419 },
420 }
421 Ok((start, ozids))
422 }
423
424 /// Collect replies within a given time, until an `OzoneMsg::Finish` message is received.
425 pub fn recv_all(
426 &self,
427 wait: Wait,
428 )
429 -> Outcome<(Instant, bool, Vec<OzoneMsg<UIDL, UID, ENC, KH>>)>
430 {
431 let mut complete = false;
432 let mut msgs = Vec::new();
433 let start = Instant::now();
434 match self.channel() {
435 None => return Err(err!(
436 "This responder does not have a channel.";
437 Missing, Data)),
438 Some(chan) => {
439 loop {
440 match chan.recv_timeout(wait.check_interval) {
441 Recv::Empty => (),
442 Recv::Result(Err(e)) => return Err(err!(e,
443 "Could not read from responder channel.";
444 Channel, Read)),
445 Recv::Result(Ok(OzoneMsg::Finish)) => {
446 complete = true;
447 break;
448 }
449 Recv::Result(Ok(msg)) => {
450 msgs.push(msg);
451 }
452 }
453 if start.elapsed() > wait.max_wait {
454 break;
455 }
456 }
457 },
458 }
459 Ok((start, complete, msgs))
460 }
461}
462
463pub struct Wait {
464 pub max_wait: Duration,
465 pub check_interval: Duration,
466}
467
468impl Default for Wait {
469 fn default() -> Self {
470 Self {
471 max_wait: constant::USER_REQUEST_TIMEOUT,
472 check_interval: constant::CHECK_INTERVAL,
473 }
474 }
475}
476
477impl Wait {
478
479 pub fn new(
480 max_wait: Duration,
481 check_interval: Duration,
482 )
483 -> Outcome<Self>
484 {
485 if check_interval > max_wait {
486 return Err(err!(
487 "The given check interval, {:?}, should not be larger than the \
488 given max wait, {:?}.", check_interval, max_wait;
489 Invalid, Input));
490 }
491 Ok(Self {
492 max_wait,
493 check_interval,
494 })
495 }
496
497 pub const fn new_default() -> Self {
498 Self {
499 max_wait: constant::USER_REQUEST_TIMEOUT,
500 check_interval: constant::CHECK_INTERVAL,
501 }
502 }
503
504 pub fn timeout(
505 max_wait: Duration,
506 )
507 -> Self
508 {
509 Self {
510 max_wait,
511 check_interval: Duration::default(),
512 }
513 }
514}