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}