Oregami
Repositories/oxedyne/fe2o3

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

1use 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
53use oxedyne_fe2o3_jdat::{
54 prelude::*,
55 chunk::PartKey,
56 id::NumIdDat,
57};
58use oxedyne_fe2o3_hash::{
59 csum::{
60 ChecksummerDefAlt,
61 ChecksumScheme,
62 },
63 hash::HashScheme,
64};
65use oxedyne_fe2o3_iop_hash::api::HashForm;
66use oxedyne_fe2o3_iop_db::api::{
67 Meta,
68 RestSchemesOverride,
69 ScanOpts,
70};
71use oxedyne_fe2o3_namex::id::{
72 InNamex,
73 NamexId,
74};
75
76use 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)]
90pub 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.
108impl<
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
1695impl<
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}