Skip to main content

mkit_server/takedown/
local.rs

1//! Route the firing partition through TimerCtx.store, never a self-DO request.
2use crate::{
3    Batch, BatchOutcome, Cursor, Key, NamespaceStore, Partition, PartitionStats, ScanPage,
4    StoreCapabilities, StoreError, Value,
5};
6
7/// A borrowed local timer store plus the existing remote partition adapter.
8#[derive(Debug)]
9pub struct LocalStore<'a, L, R> {
10    local: &'a L,
11    partition: &'a Partition,
12    remote: &'a R,
13}
14impl<L, R> Clone for LocalStore<'_, L, R> {
15    fn clone(&self) -> Self {
16        Self {
17            local: self.local,
18            partition: self.partition,
19            remote: self.remote,
20        }
21    }
22}
23impl<'a, L, R> LocalStore<'a, L, R> {
24    /// Use local SQL for exactly the firing partition and remote routing otherwise.
25    pub fn new(local: &'a L, partition: &'a Partition, remote: &'a R) -> Self {
26        Self {
27            local,
28            partition,
29            remote,
30        }
31    }
32}
33impl<L: NamespaceStore, R: NamespaceStore> NamespaceStore for LocalStore<'_, L, R> {
34    fn capabilities(&self) -> StoreCapabilities {
35        self.local.capabilities()
36    }
37    async fn get(&self, p: &Partition, k: &Key) -> Result<Option<Value>, StoreError> {
38        if p == self.partition {
39            self.local.get(p, k).await
40        } else {
41            self.remote.get(p, k).await
42        }
43    }
44    async fn get_many(&self, p: &Partition, k: &[Key]) -> Result<Vec<Option<Value>>, StoreError> {
45        if p == self.partition {
46            self.local.get_many(p, k).await
47        } else {
48            self.remote.get_many(p, k).await
49        }
50    }
51    async fn scan(
52        &self,
53        p: &Partition,
54        start: &Key,
55        end: &Key,
56        after: Option<&Cursor>,
57        limit: u32,
58    ) -> Result<ScanPage, StoreError> {
59        if p == self.partition {
60            self.local.scan(p, start, end, after, limit).await
61        } else {
62            self.remote.scan(p, start, end, after, limit).await
63        }
64    }
65    async fn apply(&self, p: &Partition, b: Batch) -> Result<BatchOutcome, StoreError> {
66        if p == self.partition {
67            self.local.apply(p, b).await
68        } else {
69            self.remote.apply(p, b).await
70        }
71    }
72    async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
73        if p == self.partition {
74            self.local.stats(p).await
75        } else {
76            self.remote.stats(p).await
77        }
78    }
79    async fn probe(&self) -> Result<(), StoreError> {
80        self.local.probe().await
81    }
82}