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
| 1 | use 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 | |
| 14 | use oxedyne_fe2o3_core::{ |
| 15 | alt::Override, |
| 16 | channels::{ |
| 17 | Recv, |
| 18 | Simplex, |
| 19 | }, |
| 20 | }; |
| 21 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 22 | use oxedyne_fe2o3_iop_crypto::enc::EncrypterDefAlt; |
| 23 | use oxedyne_fe2o3_iop_db::api::Meta; |
| 24 | use oxedyne_fe2o3_jdat::{ |
| 25 | prelude::*, |
| 26 | id::NumIdDat, |
| 27 | }; |
| 28 | |
| 29 | use std::{ |
| 30 | collections::HashSet, |
| 31 | time::{ |
| 32 | Duration, |
| 33 | Instant, |
| 34 | }, |
| 35 | }; |
| 36 | |
| 37 | #[derive(Clone, Debug)] |
| 38 | pub 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 | |
| 49 | impl< |
| 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 | |
| 463 | pub struct Wait { |
| 464 | pub max_wait: Duration, |
| 465 | pub check_interval: Duration, |
| 466 | } |
| 467 | |
| 468 | impl 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 | |
| 477 | impl 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 | } |