Skip to main content

vmgs_broker/
broker.rs

1// Copyright (c) Microsoft Corporation.
2// Licensed under the MIT License.
3
4use mesh::MeshPayload;
5use mesh::Receiver;
6use mesh::error::RemoteError;
7use mesh::payload::Protobuf;
8use mesh::rpc::Rpc;
9use thiserror::Error;
10use vmgs::Vmgs;
11use vmgs::VmgsFileInfo;
12use vmgs_format::FileId;
13
14/// An error returned by a VMGS broker operation.
15#[derive(Protobuf, Error, Debug)]
16pub enum VmgsBrokerError {
17    /// The requested file has no allocated bytes (i.e. does not exist).
18    #[error("no allocated bytes for file id being read")]
19    FileInfoNotAllocated,
20    /// Another VMGS error.
21    #[error(transparent)]
22    Other(RemoteError),
23}
24
25impl From<vmgs::Error> for VmgsBrokerError {
26    fn from(value: vmgs::Error) -> Self {
27        match value {
28            vmgs::Error::FileInfoNotAllocated(_) => VmgsBrokerError::FileInfoNotAllocated,
29            other => VmgsBrokerError::Other(RemoteError::new(other)),
30        }
31    }
32}
33
34#[derive(Protobuf)]
35pub struct BrokerFileId(u32);
36
37impl From<FileId> for BrokerFileId {
38    fn from(value: FileId) -> Self {
39        BrokerFileId(value.0)
40    }
41}
42
43impl From<BrokerFileId> for FileId {
44    fn from(value: BrokerFileId) -> Self {
45        FileId(value.0)
46    }
47}
48
49#[derive(MeshPayload)]
50pub enum VmgsBrokerRpc {
51    Inspect(inspect::Deferred),
52    GetFileInfo(Rpc<BrokerFileId, Result<VmgsFileInfo, VmgsBrokerError>>),
53    ReadFile(Rpc<BrokerFileId, Result<Vec<u8>, VmgsBrokerError>>),
54    WriteFile(Rpc<(BrokerFileId, Vec<u8>), Result<(), VmgsBrokerError>>),
55    #[cfg(feature = "encryption")]
56    WriteFileEncrypted(Rpc<(BrokerFileId, Vec<u8>), Result<(), VmgsBrokerError>>),
57    Save(Rpc<(), vmgs::save_restore::state::SavedVmgsState>),
58    DeleteFile(Rpc<BrokerFileId, Result<(), VmgsBrokerError>>),
59}
60
61pub struct VmgsBrokerTask {
62    vmgs: Vmgs,
63}
64
65impl VmgsBrokerTask {
66    /// Initialize the data store with the underlying block storage interface.
67    pub fn new(vmgs: Vmgs) -> VmgsBrokerTask {
68        VmgsBrokerTask { vmgs }
69    }
70
71    pub async fn run(&mut self, mut recv: Receiver<VmgsBrokerRpc>) {
72        loop {
73            match recv.recv().await {
74                Ok(message) => self.process_message(message).await,
75                Err(_) => return, // all mpsc senders went away
76            }
77        }
78    }
79
80    async fn process_message(&mut self, message: VmgsBrokerRpc) {
81        match message {
82            VmgsBrokerRpc::Inspect(req) => {
83                req.inspect(&self.vmgs);
84            }
85            VmgsBrokerRpc::GetFileInfo(rpc) => rpc
86                .handle_sync(|file_id| self.vmgs.get_file_info(file_id.into()).map_err(Into::into)),
87            VmgsBrokerRpc::ReadFile(rpc) => {
88                rpc.handle(async |file_id| {
89                    self.vmgs
90                        .read_file(file_id.into())
91                        .await
92                        .map_err(Into::into)
93                })
94                .await
95            }
96            VmgsBrokerRpc::WriteFile(rpc) => {
97                rpc.handle(async |(file_id, buf)| {
98                    self.vmgs
99                        .write_file(file_id.into(), &buf)
100                        .await
101                        .map_err(Into::into)
102                })
103                .await
104            }
105            #[cfg(feature = "encryption")]
106            VmgsBrokerRpc::WriteFileEncrypted(rpc) => {
107                rpc.handle(async |(file_id, buf)| {
108                    self.vmgs
109                        .write_file_encrypted(file_id.into(), &buf)
110                        .await
111                        .map_err(Into::into)
112                })
113                .await
114            }
115            VmgsBrokerRpc::Save(rpc) => rpc.handle_sync(|()| self.vmgs.save()),
116            VmgsBrokerRpc::DeleteFile(rpc) => {
117                rpc.handle(async |file_id| {
118                    self.vmgs
119                        .delete_file(file_id.into())
120                        .await
121                        .map_err(Into::into)
122                })
123                .await
124            }
125        }
126    }
127}