little_durable_objects/
state_log.rs1use 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 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}