oxedyne/fe2o3/fe2o3_o3db_sync/src/api.rs
68.2 KiB, 362 runs
created by r1870400018:715, 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::{ |
| 4 | constant, |
| 5 | id::{ |
| 6 | self, |
| 7 | OzoneBotId, |
| 8 | }, |
| 9 | index::{ |
| 10 | WorkerInd, |
| 11 | ZoneInd, |
| 12 | }, |
| 13 | }, |
| 14 | bots::{ |
| 15 | bot_zone::ZoneState, |
| 16 | worker::{ |
| 17 | bot::WorkerType, |
| 18 | bot_file::GcControl, |
| 19 | bot_reader::ReadResult, |
| 20 | }, |
| 21 | }, |
| 22 | comm::{ |
| 23 | channels::{ |
| 24 | BotChannels, |
| 25 | ChooseBot, |
| 26 | OzoneMsgCount, |
| 27 | }, |
| 28 | msg::OzoneMsg, |
| 29 | response::{ |
| 30 | Responder, |
| 31 | Wait, |
| 32 | }, |
| 33 | }, |
| 34 | data::{ |
| 35 | cache::{ |
| 36 | CacheEntry, |
| 37 | KeyVal, |
| 38 | }, |
| 39 | choose::ChooseCache, |
| 40 | core::{ |
| 41 | Encode, |
| 42 | Key, |
| 43 | RestSchemes, |
| 44 | Value, |
| 45 | }, |
| 46 | }, |
| 47 | file::{ |
| 48 | state::FileStateMap, |
| 49 | zdir::ZoneDir, |
| 50 | }, |
| 51 | }; |
| 52 | |
| 53 | use oxedyne_fe2o3_jdat::{ |
| 54 | prelude::*, |
| 55 | chunk::PartKey, |
| 56 | id::NumIdDat, |
| 57 | }; |
| 58 | use oxedyne_fe2o3_hash::{ |
| 59 | csum::{ |
| 60 | ChecksummerDefAlt, |
| 61 | ChecksumScheme, |
| 62 | }, |
| 63 | hash::HashScheme, |
| 64 | }; |
| 65 | use oxedyne_fe2o3_iop_hash::api::HashForm; |
| 66 | use oxedyne_fe2o3_iop_db::api::{ |
| 67 | Meta, |
| 68 | RestSchemesOverride, |
| 69 | ScanOpts, |
| 70 | }; |
| 71 | use oxedyne_fe2o3_namex::id::{ |
| 72 | InNamex, |
| 73 | NamexId, |
| 74 | }; |
| 75 | |
| 76 | use std::{ |
| 77 | collections::BTreeMap, |
| 78 | path::{ |
| 79 | Path, |
| 80 | PathBuf, |
| 81 | }, |
| 82 | time::{ |
| 83 | Duration, |
| 84 | Instant, |
| 85 | }, |
| 86 | }; |
| 87 | |
| 88 | |
| 89 | #[derive(Clone, Debug)] |
| 90 | pub struct OzoneApi< |
| 91 | // Data at rest. |
| 92 | const UIDL: usize, // User id byte length. |
| 93 | UID: NumIdDat<UIDL>, // User id. |
| 94 | ENC: Encrypter, // Symmetric encryption of data at rest. |
| 95 | KH: Hasher, // Hashes database keys. |
| 96 | PR: Hasher, // Pseudo-randomiser hash to distribute cache data. |
| 97 | CS: Checksummer, // Checks integrity of data at rest. |
| 98 | >{ |
| 99 | pub ozid: OzoneBotId, |
| 100 | pub db_root: PathBuf, |
| 101 | pub cfg: OzoneConfig, |
| 102 | pub chans: BotChannels<UIDL, UID, ENC, KH>, |
| 103 | pub schms: RestSchemes<ENC, KH, PR, CS>, |
| 104 | } |
| 105 | |
| 106 | /// The `'static` requirement for UID, which propagates through the code base, is initially driven |
| 107 | /// by the channel send methods in this implementation. |
| 108 | impl< |
| 109 | const UIDL: usize, |
| 110 | UID: NumIdDat<UIDL> + 'static, |
| 111 | ENC: Encrypter + 'static, |
| 112 | KH: Hasher + 'static, |
| 113 | PR: Hasher, |
| 114 | CS: Checksummer, |
| 115 | > |
| 116 | OzoneApi<UIDL, UID, ENC, KH, PR, CS> |
| 117 | { |
| 118 | /// Create a new Ozone database API instance. |
| 119 | pub fn new( |
| 120 | ozid: OzoneBotId, |
| 121 | db_root: PathBuf, |
| 122 | cfg: OzoneConfig, |
| 123 | chans: BotChannels<UIDL, UID, ENC, KH>, |
| 124 | schms: RestSchemes<ENC, KH, PR, CS>, |
| 125 | ) |
| 126 | -> Self |
| 127 | { |
| 128 | Self { |
| 129 | ozid, |
| 130 | db_root, |
| 131 | cfg, |
| 132 | chans, |
| 133 | schms, |
| 134 | } |
| 135 | } |
| 136 | |
| 137 | pub fn ozid(&self) -> &OzoneBotId { &self.ozid } |
| 138 | pub fn db_root(&self) -> &Path { &self.db_root } |
| 139 | pub fn cfg(&self) -> &OzoneConfig { &self.cfg } |
| 140 | pub fn schemes(&self) -> &RestSchemes<ENC, KH, PR, CS> { &self.schms } |
| 141 | pub fn chans(&self) -> &BotChannels<UIDL, UID, ENC, KH> { &self.chans } |
| 142 | |
| 143 | // Convenience. |
| 144 | pub fn responder(&self) -> Responder<UIDL, UID, ENC, KH> { Responder::new(Some(&self.ozid())) } |
| 145 | pub fn no_responder() -> Responder<UIDL, UID, ENC, KH> { Responder::none(None) } |
| 146 | |
| 147 | // Key API. |
| 148 | |
| 149 | /// Prepare a key for the Ozone database. |
| 150 | /// |
| 151 | /// # Arguments |
| 152 | /// * `k` - key `Dat` to be transformed into an Ozone key. |
| 153 | /// * `enc` - optional `Encypter` which can contain a `Hasher`. If this exists, this is the |
| 154 | /// hash function that is used instead of the default. |
| 155 | /// |
| 156 | /// Returns the key bytes, the zone and a hash used to deterministically select bots. |
| 157 | /// |
| 158 | /// # Local errors |
| 159 | /// * The encoded key length cannot be zero. |
| 160 | pub fn ozone_key_dat( |
| 161 | &self, |
| 162 | k: &Dat, |
| 163 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 164 | ) |
| 165 | -> Outcome<(Vec<u8>, WorkerInd, alias::ChooseHash)> |
| 166 | { |
| 167 | self.ozone_key(res!(k.as_bytes()), schms2) |
| 168 | } |
| 169 | |
| 170 | /// Derives the chunk set identifier for a value from its key bytes. A chunked value's chunk |
| 171 | /// records are addressed by `Tup5u64([set_id, index, ..])`; deriving `set_id` from the key -- |
| 172 | /// rather than from a fresh random ticket per operation -- makes a later overwrite of the same |
| 173 | /// key write its chunks under the same addresses, so the ordinary supersession path flags the |
| 174 | /// superseded chunk records old and the collector reclaims them. Seahash gives a stable |
| 175 | /// 64-bit value across runs and builds; the dedicated salt keeps it distinct from the routing |
| 176 | /// hash. The collision probability matches the random ticket it replaces (~2^-64). |
| 177 | pub fn chunk_set_id(kbuf: &[u8]) -> u64 { |
| 178 | match HashScheme::new_seahash().hash(&[kbuf], constant::CHUNK_SET_ID_SALT).as_hashform() { |
| 179 | HashForm::U64(h) => h, |
| 180 | // Seahash always yields a U64; fold any other form defensively into one. |
| 181 | other => { |
| 182 | let v = other.as_vec(); |
| 183 | let mut buf = [0u8; 8]; |
| 184 | for (i, b) in v.iter().take(8).enumerate() { |
| 185 | buf[i] = *b; |
| 186 | } |
| 187 | u64::from_be_bytes(buf) |
| 188 | }, |
| 189 | } |
| 190 | } |
| 191 | |
| 192 | pub fn ozone_key( |
| 193 | &self, |
| 194 | kbuf: Vec<u8>, |
| 195 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 196 | ) |
| 197 | -> Outcome<(Vec<u8>, WorkerInd, alias::ChooseHash)> |
| 198 | { |
| 199 | self.keygen( |
| 200 | kbuf, |
| 201 | schms2, |
| 202 | self.cfg().num_zones, |
| 203 | self.cfg().num_cbots_per_zone, |
| 204 | ) |
| 205 | } |
| 206 | |
| 207 | pub fn keygen( |
| 208 | &self, |
| 209 | kbuf: Vec<u8>, |
| 210 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 211 | nz: u16, |
| 212 | nc: u16, |
| 213 | ) |
| 214 | -> Outcome<(Vec<u8>, WorkerInd, alias::ChooseHash)> |
| 215 | { |
| 216 | if kbuf.len() == 0 { |
| 217 | return Err(err!("Key has length zero."; Input, Invalid)); |
| 218 | } |
| 219 | |
| 220 | // Compute the routing hash. This is *only* a routing signal |
| 221 | // used to select the owning cbot and zone -- it does not |
| 222 | // become the stored key form. The bytes that flow through |
| 223 | // the rest of the write path, into the cache, and onto disk |
| 224 | // are the original plaintext `kbuf`, so a later `scan` can |
| 225 | // recover the user's original `Dat` key by parsing those |
| 226 | // bytes with `Dat::from_bytes`. |
| 227 | let hash = self.schemes().key_hasher() |
| 228 | .or_hash(&[&kbuf], constant::KEY_HASH_SALT, schms2.map(|s| s.key_hasher())) |
| 229 | .as_hashform(); |
| 230 | let (cbwind, chash) = res!(ChooseCache::<PR>::choose_cbot( |
| 231 | &hash, |
| 232 | nz, |
| 233 | nc, |
| 234 | )); |
| 235 | |
| 236 | Ok(( |
| 237 | kbuf, |
| 238 | cbwind, |
| 239 | chash, |
| 240 | )) |
| 241 | } |
| 242 | |
| 243 | // Write API, for general public use. |
| 244 | |
| 245 | /// Insert key-value `Dat`icles using the given data scheme overrides. A `Responder` channel |
| 246 | /// is returned, carrying the answers `store_dat_using_responder` describes; wait on them with |
| 247 | /// `Responder::recv_store_ack`. |
| 248 | /// |
| 249 | /// # Arguments |
| 250 | /// * `k` - key `Dat`cle. |
| 251 | /// * `enc` - An optional `EncryptionScheme` that was used to store the value. An error will |
| 252 | /// be returned if the decryption does not yield a valid `Dat`icle. |
| 253 | /// |
| 254 | pub fn put( |
| 255 | &self, |
| 256 | key: Dat, |
| 257 | val: Dat, |
| 258 | user: UID, |
| 259 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 260 | ) |
| 261 | -> Outcome<Responder<UIDL, UID, ENC, KH>> |
| 262 | { |
| 263 | let resp = self.responder(); |
| 264 | let sbots = self.chans().all_sbots(); |
| 265 | let (bot, bpind) = sbots.choose_bot(&ChooseBot::Randomly); |
| 266 | match bot.send(OzoneMsg::Put { |
| 267 | key, |
| 268 | val, |
| 269 | user, |
| 270 | schms2: schms2.cloned(), |
| 271 | resp: resp.clone(), |
| 272 | }) { |
| 273 | Err(e) => Err(err!(e, |
| 274 | "{}: While sending put request to sbot {}.", self.ozid(), bpind; |
| 275 | Channel, Write)), |
| 276 | _ => Ok(resp), |
| 277 | } |
| 278 | } |
| 279 | |
| 280 | // Write API, high level, used by ServerBots. |
| 281 | |
| 282 | /// The simplest storage entry point. A default responder is automatically returned, which can |
| 283 | /// be used to gain feedback on the operation (i.e. when it is successfully completed, and |
| 284 | /// whether the key was present). |
| 285 | /// |
| 286 | /// # Arguments |
| 287 | /// * `k` - key, a reference to a `Dat`. |
| 288 | /// * `v` - value, as any type that can be converted to a `Dat` via `From`. |
| 289 | /// * `user` - `User` number responsible for request. |
| 290 | /// |
| 291 | /// Returns a default `Responder` that contains the number of chunks (0 if not chunked). |
| 292 | pub fn store( |
| 293 | &self, |
| 294 | k: Dat, |
| 295 | v: Dat, |
| 296 | user: UID, |
| 297 | ) |
| 298 | -> Outcome<Responder<UIDL, UID, ENC, KH>> |
| 299 | { |
| 300 | let resp = self.responder(); |
| 301 | res!(self.store_dat_using_responder(k, v, user, None, resp.clone())); |
| 302 | Ok(resp) |
| 303 | } |
| 304 | |
| 305 | /// Store using data schemes overrides. A default responder is automatically returned, which |
| 306 | /// can be used to gain feedback on the operation (i.e. when it is successfully completed, and |
| 307 | /// whether the key was present). |
| 308 | /// |
| 309 | /// # Arguments |
| 310 | /// * `k` - key, a reference to a `Dat`. |
| 311 | /// * `v` - value, as any type that can be converted to a `Dat` via `From`. |
| 312 | /// * `user` - `User` number responsible for request. |
| 313 | /// * `schms2` - `RestSchemesOverride` overrides database schemes (e.g. key hashing, encryption). |
| 314 | /// |
| 315 | /// Returns a default `Responder` that contains the number of chunks (0 if not chunked). |
| 316 | pub fn store_using_schemes( |
| 317 | &self, |
| 318 | k: Dat, |
| 319 | v: Dat, |
| 320 | user: UID, |
| 321 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 322 | ) |
| 323 | -> Outcome<Responder<UIDL, UID, ENC, KH>> |
| 324 | { |
| 325 | let resp = self.responder(); |
| 326 | res!(self.store_dat_using_responder(k, v, user, schms2, resp.clone())); |
| 327 | Ok(resp) |
| 328 | } |
| 329 | |
| 330 | /// Store a value without a responder to provide feedback on the operation. |
| 331 | /// |
| 332 | /// # Arguments |
| 333 | /// * `k` - key, a reference to a `Dat`. |
| 334 | /// * `v` - value, as any type that can be converted to a `Dat` via `From`. |
| 335 | /// * `user` - `User` number responsible for request. |
| 336 | /// * `schms2` - `RestSchemesOverride` overrides database schemes (e.g. key hashing, encryption). |
| 337 | pub fn store_blindly( |
| 338 | &self, |
| 339 | k: Dat, |
| 340 | v: Dat, |
| 341 | user: UID, |
| 342 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 343 | ) |
| 344 | -> Outcome<()> |
| 345 | { |
| 346 | let resp = Self::no_responder(); |
| 347 | res!(self.store_dat_using_responder(k, v, user, schms2, resp)); |
| 348 | Ok(()) |
| 349 | } |
| 350 | |
| 351 | /// Store value with a default responder and a custom chunk size. The encoded value will be |
| 352 | /// chunked if its size exceed the specified chunk size. |
| 353 | /// |
| 354 | /// # Arguments |
| 355 | /// * `k` - key, a reference to a `Dat`. |
| 356 | /// * `v` - value, as any type that can be converted to a `Dat` via `From`. |
| 357 | /// * `user` - `User` number responsible for request. |
| 358 | /// * `schms2` - `RestSchemesOverride` overrides database schemes (e.g. key hashing, |
| 359 | /// encryption), including the chunker confgiuration. |
| 360 | /// |
| 361 | /// Returns a default `Responder` that contains the number of chunks (0 if not chunked). |
| 362 | pub fn store_chunked( |
| 363 | &self, |
| 364 | k: Dat, |
| 365 | v: Dat, |
| 366 | user: UID, |
| 367 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 368 | ) |
| 369 | -> Outcome<(Responder<UIDL, UID, ENC, KH>, usize)> |
| 370 | { |
| 371 | let resp = self.responder(); |
| 372 | let num_chunks = res!(self.store_dat_using_responder(k, v, user, schms2, resp.clone())); |
| 373 | Ok((resp, num_chunks)) |
| 374 | } |
| 375 | |
| 376 | /// The primary method for storing a key-value pair. |
| 377 | /// |
| 378 | /// # Data chunking |
| 379 | /// Large values are chunked and spread across multiple zones. The key points to a "bunch key" |
| 380 | /// which contains the essential chunk information. The bunch key must have an index of 0, and |
| 381 | /// the chunks are indexed from 1. Acquiring the bunch key thus allows all chunk keys to be |
| 382 | /// reconstructed and retrieved. The chunk values are stored as raw bytes wrapped in a |
| 383 | /// `Dat::BU64`. Chunking is performed when the size of the value to be stored exceeds |
| 384 | /// some fraction of the maximum data file size, which can be customised via the given |
| 385 | /// `Responder`. Without encryption, the final chunk may be smaller than the rest. With |
| 386 | /// encryption this final chunk is padded with random bytes to the uniform size. |
| 387 | /// |
| 388 | /// ```ignore |
| 389 | /// How data chunks are stored: |
| 390 | /// +-- this is the |
| 391 | /// | "bunch key" |
| 392 | /// "Store 'hello' -> 100 bytes" where chunk size is 30 bytes... | |
| 393 | /// v |
| 394 | /// +---------------------------------------+ +---------------------------------------+ |
| 395 | /// | Dat::Str("hello") | -------> | PartKey((<id>, 0, 100, 4, 30)) | |
| 396 | /// +---------------------------------------+ +---------------------------------------+ |
| 397 | /// +---------------------------------------+ +---------------------------------------+ |
| 398 | /// | PartKey((<id>, 1, 100, 4, 30)) | -------> | Dat::BU64(<30 bytes>) | |
| 399 | /// +---------------------------------------+ +---------------------------------------+ |
| 400 | /// +---------------------------------------+ +---------------------------------------+ |
| 401 | /// | PartKey((<id>, 2, 100, 4, 30)) | -------> | Dat::BU64(<30 bytes>) | |
| 402 | /// +---------------------------------------+ +---------------------------------------+ |
| 403 | /// +---------------------------------------+ +---------------------------------------+ |
| 404 | /// | PartKey((<id>, 3, 100, 4, 30)) | -------> | Dat::BU64(<30 bytes>) | |
| 405 | /// +---------------------------------------+ +---------------------------------------+ |
| 406 | /// +---------------------------------------+ +---------------------------------------+ |
| 407 | /// | PartKey((<id>, 4, 100, 4, 30)) | -------> | Dat::BU64(<10 bytes>) | |
| 408 | /// +---------------------------------------+ +---------------------------------------+ |
| 409 | /// |
| 410 | /// ``` |
| 411 | /// Placing chunks inside (unencrypted) `Dat::BU64` wrappers allows the size to be known |
| 412 | /// during cache initialisation if there is missing index file data. |
| 413 | /// |
| 414 | /// # Arguments |
| 415 | /// * `k` - key, a reference to a `Dat`. |
| 416 | /// * `v` - an owned value, any type that can be converted to a `Dat` via `From`. |
| 417 | /// * `user` - `User` number responsible for request. |
| 418 | /// * `schms2` - `RestSchemesOverride` overrides database schemes (e.g. key hashing, encryption). |
| 419 | /// * `resp` - a `Responder` channel. |
| 420 | /// |
| 421 | /// Returns the number of chunks. The first message in the `Responder` will be an |
| 422 | /// `OzoneMsg::Chunks` containing the number of chunks. Each record written is then answered |
| 423 | /// twice: `OzoneMsg::Written` once it is appended, and a final answer once it is durable under |
| 424 | /// the sync policy and readable. Without chunking the final answer is an `OzoneMsg::KeyExists` |
| 425 | /// saying whether the key was present; with chunking it is an `OzoneMsg::KeyChunkExists` for |
| 426 | /// the bunch key (with index 0) and for each chunk. The answers to different records arrive |
| 427 | /// in no particular order. A writer that fails sends `OzoneMsg::Error` instead. |
| 428 | /// `Responder::recv_store_ack` waits on all of it, each stage to its own deadline. |
| 429 | /// |
| 430 | /// # Local errors |
| 431 | /// * The key must be transformable into an Ozone key. |
| 432 | /// * The encoded value length must exceed zero. This should not occur. |
| 433 | /// * The chunk size in the responder cannot be zero. |
| 434 | pub fn store_dat_using_responder( |
| 435 | &self, |
| 436 | k: Dat, |
| 437 | v: Dat, |
| 438 | user: UID, |
| 439 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 440 | resp: Responder<UIDL, UID, ENC, KH>, |
| 441 | ) |
| 442 | -> Outcome<usize> |
| 443 | { |
| 444 | let msgs = res!(self.prepare_write_dat( |
| 445 | k, |
| 446 | v, |
| 447 | user, |
| 448 | schms2, |
| 449 | resp.clone(), |
| 450 | None, |
| 451 | )); |
| 452 | let nchunks = msgs.len(); |
| 453 | if resp.is_some() { |
| 454 | res!(resp.send(OzoneMsg::Chunks(nchunks))); |
| 455 | } |
| 456 | res!(self.store_bytes(msgs)); |
| 457 | Ok(nchunks) |
| 458 | } |
| 459 | |
| 460 | /// Store forcing the chunk set identifier rather than deriving it from the key. Test and |
| 461 | /// migration support: it reproduces a value as an earlier build wrote it (a random |
| 462 | /// per-operation set_id), so that reads of such a value can be exercised after the switch to |
| 463 | /// key-derived identifiers. Production writes never take this path. |
| 464 | pub fn store_dat_using_responder_forcing_set_id( |
| 465 | &self, |
| 466 | k: Dat, |
| 467 | v: Dat, |
| 468 | user: UID, |
| 469 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 470 | resp: Responder<UIDL, UID, ENC, KH>, |
| 471 | set_id: u64, |
| 472 | ) |
| 473 | -> Outcome<usize> |
| 474 | { |
| 475 | let msgs = res!(self.prepare_write_dat( |
| 476 | k, |
| 477 | v, |
| 478 | user, |
| 479 | schms2, |
| 480 | resp.clone(), |
| 481 | Some(set_id), |
| 482 | )); |
| 483 | let nchunks = msgs.len(); |
| 484 | if resp.is_some() { |
| 485 | res!(resp.send(OzoneMsg::Chunks(nchunks))); |
| 486 | } |
| 487 | res!(self.store_bytes(msgs)); |
| 488 | Ok(nchunks) |
| 489 | } |
| 490 | |
| 491 | /// The key and value `Dat`icles are serialised here and then sent for final processing. A |
| 492 | /// `set_id_override` of `None` derives the chunk set identifier from the key (the ordinary |
| 493 | /// path); `Some` forces it, for reproducing an earlier build's random-keyed values. |
| 494 | pub fn prepare_write_dat( |
| 495 | &self, |
| 496 | k: Dat, |
| 497 | v: Dat, |
| 498 | user: UID, |
| 499 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 500 | resp: Responder<UIDL, UID, ENC, KH>, |
| 501 | set_id_override: Option<u64>, |
| 502 | ) |
| 503 | -> Outcome<Vec<(OzoneMsg<UIDL, UID, ENC, KH>, ZoneInd)>> |
| 504 | { |
| 505 | //let (kbuf, vbuf) = res!(Encode::encode_dat(k, v)); |
| 506 | let (kbuf, vbuf) = res!(Encode::encode_dat(k.clone(), v.clone())); |
| 507 | self.prepare_write( |
| 508 | kbuf, |
| 509 | vbuf, |
| 510 | user, |
| 511 | schms2, |
| 512 | resp, |
| 513 | set_id_override, |
| 514 | ) |
| 515 | } |
| 516 | |
| 517 | /// This is the last step in preparing the serialised data for storage, where |
| 518 | /// `RestSchemesOverride` is finally invoked. This influences how the key is hashed, if and |
| 519 | /// how the data is chunked, and if and how those chunks are encrypted. `OzoneMsg`s are |
| 520 | /// returned, ready for sending to `WriteBot`s. |
| 521 | pub fn prepare_write( |
| 522 | &self, |
| 523 | k: Vec<u8>, |
| 524 | mut vbuf: Vec<u8>, |
| 525 | user: UID, |
| 526 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 527 | resp: Responder<UIDL, UID, ENC, KH>, |
| 528 | set_id_override: Option<u64>, |
| 529 | ) |
| 530 | -> Outcome<Vec<(OzoneMsg<UIDL, UID, ENC, KH>, ZoneInd)>> |
| 531 | { |
| 532 | if vbuf.len() == 0 { |
| 533 | return Err(err!( |
| 534 | "{}: For key {:?}, the given value encoded length is zero.", |
| 535 | self.ozid(), k; |
| 536 | Input, Invalid)); |
| 537 | } |
| 538 | |
| 539 | // 1. Normalise the key. |
| 540 | let (kbuf, cbwind, chash) = res!(self.ozone_key(k, schms2)); |
| 541 | |
| 542 | // 3. Define chunking. |
| 543 | let chunk_config = match schms2 { |
| 544 | Some(schms2) => match schms2.chunk_config() { |
| 545 | Some(cfg) => cfg.clone(), |
| 546 | None => self.cfg().chunk_config(), |
| 547 | }, |
| 548 | None => self.cfg().chunk_config(), |
| 549 | }; |
| 550 | let chunk_threshold = chunk_config.threshold_bytes; |
| 551 | |
| 552 | let encryption_on = !(self.schemes().encrypter().or_is_identity(schms2.map(|s| s.encrypter()))); |
| 553 | debug!(sync_log::stream(), "Encryption is on: {}", encryption_on); |
| 554 | if encryption_on { |
| 555 | vbuf = res!( |
| 556 | self.schemes().encrypter().or_encrypt(&mut vbuf, schms2.map(|s| s.encrypter())) |
| 557 | ); |
| 558 | } |
| 559 | |
| 560 | let mut msgs = Vec::new(); |
| 561 | let mut meta = Meta::new(user); |
| 562 | res!(meta.stamp_time_now()); |
| 563 | |
| 564 | // 4. Package the value, breaking into chunks if it is too big. |
| 565 | if vbuf.len() >= chunk_threshold { |
| 566 | let chunker = OzoneConfig::chunker(chunk_config); |
| 567 | // 4.1 Chunk data. |
| 568 | let (chunks, chunk_state) = res!(chunker.chunk(&vbuf)); |
| 569 | // Address the chunks by a key-derived identifier, not the per-operation ticket, so an |
| 570 | // overwrite of the same key supersedes the prior value's chunk records in place. A |
| 571 | // forced identifier (test/migration only) reproduces an earlier build's random keys. |
| 572 | let set_id = match set_id_override { |
| 573 | Some(id) => id, |
| 574 | None => Self::chunk_set_id(&kbuf), |
| 575 | }; |
| 576 | let datkeys = res!(chunker.keys(set_id, &chunk_state)); |
| 577 | |
| 578 | // 4.2 Store main key -> bunch key. |
| 579 | let mut bkbuf = res!(datkeys[0].as_bytes()); |
| 580 | if encryption_on { |
| 581 | bkbuf = res!(self.schemes().encrypter().or_encrypt(&bkbuf, schms2.map(|s| s.encrypter()))); |
| 582 | bkbuf = res!(Dat::wrap_bytes_var(bkbuf)); |
| 583 | } |
| 584 | msgs.push(( |
| 585 | res!(Self::package_write( |
| 586 | KeyVal { |
| 587 | key: Key::Chunk(kbuf, 0), |
| 588 | val: bkbuf, |
| 589 | chash, |
| 590 | meta: meta.clone(), |
| 591 | cbpind: **cbwind.bpind(), |
| 592 | }, |
| 593 | resp.clone(), |
| 594 | self.schemes().checksummer().clone(), |
| 595 | )), |
| 596 | *cbwind.zind(), |
| 597 | )); |
| 598 | |
| 599 | // 4.3 Store chunk keys -> chunk bytes. |
| 600 | for (i, chunk) in chunks.into_iter().enumerate() { |
| 601 | let (ckbuf, ccbwind, cchash) = |
| 602 | res!(self.ozone_key_dat(&datkeys[i+1], schms2)); |
| 603 | msgs.push(( |
| 604 | res!(Self::package_write( |
| 605 | KeyVal { |
| 606 | key: Key::Chunk(ckbuf, i + 1), |
| 607 | val: chunk, |
| 608 | chash: cchash, |
| 609 | meta: meta.clone(), |
| 610 | cbpind: **ccbwind.bpind(), |
| 611 | }, |
| 612 | resp.clone(), |
| 613 | self.schemes().checksummer().clone(), |
| 614 | )), |
| 615 | *ccbwind.zind(), |
| 616 | )); |
| 617 | } |
| 618 | } else { |
| 619 | // 3.1 No chunking, just a single block of data. |
| 620 | if encryption_on { |
| 621 | vbuf = res!(Dat::wrap_bytes_var(vbuf)); |
| 622 | } |
| 623 | msgs.push(( |
| 624 | res!(Self::package_write( |
| 625 | KeyVal { |
| 626 | key: Key::Complete(kbuf), |
| 627 | val: vbuf, |
| 628 | chash, |
| 629 | meta: meta.clone(), |
| 630 | cbpind: **cbwind.bpind(), |
| 631 | }, |
| 632 | resp, |
| 633 | self.schemes().checksummer().clone(), |
| 634 | )), |
| 635 | *cbwind.zind(), |
| 636 | )); |
| 637 | } |
| 638 | Ok(msgs) |
| 639 | } |
| 640 | |
| 641 | /// This is the write dispatch method, where `WriterBots` are chosen randomly. Callers must |
| 642 | /// ensure the value is wrapped in a `Dat::BU64`. |
| 643 | /// |
| 644 | /// # Local errors |
| 645 | /// * The write request message cannot be sent via a `WriterBot` channel. |
| 646 | pub fn store_bytes( |
| 647 | &self, |
| 648 | msgs: Vec<(OzoneMsg<UIDL, UID, ENC, KH>, ZoneInd)>, |
| 649 | ) |
| 650 | -> Outcome<()> |
| 651 | { |
| 652 | for (msg, zind) in msgs { |
| 653 | let wbots = res!(self.chans().get_workers_of_type_in_zone(&WorkerType::Writer, &zind)); |
| 654 | let (bot, bpind) = wbots.choose_bot(&ChooseBot::Randomly); |
| 655 | match bot.send(msg) { |
| 656 | Err(e) => return Err(err!(e, |
| 657 | "{}: While sending write request to wbot {}.", |
| 658 | self.ozid(), WorkerInd::new(zind, bpind); |
| 659 | Channel, Write)), |
| 660 | _ => (), |
| 661 | } |
| 662 | } |
| 663 | Ok(()) |
| 664 | } |
| 665 | |
| 666 | pub fn package_write( |
| 667 | kv: KeyVal<UIDL, UID>, |
| 668 | resp: Responder<UIDL, UID, ENC, KH>, |
| 669 | csummer: ChecksummerDefAlt<ChecksumScheme, CS>, |
| 670 | ) |
| 671 | -> Outcome<OzoneMsg<UIDL, UID, ENC, KH>> |
| 672 | { |
| 673 | let klen_cache = kv.key.len(); |
| 674 | let (kstored, vstored, cind, meta, cbpind, _, _) = res!(Encode::encode(kv, csummer)); |
| 675 | |
| 676 | Ok(OzoneMsg::Write{ |
| 677 | kstored, |
| 678 | vstored, |
| 679 | klen_cache, |
| 680 | cind, |
| 681 | meta, |
| 682 | cbpind, |
| 683 | resp, |
| 684 | }) |
| 685 | } |
| 686 | |
| 687 | pub fn delete_using_responder( |
| 688 | &self, |
| 689 | k: &Dat, |
| 690 | user: UID, |
| 691 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 692 | resp: Responder<UIDL, UID, ENC, KH>, |
| 693 | ) |
| 694 | -> Outcome<()> |
| 695 | { |
| 696 | // 1. Normalise the key. The cache hash belongs to the stored record, not merely to the |
| 697 | // routing decision, so it is carried through to the writer rather than dropped. |
| 698 | let (kstored, cbwind, chash) = res!(self.ozone_key_dat(k, schms2)); |
| 699 | |
| 700 | // 1a. If the value is chunked, the tombstone on the user key below supersedes only the |
| 701 | // bunch key; the chunk records live under their own keys and would leak forever (the |
| 702 | // whole reason a chunked value's chunks are never rewritten on delete). So read the |
| 703 | // current bunch key, reconstruct each chunk key from the set_id it stores -- random |
| 704 | // for a pre-upgrade value, key-derived for a new one, either way exactly what |
| 705 | // `fetch_chunks` reconstructs to read them -- and tombstone each so the ordinary |
| 706 | // supersession path reclaims them. These carry no responder: the caller waits only |
| 707 | // on the single bunch-key delete below. The read is confined to the delete path, |
| 708 | // which is rare relative to writes, and only chunked values pay the fan-out. |
| 709 | res!(self.reclaim_chunks_on_delete(k, user, schms2)); |
| 710 | |
| 711 | // 2. The value we use to indicate deletion is an unencrypted custom usr type. |
| 712 | let v = Dat::Usr(id::usr_kind_id_deleted(), Some(Box::new(Dat::Empty))); |
| 713 | let vstored = res!(v.as_bytes()); |
| 714 | |
| 715 | // 3. Select a zone writer bot. |
| 716 | let wbots = res!(self.chans().get_workers_of_type_in_zone(&WorkerType::Writer, cbwind.zind())); |
| 717 | let (bot, bpind) = wbots.choose_bot(&ChooseBot::Randomly); |
| 718 | |
| 719 | // 4. Create the metadata. |
| 720 | let mut meta = Meta::new(user); |
| 721 | res!(meta.stamp_time_now()); |
| 722 | |
| 723 | // 5. Frame the tombstone through the same encoder an insertion goes through. |
| 724 | let msg = res!(Self::package_write( |
| 725 | KeyVal { |
| 726 | key: Key::Complete(kstored), |
| 727 | val: vstored, |
| 728 | chash, |
| 729 | meta, |
| 730 | cbpind: **cbwind.bpind(), |
| 731 | }, |
| 732 | resp, |
| 733 | self.schemes().checksummer().clone(), |
| 734 | )); |
| 735 | |
| 736 | // 6. Send write request, with responder. |
| 737 | match bot.send(msg) { |
| 738 | Err(e) => return Err(err!(e, |
| 739 | "{}: While sending delete request to wbot {}.", |
| 740 | self.ozid(), WorkerInd::new(*cbwind.zind(), bpind); |
| 741 | Channel, Write)), |
| 742 | _ => Ok(()), |
| 743 | } |
| 744 | } |
| 745 | |
| 746 | /// Reads the current value at `k` and, if it is chunked, tombstones every chunk record so the |
| 747 | /// collector reclaims them. A no-op for an unchunked or absent value. Chunk keys are |
| 748 | /// reconstructed from the part key exactly as `fetch_chunks` does, so this works for values |
| 749 | /// written under either the old random set_id or the new key-derived one. |
| 750 | fn reclaim_chunks_on_delete( |
| 751 | &self, |
| 752 | k: &Dat, |
| 753 | user: UID, |
| 754 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 755 | ) |
| 756 | -> Outcome<()> |
| 757 | { |
| 758 | let enc = self.schemes().encrypter(); |
| 759 | let or_enc = schms2.map(|s| s.encrypter()); |
| 760 | |
| 761 | let resp = res!(self.fetch_using_schemes(k, schms2)); |
| 762 | let pkey = match res!(resp.recv_daticle(enc, or_enc)) { |
| 763 | (Some((Dat::Tup5u64(tup), _)), _) => PartKey(tup), |
| 764 | _ => return Ok(()), // Not chunked, or the key is absent: nothing extra to reclaim. |
| 765 | }; |
| 766 | |
| 767 | for i in 1..(pkey.num_parts() + 1) { |
| 768 | let ck = Dat::Tup5u64([ |
| 769 | pkey.set_id(), |
| 770 | i, |
| 771 | pkey.data_len(), |
| 772 | pkey.num_parts(), |
| 773 | pkey.part_size(), |
| 774 | ]); |
| 775 | // No responder: the caller waits only on the single bunch-key delete. |
| 776 | res!(self.tombstone_chunk_key(&ck, user, schms2, Self::no_responder())); |
| 777 | } |
| 778 | Ok(()) |
| 779 | } |
| 780 | |
| 781 | /// Writes an unencrypted deleted-kind tombstone at the `Key::Complete` form of a chunk-data |
| 782 | /// key, dispatched to the writer of the key's routed zone under the given responder. Because |
| 783 | /// the stored-key bytes are the chunk key's bytes either way, the tombstone supersedes the |
| 784 | /// chunk record at the same cache key, so the ordinary supersession collector flags the |
| 785 | /// chunk's bytes old and reclaims them -- the only in-place way to retire a chunk record, since |
| 786 | /// the store has no primitive that forgets a key without writing something at it. Shared by |
| 787 | /// the delete path, which reclaims a deleted value's chunks, and the orphan sweep, which |
| 788 | /// reclaims chunk records no live bunch key references. |
| 789 | /// |
| 790 | /// `ck` must be the chunk's `Dat::Tup5u64` part key. The tombstone is left unencrypted, |
| 791 | /// exactly as an ordinary key delete leaves it, so a reader recognises it without the at-rest |
| 792 | /// key. |
| 793 | pub fn tombstone_chunk_key( |
| 794 | &self, |
| 795 | ck: &Dat, |
| 796 | user: UID, |
| 797 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 798 | resp: Responder<UIDL, UID, ENC, KH>, |
| 799 | ) |
| 800 | -> Outcome<()> |
| 801 | { |
| 802 | let (ckbuf, ccbwind, cchash) = res!(self.ozone_key_dat(ck, schms2)); |
| 803 | let tomb = Dat::Usr(id::usr_kind_id_deleted(), Some(Box::new(Dat::Empty))); |
| 804 | let tvstored = res!(tomb.as_bytes()); |
| 805 | let mut cmeta = Meta::new(user); |
| 806 | res!(cmeta.stamp_time_now()); |
| 807 | let msg = res!(Self::package_write( |
| 808 | KeyVal { |
| 809 | key: Key::Complete(ckbuf), |
| 810 | val: tvstored, |
| 811 | chash: cchash, |
| 812 | meta: cmeta, |
| 813 | cbpind: **ccbwind.bpind(), |
| 814 | }, |
| 815 | resp, |
| 816 | self.schemes().checksummer().clone(), |
| 817 | )); |
| 818 | let cwbots = res!(self.chans().get_workers_of_type_in_zone(&WorkerType::Writer, ccbwind.zind())); |
| 819 | let (cbot, cbpind) = cwbots.choose_bot(&ChooseBot::Randomly); |
| 820 | match cbot.send(msg) { |
| 821 | Err(e) => Err(err!(e, |
| 822 | "{}: While sending chunk tombstone for {:?} to wbot {}.", |
| 823 | self.ozid(), ck, WorkerInd::new(*ccbwind.zind(), cbpind); |
| 824 | Channel, Write)), |
| 825 | _ => Ok(()), |
| 826 | } |
| 827 | } |
| 828 | |
| 829 | // Read API, for general public use. |
| 830 | |
| 831 | /// Get a `Dat`icle value using the given key and data scheme overrides. The result is |
| 832 | /// available asynchronously in the returned `Responder` channel. |
| 833 | /// |
| 834 | /// # Arguments |
| 835 | /// * `k` - key `Dat`cle. |
| 836 | /// * `enc` - An optional `EncryptionScheme` that was used to store the value. An error will be returned if the decryption does not yield a valid `Dat`icle. |
| 837 | /// |
| 838 | pub fn get( |
| 839 | &self, |
| 840 | key: &Dat, |
| 841 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 842 | ) |
| 843 | -> Outcome<Responder<UIDL, UID, ENC, KH>> |
| 844 | { |
| 845 | let resp = self.responder(); |
| 846 | let sbots = self.chans().all_sbots(); |
| 847 | let (bot, bpind) = sbots.choose_bot(&ChooseBot::Randomly); |
| 848 | // Send read request, with responder. |
| 849 | match bot.send(OzoneMsg::Get { |
| 850 | key: key.clone(), |
| 851 | schms2: schms2.cloned(), |
| 852 | resp: resp.clone(), |
| 853 | }) { |
| 854 | Err(e) => Err(err!(e, |
| 855 | "{}: While sending get request to sbot {}.", |
| 856 | self.ozid(), bpind; |
| 857 | Channel, Write)), |
| 858 | _ => Ok(resp), |
| 859 | } |
| 860 | } |
| 861 | |
| 862 | // Read API, high level, used by ServerBots. |
| 863 | // |
| 864 | /// Blocking retrieval of a `Dat`icle value using the given key and data scheme overrides. |
| 865 | /// |
| 866 | /// # Arguments |
| 867 | /// * `k` - key `Dat` to be transformed into an Ozone key. |
| 868 | /// * `schms2` - `RestSchemesOverride` overrides database schemes (e.g. key hashing, encryption). |
| 869 | /// |
| 870 | pub fn get_wait( |
| 871 | &self, |
| 872 | k: &Dat, |
| 873 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 874 | ) |
| 875 | -> Outcome<Option<(Dat, Meta<UIDL, UID>)>> |
| 876 | { |
| 877 | let enc = self.schemes().encrypter(); |
| 878 | let or_enc = schms2.map(|s| s.encrypter()); |
| 879 | |
| 880 | let resp = res!(self.fetch_using_schemes(k, schms2)); |
| 881 | match res!(resp.recv_daticle(enc, or_enc)) { |
| 882 | (None, _) => Ok(None), // The key was not found. |
| 883 | // The value was too large for a single record, so what is stored under this key is a |
| 884 | // part key naming its chunks. `fetch_chunks` gathers them, rejoins the bytes, |
| 885 | // decrypts them and decodes them, so what it hands back is already the caller's |
| 886 | // value -- fully formed, of whatever kind they stored. |
| 887 | // |
| 888 | // It was previously taken for raw bytes and decoded a SECOND time, which no chunked |
| 889 | // value survives: a list, a map or a string fell through to a catch-all and came back |
| 890 | // as an error, and a byte string -- the one kind the arms matched -- had its payload |
| 891 | // read as though it were itself an encoding. So every value large enough to be |
| 892 | // chunked was written perfectly well and could not be read: an accumulating value, |
| 893 | // such as a ledger, worked until the day it crossed the chunk size and then failed |
| 894 | // for good. |
| 895 | (Some((Dat::Tup5u64(tup), meta)), _) => |
| 896 | Ok(Some((res!(self.fetch_chunks(&Dat::Tup5u64(tup), schms2)), meta))), |
| 897 | // The data received was in a single piece. |
| 898 | (Some((dat, meta)), _) => Ok(Some((dat, meta))), |
| 899 | } |
| 900 | } |
| 901 | |
| 902 | // Read API, lower level. |
| 903 | |
| 904 | /// Fetch a value using the given key. This is just a caller of `OzoneApi::fetch_using_responder` |
| 905 | /// that provides a default `Responder`. Default database schemes (e.g. encryption) are used. |
| 906 | /// |
| 907 | /// # Arguments |
| 908 | /// * `k` - key `Dat` to be transformed into an Ozone key. |
| 909 | /// |
| 910 | /// Returns a default `Responder`. |
| 911 | pub fn fetch( |
| 912 | &self, |
| 913 | k: &Dat, |
| 914 | ) |
| 915 | -> Outcome<Responder<UIDL, UID, ENC, KH>> |
| 916 | { |
| 917 | let resp = self.responder(); |
| 918 | res!(self.fetch_using_responder(k, None, resp.clone())); |
| 919 | Ok(resp) |
| 920 | } |
| 921 | |
| 922 | /// Fetch a value using the given key and data schemes override. This is just a caller of |
| 923 | /// `OzoneApi::fetch_using_responder` that provides a default `Responder`. |
| 924 | /// |
| 925 | /// # Arguments |
| 926 | /// * `k` - key `Dat` to be transformed into an Ozone key. |
| 927 | /// * `schms2` - `RestSchemesOverride` overrides database schemes (e.g. key hashing, encryption). |
| 928 | /// |
| 929 | /// Returns a default `Responder`. |
| 930 | pub fn fetch_using_schemes( |
| 931 | &self, |
| 932 | k: &Dat, |
| 933 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 934 | ) |
| 935 | -> Outcome<Responder<UIDL, UID, ENC, KH>> |
| 936 | { |
| 937 | let resp = self.responder(); |
| 938 | res!(self.fetch_using_responder(k, schms2, resp.clone())); |
| 939 | Ok(resp) |
| 940 | } |
| 941 | |
| 942 | pub fn fetch_using_key( |
| 943 | &self, |
| 944 | key: Key, |
| 945 | cbwind: WorkerInd, |
| 946 | ) |
| 947 | -> Outcome<Responder<UIDL, UID, ENC, KH>> |
| 948 | { |
| 949 | let resp = self.responder(); |
| 950 | res!(self.fetch_using_key_and_responder( |
| 951 | key, |
| 952 | cbwind, |
| 953 | resp.clone(), |
| 954 | )); |
| 955 | Ok(resp) |
| 956 | } |
| 957 | |
| 958 | /// Fetch a value using the given key, scheme overrides and a customisable `Responder`. The |
| 959 | /// caller can use the `Responder` to wait for a single value, or an error. An error will |
| 960 | /// result if the value cannot be decoded into a `Dat`. This can occur if the value was |
| 961 | /// improperly stored or the given decrypter does not match the original encrypter. If the |
| 962 | /// value was chunked, a `PartKey` "bunch key" will be returned, which can be passed to |
| 963 | /// `OzoneApi::fetch_chunks` to collect the chunks and re-assemble the value. |
| 964 | /// |
| 965 | /// # Arguments |
| 966 | /// * `k` - key `Dat` to be transformed into an Ozone key. |
| 967 | /// * `schms2` - `RestSchemesOverride` overrides database schemes (e.g. key hashing, encryption). |
| 968 | /// |
| 969 | /// # Local errors |
| 970 | /// * An error will result if the request cannot be sent to the randomly chosen `ReaderBot`. |
| 971 | pub fn fetch_using_responder( |
| 972 | &self, |
| 973 | k: &Dat, |
| 974 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 975 | resp: Responder<UIDL, UID, ENC, KH>, |
| 976 | ) |
| 977 | -> Outcome<()> |
| 978 | { |
| 979 | // Normalise the key. |
| 980 | let (kbuf, cbwind, _chash) = res!(self.ozone_key_dat(k, schms2)); |
| 981 | let key = match k { |
| 982 | Dat::Tup5u64(tup) => Key::Chunk(kbuf, try_into!(usize, PartKey(*tup).index())), |
| 983 | _ => Key::Complete(kbuf), |
| 984 | }; |
| 985 | self.fetch_using_key_and_responder( |
| 986 | key, |
| 987 | cbwind, |
| 988 | resp, |
| 989 | ) |
| 990 | } |
| 991 | |
| 992 | pub fn fetch_using_key_and_responder( |
| 993 | &self, |
| 994 | key: Key, |
| 995 | cbwind: WorkerInd, // CacheBot worker index. |
| 996 | resp: Responder<UIDL, UID, ENC, KH>, |
| 997 | ) |
| 998 | -> Outcome<()> |
| 999 | { |
| 1000 | // Select a zone reader bot. |
| 1001 | let rbots = res!(self.chans().get_workers_of_type_in_zone(&WorkerType::Reader, cbwind.zind())); |
| 1002 | let (bot, bpind) = rbots.choose_bot(&ChooseBot::Randomly); |
| 1003 | // Send read request, with responder. |
| 1004 | match bot.send(OzoneMsg::Read(key, **cbwind.bpind(), resp)) { |
| 1005 | Err(e) => return Err(err!(e, |
| 1006 | "{}: While sending read request to rbot {}.", |
| 1007 | self.ozid(), WorkerInd::new(*cbwind.zind(), bpind); |
| 1008 | Channel, Write)), |
| 1009 | _ => Ok(()), |
| 1010 | } |
| 1011 | } |
| 1012 | |
| 1013 | /// Data is automatically chunked when stored, such that each chunk is accessed via its own |
| 1014 | /// `PartKey`. However chunked data is not automatically reassembled. A valid `PartKey` |
| 1015 | /// ("bunch key") passed to this method will perform the collection and reassembly. |
| 1016 | /// |
| 1017 | /// # Arguments |
| 1018 | /// * `k` - bunch key `PartKey` which provides all necessary chunk metrics. |
| 1019 | /// * `schms2` - `RestSchemesOverride` overrides database schemes (e.g. key hashing, encryption). |
| 1020 | /// |
| 1021 | /// # Local errors |
| 1022 | /// * The key must be a `PartKey`. |
| 1023 | /// * The `PartKey` part size must exceed zero. |
| 1024 | /// * The `PartKey` number of parts must exceed zero. |
| 1025 | /// * The `PartKey` index must be zero. |
| 1026 | /// * A read request cannot be sent to a randomly chosen `ReaderBot`. |
| 1027 | /// * A `Dat` value cannot be received from the `Responder`. |
| 1028 | /// * The `PartKey` key for a chunk cannot have an index value of zero. |
| 1029 | /// * The index of a chunk cannot exceed the expected number of chunks. This can occur if the |
| 1030 | /// chunks were incorrectly stored. |
| 1031 | /// * The chunk value must be wrapped in a `Dat::BU64`. |
| 1032 | /// * The unwrapped length of all chunks must match the bunch key part size, except for the final |
| 1033 | /// chunk. Note that `store_using_responder` pads the final chunk to the uniform size when |
| 1034 | /// encryption is used. |
| 1035 | /// * If the final chunk length differs from the part size, it must not exceed the part size. |
| 1036 | /// * An error will be raised if the total length of the chunk data received exceeds the |
| 1037 | /// expected capacity of the receiving receptable. This error should not occur. |
| 1038 | /// * An error will occur if a chunk cannot be found. |
| 1039 | /// * An error will occur if the `ReaderBot` responds with an unexpected message. |
| 1040 | /// * The re-assembled data value must be an encoded `Dat`. |
| 1041 | /// |
| 1042 | pub fn fetch_chunks( |
| 1043 | &self, |
| 1044 | k: &Dat, |
| 1045 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 1046 | ) |
| 1047 | -> Outcome<Dat> |
| 1048 | { |
| 1049 | let self_id = self.ozid().clone(); |
| 1050 | let enc = self.schemes().encrypter(); |
| 1051 | let or_enc = schms2.map(|s| s.encrypter()); |
| 1052 | let encryption_on = !(enc.or_is_identity(or_enc)); |
| 1053 | |
| 1054 | match k { |
| 1055 | Dat::Tup5u64(tup) => { |
| 1056 | let pkey = PartKey(*tup); |
| 1057 | if pkey.part_size() == 0 { |
| 1058 | return Err(err!( |
| 1059 | "{}: Chunk size must exceed zero.", self_id; |
| 1060 | Input, Invalid)); |
| 1061 | } |
| 1062 | if pkey.num_parts() == 0 { |
| 1063 | return Err(err!( |
| 1064 | "{}: Number of chunks must exceed zero.", self_id; |
| 1065 | Input, Invalid)); |
| 1066 | } |
| 1067 | if pkey.index() != 0 { |
| 1068 | return Err(err!( |
| 1069 | "{}: Index in bunch key must be zero.", self_id; |
| 1070 | Input, Invalid)); |
| 1071 | } |
| 1072 | // 1. Send requests for chunks. |
| 1073 | let data_len = try_into!(usize, pkey.data_len()); |
| 1074 | let chunk_size = try_into!(usize, pkey.part_size()); |
| 1075 | let num_chunks = try_into!(usize, pkey.num_parts()); |
| 1076 | let resp = self.responder(); |
| 1077 | for i in 1..(pkey.num_parts() + 1) { |
| 1078 | let k = Dat::Tup5u64([ |
| 1079 | pkey.set_id(), |
| 1080 | i, |
| 1081 | pkey.data_len(), |
| 1082 | pkey.num_parts(), |
| 1083 | pkey.part_size(), |
| 1084 | ]); |
| 1085 | let (kbuf, cbwind, _chash) = res!(self.ozone_key_dat(&k, schms2)); |
| 1086 | |
| 1087 | let rbots = res!(self.chans().get_workers_of_type_in_zone(&WorkerType::Reader, cbwind.zind())); |
| 1088 | let (bot, bpind) = rbots.choose_bot(&ChooseBot::Randomly); |
| 1089 | let key = Key::Chunk(kbuf, try_into!(usize, i)); |
| 1090 | match bot.send(OzoneMsg::Read(key, **cbwind.bpind(), resp.clone())) { |
| 1091 | Err(e) => return Err(err!(e, |
| 1092 | "{}: While sending chunk {} read request to rbot {}.", |
| 1093 | self_id, i, WorkerInd::new(*cbwind.zind(), bpind); |
| 1094 | Channel, Write)), |
| 1095 | _ => (), |
| 1096 | } |
| 1097 | } |
| 1098 | // 2. Collection and reassembly. |
| 1099 | let capacity = num_chunks * chunk_size; |
| 1100 | let mut joined = vec![0; capacity]; |
| 1101 | for _ in 0..num_chunks { |
| 1102 | match resp.recv_timeout(constant::USER_REQUEST_TIMEOUT) { |
| 1103 | Err(e) => return Err(err!(e, |
| 1104 | "{}: Could not read from chunk collection responder channel.", self_id; |
| 1105 | IO, Channel, Read)), |
| 1106 | Ok(OzoneMsg::Value(Value::Chunk(Some((Dat::BU8(v), _)), i, _))) | |
| 1107 | Ok(OzoneMsg::Value(Value::Chunk(Some((Dat::BU16(v), _)), i, _))) | |
| 1108 | Ok(OzoneMsg::Value(Value::Chunk(Some((Dat::BU32(v), _)), i, _))) | |
| 1109 | Ok(OzoneMsg::Value(Value::Chunk(Some((Dat::BU64(v), _)), i, _))) => { |
| 1110 | if i == 0 { |
| 1111 | return Err(err!( |
| 1112 | "{}: For key {:?}, data chunk of size {} has an invalid \ |
| 1113 | index of zero amongst an expected total of {} chunks.", |
| 1114 | self_id, k, v.len(), num_chunks; |
| 1115 | Invalid, Input)); |
| 1116 | } |
| 1117 | if i > num_chunks { |
| 1118 | return Err(err!( |
| 1119 | "{}: For key {:?}, data chunk of size {} with index {} \ |
| 1120 | exceeds the expected number of chunks, {}.", |
| 1121 | self_id, k, v.len(), i, num_chunks; |
| 1122 | Invalid, Input)); |
| 1123 | } |
| 1124 | let mut end = chunk_size * i; |
| 1125 | let mut start = end - chunk_size; |
| 1126 | if v.len() != chunk_size { |
| 1127 | if i < num_chunks { |
| 1128 | return Err(err!( |
| 1129 | "{}: For key {:?}, data chunk {} of {} size of {} does \ |
| 1130 | not match the size of {} specified by the \ |
| 1131 | PartKey.", |
| 1132 | self_id, k, i, num_chunks, v.len(), chunk_size; |
| 1133 | Input, Size, Mismatch)); |
| 1134 | } else { |
| 1135 | if v.len() > chunk_size { |
| 1136 | return Err(err!( |
| 1137 | "{}: For key {:?}, the final data chunk {} size \ |
| 1138 | of {} must be less than the size of the {} other \ |
| 1139 | chunks, {} bytes.", |
| 1140 | self_id, k, i, v.len(), num_chunks-1, chunk_size; |
| 1141 | Input, Size, Invalid)); |
| 1142 | } else { |
| 1143 | start = chunk_size * (i-1); |
| 1144 | end = start + v.len(); |
| 1145 | } |
| 1146 | } |
| 1147 | } |
| 1148 | if end > capacity { |
| 1149 | return Err(err!( |
| 1150 | "{}: For key {:?}, end location {} for retrieved data \ |
| 1151 | (chunk {} of {}) of length {} exceeds the end location \ |
| 1152 | of the expected reassembled data, {}.", |
| 1153 | self_id, k, end, i, num_chunks, chunk_size, capacity; |
| 1154 | Bug, Input, Size, Mismatch)); |
| 1155 | } |
| 1156 | joined[start..end].copy_from_slice(&v[..]); |
| 1157 | }, |
| 1158 | Ok(OzoneMsg::Value(Value::Chunk(None, i, _))) => return Err(err!( |
| 1159 | "{}: For key {:?}, data chunk {} of {} was not found.", |
| 1160 | self_id, k, i, num_chunks; |
| 1161 | Missing, Data)), |
| 1162 | Ok(msg) => return Err(err!( |
| 1163 | "{}: Unrecognised chunk request response: {:?}", self_id, msg; |
| 1164 | Invalid, Input)), |
| 1165 | } |
| 1166 | } |
| 1167 | if encryption_on { |
| 1168 | joined = res!(enc.or_decrypt(&joined[..data_len], or_enc)); |
| 1169 | } |
| 1170 | match Dat::from_bytes(&joined) { |
| 1171 | Err(e) => return Err(err!(e, |
| 1172 | "{}: For key {:?}, a Dat could not be formed from the value bytes. \ |
| 1173 | This could mean the data was not originally stored as a Dat, or the \ |
| 1174 | encrypter, {}, differs from that used to store the original data.", |
| 1175 | self_id, k, enc.or_debug(or_enc); |
| 1176 | Decode, Bytes)), |
| 1177 | Ok((dat, _)) => return Ok(dat), |
| 1178 | } |
| 1179 | }, |
| 1180 | _ => return Err(err!("{}: Key must be a PartKey.", self_id; Input, Invalid)), |
| 1181 | } |
| 1182 | } |
| 1183 | |
| 1184 | /// Walks every index file in every zone, so the cost is the size of the store rather than |
| 1185 | /// the size of the result. The deadline is the ordinary user request one, which is what |
| 1186 | /// keeps "nothing on a request path may scan" enforceable rather than advisory: a walk big |
| 1187 | /// enough to matter fails here instead of stalling the caller. A background walk that has |
| 1188 | /// deliberately accepted the cost says so at its call site with `scan_with_wait`. |
| 1189 | pub fn scan( |
| 1190 | &self, |
| 1191 | opts: &ScanOpts, |
| 1192 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 1193 | ) |
| 1194 | -> Outcome<Vec<(Dat, Dat, Meta<UIDL, UID>)>> |
| 1195 | { |
| 1196 | self.scan_with_wait(opts, schms2, constant::USER_REQUEST_WAIT) |
| 1197 | } |
| 1198 | |
| 1199 | /// `scan` with the deadline named by the caller instead of taken from `USER_REQUEST_WAIT`. |
| 1200 | /// |
| 1201 | /// # Arguments |
| 1202 | /// |
| 1203 | /// * `wait` - how long every zone has, in total, to return its entries. Only a caller that |
| 1204 | /// knows it is off the request path should lengthen this. |
| 1205 | pub fn scan_with_wait( |
| 1206 | &self, |
| 1207 | opts: &ScanOpts, |
| 1208 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 1209 | wait: Wait, |
| 1210 | ) |
| 1211 | -> Outcome<Vec<(Dat, Dat, Meta<UIDL, UID>)>> |
| 1212 | { |
| 1213 | let nz = self.cfg().num_zones(); |
| 1214 | let max_wait = wait.max_wait; |
| 1215 | let resp = self.responder(); |
| 1216 | |
| 1217 | // Send one ScanRequest to a scan bot of every zone. |
| 1218 | for z in 0..nz { |
| 1219 | let zind = ZoneInd::new(z); |
| 1220 | let scbots = res!(self.chans().get_workers_of_type_in_zone( |
| 1221 | &WorkerType::Scan, |
| 1222 | &zind, |
| 1223 | )); |
| 1224 | let (bot, _) = scbots.choose_bot(&ChooseBot::Randomly); |
| 1225 | if let Err(e) = bot.send(OzoneMsg::ScanRequest { |
| 1226 | opts: opts.clone(), |
| 1227 | schms2: schms2.cloned(), |
| 1228 | resp: resp.clone(), |
| 1229 | }) { |
| 1230 | return Err(err!(e, |
| 1231 | "{}: Cannot send scan request to scan bot in zone {}.", |
| 1232 | self.ozid(), z; |
| 1233 | Channel, Write)); |
| 1234 | } |
| 1235 | } |
| 1236 | |
| 1237 | // Gather one ScanEntries response per zone. |
| 1238 | let (_, msgs) = match resp.recv_number(nz, wait) { |
| 1239 | Ok(v) => v, |
| 1240 | Err(e) => return Err(err!(e, |
| 1241 | "{}: A scan of all {} zones did not finish within {:?}. A scan walks every \ |
| 1242 | index file in every zone, so it takes longer the larger the store is, however \ |
| 1243 | few entries match; a scan on a request path is what this deadline exists to \ |
| 1244 | catch. A caller that is deliberately off the request path should name its own \ |
| 1245 | deadline with scan_with_wait rather than lengthening \ |
| 1246 | constant::USER_REQUEST_TIMEOUT, which every user request shares.", |
| 1247 | self.ozid(), nz, max_wait; |
| 1248 | Channel, Timeout)), |
| 1249 | }; |
| 1250 | let mut out: Vec<(Dat, Dat, Meta<UIDL, UID>)> = Vec::new(); |
| 1251 | for msg in msgs { |
| 1252 | match msg { |
| 1253 | OzoneMsg::ScanEntries(entries) => { |
| 1254 | out.extend(entries); |
| 1255 | }, |
| 1256 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1257 | "{}: Zone-level scan failure.", self.ozid(); |
| 1258 | Channel)), |
| 1259 | other => return Err(err!( |
| 1260 | "{}: Unexpected response to scan request: {:?}", |
| 1261 | self.ozid(), other; |
| 1262 | Channel, Unexpected)), |
| 1263 | } |
| 1264 | } |
| 1265 | |
| 1266 | // Apply the global limit. Per-zone limits have already been |
| 1267 | // applied inside each igbot, so this cap tightens the |
| 1268 | // cross-zone merge rather than truncating any single zone. |
| 1269 | if let Some(lim) = opts.limit { |
| 1270 | if out.len() > lim { |
| 1271 | out.truncate(lim); |
| 1272 | } |
| 1273 | } |
| 1274 | Ok(out) |
| 1275 | } |
| 1276 | |
| 1277 | /// Explains a control operation that did not complete, since the bare shortfall from |
| 1278 | /// `recv_number` names neither the operation, nor the reason a healthy database can miss |
| 1279 | /// the deadline, nor the constant to change if it should not have. |
| 1280 | fn control_failure( |
| 1281 | &self, |
| 1282 | e: Error<ErrTag>, |
| 1283 | opn: &str, // the operation attempted |
| 1284 | n: usize, // acknowledgements expected |
| 1285 | who: &str, // the bots expected to acknowledge |
| 1286 | ) |
| 1287 | -> Error<ErrTag> |
| 1288 | { |
| 1289 | err!(e, |
| 1290 | "{}: The {} failed while waiting up to {:?} for all {} {} to acknowledge it. A \ |
| 1291 | bot acknowledges a control message only once it reaches it in its queue, so the \ |
| 1292 | usual cause is a store large enough that the bots are still surveying its files \ |
| 1293 | after startup, rather than any fault in the operation. This deadline is \ |
| 1294 | constant::CONTROL_REQUEST_TIMEOUT; raise that if a store legitimately needs \ |
| 1295 | longer, and not constant::USER_REQUEST_TIMEOUT, which is deliberately short \ |
| 1296 | because every user request shares it.", |
| 1297 | self.ozid(), opn, constant::CONTROL_REQUEST_TIMEOUT, n, who; |
| 1298 | Channel, Timeout) |
| 1299 | } |
| 1300 | |
| 1301 | /// Activate garbage collection by sending a control message to the igbots via the zbots, via the supervisor. |
| 1302 | pub fn activate_gc(&self, on: bool) -> Outcome<()> { |
| 1303 | info!(sync_log::stream(), "Activating garbage collection..."); |
| 1304 | let emsg = "garbage collection activation"; |
| 1305 | let resp = self.responder(); |
| 1306 | if let Err(e) = self.chans().sup().send( |
| 1307 | OzoneMsg::GcControl(GcControl::On(on), resp.clone()) |
| 1308 | ) { |
| 1309 | return Err(err!(e, |
| 1310 | "{}: Cannot send {} to supervisor.", self.ozid(), emsg; |
| 1311 | Channel, Write)); |
| 1312 | } |
| 1313 | // A control operation, not a user request: this runs once, at startup, behind whatever |
| 1314 | // initialisation the zone bots are still doing. See constant::CONTROL_REQUEST_TIMEOUT. |
| 1315 | let nz = self.cfg().num_zones(); |
| 1316 | let (_, msgs) = match resp.recv_number(nz, constant::CONTROL_REQUEST_WAIT) { |
| 1317 | Ok(v) => v, |
| 1318 | Err(e) => return Err(self.control_failure(e, emsg, nz, "zone bots")), |
| 1319 | }; |
| 1320 | for msg in msgs { |
| 1321 | match msg { |
| 1322 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1323 | "{}: In response to {}.", self.ozid(), emsg; |
| 1324 | Channel)), |
| 1325 | OzoneMsg::Ok => (), |
| 1326 | msg => return Err(err!( |
| 1327 | "{}: Unexpected response to {}: {:?}", self.ozid(), emsg, msg; |
| 1328 | Channel)), |
| 1329 | }; |
| 1330 | } |
| 1331 | Ok(()) |
| 1332 | } |
| 1333 | |
| 1334 | // Utility methods useful for situational awareness and testing. |
| 1335 | |
| 1336 | /// Command all cbots to clear their caches. |
| 1337 | pub fn clear_cache_values(&self, wait: Wait) -> Outcome<()> { |
| 1338 | let resp = self.responder(); |
| 1339 | if let Err(e) = self.chans().sup().send( |
| 1340 | OzoneMsg::ClearCache(resp.clone()) |
| 1341 | ) { |
| 1342 | return Err(err!(e, |
| 1343 | "{}: Cannot send clear cache command to supervisor.", self.ozid(); |
| 1344 | Channel, Write)); |
| 1345 | } |
| 1346 | let n = self.cfg().num_caches(); |
| 1347 | let (_, msgs) = res!(resp.recv_number(n, wait)); |
| 1348 | for msg in msgs { |
| 1349 | match msg { |
| 1350 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1351 | "{}: In response to clear cache command.", self.ozid(); |
| 1352 | Channel)), |
| 1353 | OzoneMsg::Ok => (), |
| 1354 | msg => return Err(err!( |
| 1355 | "{}: Unexpected response to clear cache command: {:?}", self.ozid(), msg; |
| 1356 | Channel)), |
| 1357 | } |
| 1358 | } |
| 1359 | warn!(sync_log::stream(), "All {} caches successfully cleared.", n); |
| 1360 | Ok(()) |
| 1361 | } |
| 1362 | |
| 1363 | /// Dump all cache contents to the log file. |
| 1364 | pub fn dump_caches(&self, wait: Wait) -> Outcome<()> { |
| 1365 | // Gather. |
| 1366 | let resp = self.responder(); |
| 1367 | if let Err(e) = self.chans().sup().send( |
| 1368 | OzoneMsg::DumpCacheRequest(resp.clone()) |
| 1369 | ) { |
| 1370 | return Err(err!(e, |
| 1371 | "{}: Cannot send cache dump request to supervisor.", self.ozid(); |
| 1372 | Channel, Write)); |
| 1373 | } |
| 1374 | let n = self.cfg().num_caches(); |
| 1375 | let (_, msgs) = res!(resp.recv_number(n, wait)); |
| 1376 | let mut sorted = BTreeMap::new(); |
| 1377 | for msg in msgs { |
| 1378 | match msg { |
| 1379 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1380 | "{}: In response to cache dump request.", self.ozid(); |
| 1381 | Channel)), |
| 1382 | OzoneMsg::DumpCacheResponse(wind, cache) => { |
| 1383 | sorted.insert(wind, cache); |
| 1384 | }, |
| 1385 | msg => return Err(err!( |
| 1386 | "{}: Unexpected response to cache dump request: {:?}", self.ozid(), msg; |
| 1387 | Channel)), |
| 1388 | } |
| 1389 | } |
| 1390 | // Display. |
| 1391 | info!(sync_log::stream(), "Cache dump summary"); |
| 1392 | info!(sync_log::stream(), "+-----------+--------------+--------------+"); |
| 1393 | info!(sync_log::stream(), "| Cache | Entries | Size [B] |"); |
| 1394 | info!(sync_log::stream(), "+-----------+--------------+--------------+"); |
| 1395 | for (wind, cache) in &sorted { |
| 1396 | info!(sync_log::stream(), "|{:^11}|{:>13} |{:>13} |", |
| 1397 | fmt!("{}", wind), |
| 1398 | cache.map().len(), |
| 1399 | cache.get_size(), |
| 1400 | ); |
| 1401 | } |
| 1402 | info!(sync_log::stream(), "+-----------+--------------+--------------+"); |
| 1403 | for (wind, cache) in sorted { |
| 1404 | let mut total_size = 0; |
| 1405 | info!(sync_log::stream(), "{} cache dump of {} entries:", wind, cache.map().len()); |
| 1406 | if cache.map().len() == 0 { |
| 1407 | info!(sync_log::stream(), " No cache entries."); |
| 1408 | } else { |
| 1409 | for (kbyt, centry) in cache.map() { |
| 1410 | if let CacheEntry::LocatedValue(mloc, val) = centry { |
| 1411 | //let (k, _) = res!(Dat::from_bytes(&kbyt)); |
| 1412 | let vlen = match val { |
| 1413 | Some(v) => v.len(), |
| 1414 | None => 0, |
| 1415 | }; |
| 1416 | let size = |
| 1417 | kbyt.len() + |
| 1418 | vlen + |
| 1419 | cache.mloc_size(); |
| 1420 | |
| 1421 | info!(sync_log::stream(), " kbyt = {:02x?} vlen = {} floc = {:?}", |
| 1422 | kbyt, vlen, mloc.file_location(), |
| 1423 | ); |
| 1424 | total_size += size; |
| 1425 | } |
| 1426 | } |
| 1427 | } |
| 1428 | info!(sync_log::stream(), "{} cache size estimate: {} [B]", wind, total_size); |
| 1429 | } |
| 1430 | |
| 1431 | Ok(()) |
| 1432 | } |
| 1433 | |
| 1434 | /// Returns cache entry for the given key. |
| 1435 | pub fn cache_entry_info( |
| 1436 | &self, |
| 1437 | k: &Dat, |
| 1438 | schms2: Option<&RestSchemesOverride<ENC, KH>>, |
| 1439 | timeout: Duration, |
| 1440 | ) |
| 1441 | -> Outcome<ReadResult<UIDL, UID>> |
| 1442 | { |
| 1443 | let (kbuf, cbwind, _chash) = res!(self.ozone_key_dat(k, schms2)); |
| 1444 | let resp = self.responder(); |
| 1445 | |
| 1446 | let cbots = res!(self.chans().get_workers_of_type_in_zone(&WorkerType::Cache, cbwind.zind())); |
| 1447 | let bot = res!(cbots.get_bot(**cbwind.bpind())); |
| 1448 | let key = match k { |
| 1449 | Dat::Tup5u64(tup) => Key::Chunk(kbuf, try_into!(usize, PartKey(*tup).index())), |
| 1450 | _ => Key::Complete(kbuf), |
| 1451 | }; |
| 1452 | match bot.send(OzoneMsg::ReadCache(key, resp.clone())) { |
| 1453 | Err(e) => return Err(err!(e, |
| 1454 | "{}: While sending cache entry info request to cbot {}.", |
| 1455 | self.ozid(), cbwind; |
| 1456 | Channel, Write)), |
| 1457 | _ => (), |
| 1458 | } |
| 1459 | |
| 1460 | match res!(resp.recv_timeout(timeout)) { |
| 1461 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1462 | "{}: In response to cache entry info request.", self.ozid(); |
| 1463 | Channel)), |
| 1464 | OzoneMsg::ReadResult(readres) => Ok(readres), |
| 1465 | msg => Err(err!( |
| 1466 | "{}: Unexpected response to cache entry info request: {:?}", self.ozid(), msg; |
| 1467 | Channel)), |
| 1468 | } |
| 1469 | } |
| 1470 | |
| 1471 | pub fn collect_file_states( |
| 1472 | &self, |
| 1473 | wait: Wait, |
| 1474 | ) |
| 1475 | -> Outcome<BTreeMap<WorkerInd, FileStateMap>> |
| 1476 | { |
| 1477 | let resp = self.responder(); |
| 1478 | if let Err(e) = self.chans().sup().send( |
| 1479 | OzoneMsg::DumpFileStatesRequest(resp.clone()) |
| 1480 | ) { |
| 1481 | return Err(err!(e, |
| 1482 | "{}: Cannot send file state dump request to supervisor.", self.ozid(); |
| 1483 | Channel, Write)); |
| 1484 | } |
| 1485 | let n = self.cfg().num_filemaps(); |
| 1486 | let (_, msgs) = res!(resp.recv_number(n, wait)); |
| 1487 | let mut sorted = BTreeMap::new(); |
| 1488 | for msg in msgs { |
| 1489 | match msg { |
| 1490 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1491 | "{}: In response to file state dump request.", self.ozid(); |
| 1492 | Channel)), |
| 1493 | OzoneMsg::DumpFileStatesResponse(wind, fstates) => { |
| 1494 | sorted.insert(wind, fstates); |
| 1495 | }, |
| 1496 | msg => return Err(err!( |
| 1497 | "{}: Unexpected response to file states dump request: {:?}", self.ozid(), msg; |
| 1498 | Channel)), |
| 1499 | } |
| 1500 | } |
| 1501 | Ok(sorted) |
| 1502 | } |
| 1503 | |
| 1504 | /// Dump all zone file states to the log file. |
| 1505 | pub fn dump_file_states(&self, wait: Wait) -> Outcome<()> { |
| 1506 | // Gather. |
| 1507 | let sorted = res!(self.collect_file_states(wait)); |
| 1508 | // Display. |
| 1509 | for (wind, fstates) in sorted { |
| 1510 | info!(sync_log::stream(), "{} file states dump:", wind); |
| 1511 | if fstates.map().len() == 0 { |
| 1512 | info!(sync_log::stream(), " None"); |
| 1513 | } else { |
| 1514 | for (fnum, fstat) in fstates.map() { |
| 1515 | info!(sync_log::stream(), " {:10} {:?} old = {:.1}%", |
| 1516 | fnum, fstat, |
| 1517 | 100.0 * (fstat.get_old_sum() as f64) |
| 1518 | / (self.cfg().data_file_max_bytes as f64), |
| 1519 | ); |
| 1520 | } |
| 1521 | } |
| 1522 | } |
| 1523 | |
| 1524 | Ok(()) |
| 1525 | } |
| 1526 | |
| 1527 | /// Returns the number of messages in all channels for all zones. |
| 1528 | pub fn ozone_msg_count(&self) -> OzoneMsgCount { |
| 1529 | self.chans().msg_count() |
| 1530 | } |
| 1531 | |
| 1532 | /// Returns the file directory size, in-memory cache size and bot message queues for each zone, in bytes. |
| 1533 | pub fn ozone_state(&self, wait: Wait) -> Outcome<Vec<ZoneState>> { |
| 1534 | let emsg = "ozone state request"; |
| 1535 | let resp = self.responder(); |
| 1536 | if let Err(e) = self.chans().sup().send( |
| 1537 | OzoneMsg::OzoneStateRequest(resp.clone()) |
| 1538 | ) { |
| 1539 | return Err(err!(e, |
| 1540 | "{}: Cannot send {} to supervisor.", self.ozid(), emsg; |
| 1541 | Channel, Write)); |
| 1542 | } |
| 1543 | match res!(resp.recv_timeout(wait.max_wait)) { |
| 1544 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1545 | "{}: In response to {}.", self.ozid(), emsg; |
| 1546 | Channel)), |
| 1547 | OzoneMsg::OzoneStateResponse(zstats) => Ok(zstats), |
| 1548 | msg => Err(err!( |
| 1549 | "{}: Unexpected response to {}: {:?}", self.ozid(), emsg, msg; |
| 1550 | Channel)), |
| 1551 | } |
| 1552 | } |
| 1553 | |
| 1554 | /// Ping the bots for proof of life. |
| 1555 | pub fn ping_bots(&self, wait: Wait) -> Outcome<(Instant, Vec<OzoneMsg<UIDL, UID, ENC, KH>>)> { |
| 1556 | let resp = self.responder(); |
| 1557 | let self_id = self.ozid().clone(); |
| 1558 | let n = res!(self.chans().send_to_all(OzoneMsg::Ping(self_id, resp.clone()))); |
| 1559 | resp.recv_number(n, wait) |
| 1560 | } |
| 1561 | |
| 1562 | pub fn bot_error_count(&self, wait: Wait) -> Outcome<(usize, usize)> { |
| 1563 | let (_, msgs) = res!(self.ping_bots(wait)); |
| 1564 | let mut nbots: usize = 0; |
| 1565 | let mut errs: usize = 0; |
| 1566 | for msg in msgs { |
| 1567 | match msg { |
| 1568 | OzoneMsg::Pong(_ozid, n) => { |
| 1569 | nbots += 1; |
| 1570 | errs = errs.saturating_add(n); |
| 1571 | }, |
| 1572 | msg => return Err(err!( |
| 1573 | "{}: Unexpected response to bot ping: {:?}", self.ozid(), msg; |
| 1574 | Channel)), |
| 1575 | } |
| 1576 | } |
| 1577 | Ok((errs, nbots)) |
| 1578 | } |
| 1579 | |
| 1580 | pub fn list_files(&self, wait: Wait) -> Outcome<()> { |
| 1581 | |
| 1582 | info!(sync_log::stream(), "Directory listing for {} zones, key:", self.cfg().num_zones()); |
| 1583 | info!(sync_log::stream(), " Typ: f File | d Directory | s Symlink"); |
| 1584 | info!(sync_log::stream(), " Size: in bytes"); |
| 1585 | info!(sync_log::stream(), " Mod: seconds since last modified"); |
| 1586 | info!(sync_log::stream(), " Name: object label"); |
| 1587 | |
| 1588 | let emsg = "list files request"; |
| 1589 | let resp = self.responder(); |
| 1590 | if let Err(e) = self.chans().sup().send( |
| 1591 | OzoneMsg::DumpFiles(resp.clone()) |
| 1592 | ) { |
| 1593 | return Err(err!(e, |
| 1594 | "{}: Cannot send {} to supervisor.", self.ozid(), emsg; |
| 1595 | Channel, Write)); |
| 1596 | } |
| 1597 | let (_, msgs) = res!(resp.recv_number(self.cfg().num_zones(), wait)); |
| 1598 | let mut map = BTreeMap::new(); |
| 1599 | for msg in msgs { |
| 1600 | match msg { |
| 1601 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1602 | "{}: In response to {}.", self.ozid(), emsg; |
| 1603 | Channel)), |
| 1604 | OzoneMsg::Files(zind, zmap) => map.insert(zind, zmap), |
| 1605 | msg => return Err(err!( |
| 1606 | "{}: Unexpected response to {}: {:?}", self.ozid(), emsg, msg; |
| 1607 | Channel)), |
| 1608 | }; |
| 1609 | } |
| 1610 | for (zind, zmap) in map { |
| 1611 | let mut total_size = 0; |
| 1612 | info!(sync_log::stream(), "{:?} directory", zind); |
| 1613 | info!(sync_log::stream(), "+-----+--------------+--------------+-----------------------------------------------------"); |
| 1614 | info!(sync_log::stream(), "| Typ | Size [B] | Mod [s] | Name"); |
| 1615 | info!(sync_log::stream(), "+-----+--------------+--------------+-----------------------------------------------------"); |
| 1616 | for (_key, entry) in zmap { |
| 1617 | info!(sync_log::stream(), |
| 1618 | "| {} |{:>13} |{:>13} | {}", |
| 1619 | entry.typ, |
| 1620 | entry.size, |
| 1621 | entry.mods, |
| 1622 | entry.name, |
| 1623 | ); |
| 1624 | total_size += entry.size; |
| 1625 | } |
| 1626 | info!(sync_log::stream(), "+-----+--------------+--------------+-----------------------------------------------------"); |
| 1627 | info!(sync_log::stream(), "| |{:>13} | |", total_size); |
| 1628 | info!(sync_log::stream(), "+-----+--------------+--------------+-----------------------------------------------------"); |
| 1629 | } |
| 1630 | Ok(()) |
| 1631 | } |
| 1632 | |
| 1633 | pub fn get_zone_dirs(&self) -> Outcome<BTreeMap<ZoneInd, ZoneDir>> { |
| 1634 | let emsg = "zone directories request"; |
| 1635 | let resp = self.responder(); |
| 1636 | if let Err(e) = self.chans().sup().send( |
| 1637 | OzoneMsg::GetZoneDir(resp.clone()) |
| 1638 | ) { |
| 1639 | return Err(err!(e, |
| 1640 | "{}: Cannot send {} to supervisor.", self.ozid(), emsg; |
| 1641 | Channel, Write)); |
| 1642 | } |
| 1643 | // Deliberately a user request deadline and not the control one: this reports what the |
| 1644 | // zone bots already hold, it is callable at any time rather than only during startup, |
| 1645 | // and a caller wanting an answer is better told quickly that there is none. |
| 1646 | let (_, msgs) = res!(resp.recv_number(self.cfg().num_zones(), constant::USER_REQUEST_WAIT)); |
| 1647 | let mut map = BTreeMap::new(); |
| 1648 | for msg in msgs { |
| 1649 | match msg { |
| 1650 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1651 | "{}: In response to {}.", self.ozid(), emsg; |
| 1652 | Channel)), |
| 1653 | OzoneMsg::ZoneDir(zind, zdir) => map.insert(zind, zdir), |
| 1654 | msg => return Err(err!( |
| 1655 | "{}: Unexpected response to {}: {:?}", self.ozid(), emsg, msg; |
| 1656 | Channel)), |
| 1657 | }; |
| 1658 | } |
| 1659 | Ok(map) |
| 1660 | } |
| 1661 | |
| 1662 | /// Instruct the wbots to increment to their next live files, to provide a clean slate for |
| 1663 | /// testing. |
| 1664 | pub fn new_live_files(&self) -> Outcome<()> { |
| 1665 | let emsg = "new live files request"; |
| 1666 | let resp = self.responder(); |
| 1667 | if let Err(e) = self.chans().sup().send( |
| 1668 | OzoneMsg::NewLiveFile(None, resp.clone()) |
| 1669 | ) { |
| 1670 | return Err(err!(e, |
| 1671 | "{}: Cannot send {} to supervisor.", self.ozid(), emsg; |
| 1672 | Channel, Write)); |
| 1673 | } |
| 1674 | // A control operation like `activate_gc`, and queued behind the same zone bot work. |
| 1675 | let nw = self.cfg().num_wbots(); |
| 1676 | let (_, msgs) = match resp.recv_number(nw, constant::CONTROL_REQUEST_WAIT) { |
| 1677 | Ok(v) => v, |
| 1678 | Err(e) => return Err(self.control_failure(e, emsg, nw, "writer bots")), |
| 1679 | }; |
| 1680 | for msg in msgs { |
| 1681 | match msg { |
| 1682 | OzoneMsg::Error(e) => return Err(err!(e, |
| 1683 | "{}: In response to {}.", self.ozid(), emsg; |
| 1684 | Channel)), |
| 1685 | OzoneMsg::Ok => (), |
| 1686 | msg => return Err(err!( |
| 1687 | "{}: Unexpected response to {}: {:?}", self.ozid(), emsg, msg; |
| 1688 | Channel)), |
| 1689 | }; |
| 1690 | } |
| 1691 | Ok(()) |
| 1692 | } |
| 1693 | } |
| 1694 | |
| 1695 | impl< |
| 1696 | const UIDL: usize, |
| 1697 | UID: NumIdDat<UIDL> + 'static, |
| 1698 | ENC: Encrypter + 'static, |
| 1699 | KH: Hasher + 'static, |
| 1700 | PR: Hasher + 'static, |
| 1701 | CS: Checksummer + 'static, |
| 1702 | > |
| 1703 | InNamex for OzoneApi<UIDL, UID, ENC, KH, PR, CS> |
| 1704 | { |
| 1705 | fn name_id(&self) -> Outcome<NamexId> { |
| 1706 | NamexId::try_from(constant::NAMEX_ID) |
| 1707 | } |
| 1708 | } |