Skip to main content

little_durable_objects/
state_log.rs

1use anyhow::{Context, Result, ensure};
2use serde::{Deserialize, Serialize};
3use serde_json::Value;
4
5pub const MAX_ACTOR_STATE_BYTES: usize = 1024 * 1024;
6pub const MAX_STATE_LOG_BYTES: usize = 4 * 1024 * 1024;
7pub const MAX_STATE_LOG_RECORDS: usize = 64;
8
9#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
10#[serde(rename_all = "camelCase")]
11pub struct StateRecord {
12    pub state_version: u64,
13    pub owner_epoch: u64,
14    /// `None` is used only by an ownership claim before an actor has committed
15    /// its first state snapshot.
16    pub state: Option<Value>,
17}
18
19#[derive(Clone, Debug, Default, PartialEq)]
20pub struct StateLog {
21    records: Vec<StateRecord>,
22}
23
24#[derive(Clone, Copy, Debug, PartialEq, Eq)]
25pub enum StateAppend {
26    Changed,
27    Unchanged,
28}
29
30impl StateLog {
31    pub fn decode(bytes: &[u8]) -> Result<Self> {
32        ensure!(
33            bytes.len() <= MAX_STATE_LOG_BYTES + MAX_ACTOR_STATE_BYTES + 1024,
34            "actor state log is too large"
35        );
36        let mut records: Vec<StateRecord> = Vec::new();
37        for (index, line) in bytes.split(|byte| *byte == b'\n').enumerate() {
38            if line.is_empty() {
39                continue;
40            }
41            let record: StateRecord = serde_json::from_slice(line)
42                .with_context(|| format!("decode actor state record {}", index + 1))?;
43            validate_record(&record)?;
44            if let Some(previous) = records.last() {
45                ensure!(
46                    record.state_version > previous.state_version,
47                    "actor state versions must increase"
48                );
49            }
50            records.push(record);
51        }
52        ensure!(
53            records.len() <= MAX_STATE_LOG_RECORDS,
54            "actor state log contains too many records"
55        );
56        Ok(Self { records })
57    }
58
59    pub fn latest(&self) -> Option<&StateRecord> {
60        self.records.last()
61    }
62
63    pub fn latest_state(&self) -> Option<&Value> {
64        self.latest().and_then(|record| record.state.as_ref())
65    }
66
67    pub fn claim(&mut self, owner_epoch: u64) -> Result<bool> {
68        ensure!(owner_epoch > 0, "owner epoch must be positive");
69        match self.records.last_mut() {
70            Some(record) if record.owner_epoch == owner_epoch => Ok(false),
71            Some(record) => {
72                ensure!(
73                    owner_epoch > record.owner_epoch,
74                    "owner epoch must increase during a claim"
75                );
76                record.owner_epoch = owner_epoch;
77                Ok(true)
78            }
79            None => {
80                self.records.push(StateRecord {
81                    state_version: 0,
82                    owner_epoch,
83                    state: None,
84                });
85                Ok(true)
86            }
87        }
88    }
89
90    pub fn append(&mut self, owner_epoch: u64, state: Value) -> Result<StateAppend> {
91        ensure!(owner_epoch > 0, "owner epoch must be positive");
92        ensure!(state.is_object(), "actor state must be a JSON object");
93        let state_bytes = serde_json::to_vec(&state)?;
94        ensure!(
95            state_bytes.len() <= MAX_ACTOR_STATE_BYTES,
96            "actor state exceeds the {MAX_ACTOR_STATE_BYTES}-byte limit"
97        );
98        if self.latest_state() == Some(&state) {
99            return Ok(StateAppend::Unchanged);
100        }
101        if let Some(latest) = self.latest() {
102            ensure!(
103                latest.owner_epoch == owner_epoch,
104                "actor state owner epoch does not match the current claim"
105            );
106        }
107        let state_version = self
108            .latest()
109            .map_or(1, |record| record.state_version.saturating_add(1));
110        ensure!(state_version > 0, "actor state version overflow");
111        self.records.push(StateRecord {
112            state_version,
113            owner_epoch,
114            state: Some(state),
115        });
116        if self.records.len() >= MAX_STATE_LOG_RECORDS
117            || self.encode()?.len() >= MAX_STATE_LOG_BYTES
118        {
119            let latest = self.records.pop().expect("state append added a record");
120            self.records.clear();
121            self.records.push(latest);
122        }
123        Ok(StateAppend::Changed)
124    }
125
126    pub fn encode(&self) -> Result<Vec<u8>> {
127        let mut bytes = Vec::new();
128        for record in &self.records {
129            validate_record(record)?;
130            serde_json::to_writer(&mut bytes, record)?;
131            bytes.push(b'\n');
132        }
133        Ok(bytes)
134    }
135
136    pub fn record_count(&self) -> usize {
137        self.records.len()
138    }
139}
140
141fn validate_record(record: &StateRecord) -> Result<()> {
142    ensure!(record.owner_epoch > 0, "owner epoch must be positive");
143    match &record.state {
144        Some(state) => {
145            ensure!(
146                record.state_version > 0,
147                "committed actor state version must be positive"
148            );
149            ensure!(state.is_object(), "actor state must be a JSON object");
150            ensure!(
151                serde_json::to_vec(state)?.len() <= MAX_ACTOR_STATE_BYTES,
152                "actor state exceeds the {MAX_ACTOR_STATE_BYTES}-byte limit"
153            );
154        }
155        None => ensure!(
156            record.state_version == 0,
157            "only an uninitialized claim may omit state"
158        ),
159    }
160    Ok(())
161}
162
163#[cfg(test)]
164mod tests {
165    use super::*;
166    use serde_json::json;
167
168    #[test]
169    fn snapshots_round_trip_as_ndjson() -> Result<()> {
170        let mut log = StateLog::default();
171        assert!(log.claim(1)?);
172        assert_eq!(log.append(1, json!({"count": 1}))?, StateAppend::Changed);
173        assert_eq!(log.append(1, json!({"count": 1}))?, StateAppend::Unchanged);
174        assert_eq!(log.append(1, json!({"count": 2}))?, StateAppend::Changed);
175
176        let bytes = log.encode()?;
177        assert_eq!(bytes.iter().filter(|byte| **byte == b'\n').count(), 3);
178        assert!(String::from_utf8(bytes.clone())?.contains("\"stateVersion\":2"));
179        let decoded = StateLog::decode(&bytes)?;
180        assert_eq!(decoded.latest().map(|record| record.state_version), Some(2));
181        assert_eq!(decoded.latest_state(), Some(&json!({"count": 2})));
182        Ok(())
183    }
184
185    #[test]
186    fn ownership_claim_changes_epoch_without_changing_state_version() -> Result<()> {
187        let mut log = StateLog::default();
188        log.claim(3)?;
189        log.append(3, json!({"count": 7}))?;
190        assert!(log.claim(4)?);
191
192        let latest = log.latest().expect("latest state");
193        assert_eq!(latest.state_version, 1);
194        assert_eq!(latest.owner_epoch, 4);
195        assert_eq!(latest.state, Some(json!({"count": 7})));
196        Ok(())
197    }
198
199    #[test]
200    fn compacts_inline_at_the_record_limit() -> Result<()> {
201        let mut log = StateLog::default();
202        log.claim(1)?;
203        for count in 1..MAX_STATE_LOG_RECORDS {
204            log.append(1, json!({"count": count}))?;
205        }
206        assert_eq!(log.record_count(), 1);
207        assert_eq!(
208            log.latest().map(|record| record.state_version),
209            Some((MAX_STATE_LOG_RECORDS - 1) as u64)
210        );
211        Ok(())
212    }
213
214    #[test]
215    fn rejects_oversized_state() {
216        let mut log = StateLog::default();
217        log.claim(1).expect("claim");
218        let error = log
219            .append(1, json!({"value": "x".repeat(MAX_ACTOR_STATE_BYTES)}))
220            .expect_err("oversized state");
221        assert!(error.to_string().contains("actor state exceeds"));
222    }
223}