Skip to main content

mai_sdk_core/
storage.rs

1use crate::{network::PeerId, task_queue::TaskId};
2use async_channel::Sender;
3use std::{collections::HashMap, sync::Arc};
4use tokio::sync::RwLock;
5
6pub type TaskAssignments = Arc<RwLock<HashMap<TaskId, PeerId>>>;
7
8pub type RemoteTasks<Task> = Arc<RwLock<HashMap<TaskId, Task>>>;
9
10pub type OwnedTasks<Task, TaskOutput> = Arc<RwLock<HashMap<TaskId, (Task, Sender<TaskOutput>)>>>;
11
12use anyhow::{bail, Result};
13use slog::{error, info, Logger};
14
15use crate::bridge::EventBridge;
16
17#[derive(Clone, Debug)]
18pub struct DistributedKVStore {
19    logger: Logger,
20    bridge: EventBridge,
21}
22
23type Value = Vec<u8>;
24
25#[derive(Debug, Clone)]
26pub struct SetEvent {
27    pub key: String,
28    pub value: Value,
29    pub result: async_channel::Sender<Result<()>>,
30}
31
32#[derive(Debug, Clone)]
33pub struct GetEvent {
34    pub key: String,
35    pub result: async_channel::Sender<Result<Option<Value>>>,
36}
37
38impl DistributedKVStore {
39    pub fn new(logger: &Logger, bridge: &EventBridge) -> Self {
40        DistributedKVStore {
41            logger: logger.clone(),
42            bridge: bridge.clone(),
43        }
44    }
45
46    pub async fn get(&self, key: String) -> Result<Option<Value>> {
47        let (tx, rx) = async_channel::bounded(1);
48        if let Err(e) = self
49            .bridge
50            .publish(crate::bridge::PublishEvents::GetEvent(GetEvent {
51                key,
52                result: tx.clone(),
53            }))
54            .await
55        {
56            error!(self.logger, "Failed to send get event"; "error" => ?e);
57            bail!(e)
58        };
59        rx.recv().await?
60    }
61
62    pub async fn set(&self, key: String, value: Value) -> Result<()> {
63        let (tx, rx) = async_channel::bounded(1);
64        if let Err(e) = self
65            .bridge
66            .publish(crate::bridge::PublishEvents::SetEvent(SetEvent {
67                key: key.clone(),
68                value,
69                result: tx.clone(),
70            }))
71            .await
72        {
73            error!(self.logger, "Failed to send set event"; "error" => ?e);
74            bail!(e)
75        };
76        if let Err(e) = rx.recv().await {
77            error!(self.logger, "Failed to set key"; "error" => ?e);
78            bail!(e)
79        } else {
80            info!(self.logger, "Set key"; "key" => key);
81        };
82        Ok(())
83    }
84}