massa_async_pool/
changes.rs

1//! Copyright (c) 2022 MASSA LABS <info@massa.net>
2
3//! This file provides structures representing changes to the asynchronous message pool
4use std::collections::{btree_map::Entry, BTreeMap};
5
6use massa_models::{
7    async_msg::{AsyncMessage, AsyncMessageUpdate},
8    async_msg_id::AsyncMessageId,
9    types::{Applicable, SetOrKeep, SetUpdateOrDelete},
10};
11
12use serde::{Deserialize, Serialize};
13use serde_with::serde_as;
14
15/// Consolidated changes to the asynchronous message pool
16#[serde_as]
17#[derive(Default, Debug, Clone, PartialEq, Eq, Deserialize, Serialize)]
18pub struct AsyncPoolChanges(
19    #[serde_as(as = "Vec<(_, _)>")]
20    pub  BTreeMap<AsyncMessageId, SetUpdateOrDelete<AsyncMessage, AsyncMessageUpdate>>,
21);
22
23impl Applicable<AsyncPoolChanges> for AsyncPoolChanges {
24    /// extends the current `AsyncPoolChanges` with another one
25    fn apply(&mut self, changes: AsyncPoolChanges) {
26        for (id, msg_change) in changes.0 {
27            match self.0.entry(id) {
28                Entry::Occupied(mut occ) => {
29                    // apply incoming change if a change on this entry already exists
30                    occ.get_mut().apply(msg_change);
31                }
32                Entry::Vacant(vac) => {
33                    // otherwise insert the incoming change
34                    vac.insert(msg_change);
35                }
36            }
37        }
38    }
39}
40
41impl AsyncPoolChanges {
42    /// Pushes a message addition to the list of changes.
43    /// No add/delete compensations are done.
44    ///
45    /// Arguments:
46    /// * `msg_id`: ID of the message to push as added to the list of changes
47    /// * `msg`: message to push as added to the list of changes
48    pub fn push_add(&mut self, msg_id: AsyncMessageId, msg: AsyncMessage) {
49        let mut change = AsyncPoolChanges::default();
50        change.0.insert(msg_id, SetUpdateOrDelete::Set(msg));
51        self.apply(change);
52    }
53
54    /// Pushes a message deletion to the list of changes.
55    /// No add/delete compensations are done.
56    ///
57    /// Arguments:
58    /// * `msg_id`: ID of the message to push as deleted to the list of changes
59    pub fn push_delete(&mut self, msg_id: AsyncMessageId) {
60        let mut change = AsyncPoolChanges::default();
61        change.0.insert(msg_id, SetUpdateOrDelete::Delete);
62        self.apply(change);
63    }
64
65    /// Pushes a message activation to the list of changes.
66    ///
67    /// Arguments:
68    /// * `msg_id`: ID of the message to push as ready to be executed to the list of changes
69    pub fn push_activate(&mut self, msg_id: AsyncMessageId) {
70        let mut change = AsyncPoolChanges::default();
71
72        let msg_update = AsyncMessageUpdate {
73            can_be_executed: SetOrKeep::Set(true),
74            ..Default::default()
75        };
76
77        change
78            .0
79            .insert(msg_id, SetUpdateOrDelete::Update(msg_update));
80        self.apply(change);
81    }
82}
83
84#[cfg(test)]
85mod tests {
86    use std::str::FromStr;
87
88    use massa_models::types::SetUpdateOrDelete;
89    use massa_models::{
90        address::Address, amount::Amount, async_msg::AsyncMessageTrigger, slot::Slot,
91    };
92
93    use assert_matches::assert_matches;
94
95    use super::*;
96
97    fn get_message() -> AsyncMessage {
98        AsyncMessage::new(
99            Slot::new(1, 0),
100            0,
101            Address::from_str("AU12dG5xP1RDEB5ocdHkymNVvvSJmUL9BgHwCksDowqmGWxfpm93x").unwrap(),
102            Address::from_str("AU12htxRWiEm8jDJpJptr6cwEhWNcCSFWstN1MLSa96DDkVM9Y42G").unwrap(),
103            String::from("test"),
104            10000000,
105            Amount::from_str("1").unwrap(),
106            Amount::from_str("1").unwrap(),
107            Slot::new(2, 0),
108            Slot::new(3, 0),
109            vec![1, 2, 3, 4],
110            Some(AsyncMessageTrigger {
111                address: Address::from_str("AU12dG5xP1RDEB5ocdHkymNVvvSJmUL9BgHwCksDowqmGWxfpm93x")
112                    .unwrap(),
113                datastore_key: Some(vec![1, 2, 3, 4]),
114            }),
115            None,
116        )
117    }
118
119    #[test]
120    fn test_pool_changes_push() {
121        // AsyncPoolChanges, push_add/push_delete/push_activate
122
123        let message = get_message();
124        assert!(!message.can_be_executed);
125
126        let mut changes = AsyncPoolChanges::default();
127
128        changes.push_add(message.compute_id(), message.clone());
129        assert_eq!(changes.0.len(), 1);
130        assert_matches!(
131            changes.0.get(&message.compute_id()),
132            Some(&SetUpdateOrDelete::Set(..))
133        );
134
135        changes.push_activate(message.compute_id());
136        assert_eq!(changes.0.len(), 1);
137        let value = changes.0.get(&message.compute_id()).unwrap();
138        match value {
139            SetUpdateOrDelete::Set(msg) => {
140                assert!(msg.can_be_executed);
141            }
142            _ => {
143                panic!("Unexpected value");
144            }
145        }
146
147        changes.push_delete(message.compute_id());
148        // Len is still 1, but value has changed
149        assert_eq!(changes.0.len(), 1);
150        assert_eq!(
151            changes.0.get(&message.compute_id()),
152            Some(&SetUpdateOrDelete::Delete)
153        );
154    }
155}