1use std::{collections::HashMap, sync::Mutex};
2
3use anyhow::{Context, Result, ensure};
4use async_trait::async_trait;
5use serde::{Deserialize, Serialize};
6
7use crate::{actor_state::ActorStorageKey, host::HostId, postgres::PostgresDatabase};
8
9#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
10pub struct ObjectPlacement {
11 pub object: ActorStorageKey,
12 pub owner: HostId,
13 pub owner_epoch: u64,
14 pub home_region: String,
15}
16
17#[derive(Clone, Debug, PartialEq, Eq)]
18pub enum PlacementClaim {
19 Acquired(ObjectPlacement),
20 Current(ObjectPlacement),
21}
22
23#[async_trait]
24pub trait ObjectPlacementStore: Send + Sync {
25 async fn get(&self, object: &ActorStorageKey) -> Result<Option<ObjectPlacement>>;
26
27 async fn claim(
28 &self,
29 object: &ActorStorageKey,
30 expected: Option<&ObjectPlacement>,
31 owner: &HostId,
32 home_region: &str,
33 ) -> Result<PlacementClaim>;
34}
35
36#[derive(Default)]
37pub struct LocalObjectPlacementStore {
38 placements: Mutex<HashMap<ActorStorageKey, ObjectPlacement>>,
39}
40
41#[async_trait]
42impl ObjectPlacementStore for LocalObjectPlacementStore {
43 async fn get(&self, object: &ActorStorageKey) -> Result<Option<ObjectPlacement>> {
44 Ok(self
45 .placements
46 .lock()
47 .map_err(|_| anyhow::anyhow!("object placement lock poisoned"))?
48 .get(object)
49 .cloned())
50 }
51
52 async fn claim(
53 &self,
54 object: &ActorStorageKey,
55 expected: Option<&ObjectPlacement>,
56 owner: &HostId,
57 home_region: &str,
58 ) -> Result<PlacementClaim> {
59 validate_region(home_region)?;
60 let mut placements = self
61 .placements
62 .lock()
63 .map_err(|_| anyhow::anyhow!("object placement lock poisoned"))?;
64 match placements.get(object) {
65 None if expected.is_none() => {
66 let placement = ObjectPlacement {
67 object: object.clone(),
68 owner: owner.clone(),
69 owner_epoch: 1,
70 home_region: home_region.to_owned(),
71 };
72 placements.insert(object.clone(), placement.clone());
73 Ok(PlacementClaim::Acquired(placement))
74 }
75 Some(current) if expected == Some(current) => {
76 ensure!(
77 current.home_region == home_region,
78 "object home region cannot change"
79 );
80 if ¤t.owner == owner {
81 return Ok(PlacementClaim::Current(current.clone()));
82 }
83 let placement = ObjectPlacement {
84 object: object.clone(),
85 owner: owner.clone(),
86 owner_epoch: current
87 .owner_epoch
88 .checked_add(1)
89 .context("object owner epoch overflow")?,
90 home_region: home_region.to_owned(),
91 };
92 placements.insert(object.clone(), placement.clone());
93 Ok(PlacementClaim::Acquired(placement))
94 }
95 Some(current) => Ok(PlacementClaim::Current(current.clone())),
96 None => anyhow::bail!("expected object placement no longer exists"),
97 }
98 }
99}
100
101pub struct PostgresObjectPlacementStore {
102 database: PostgresDatabase,
103}
104
105impl PostgresObjectPlacementStore {
106 pub async fn connect(url: &str) -> Result<Self> {
107 Ok(Self::from_database(PostgresDatabase::connect(url).await?))
108 }
109
110 pub(crate) fn from_database(database: PostgresDatabase) -> Self {
111 Self { database }
112 }
113}
114
115#[async_trait]
116impl ObjectPlacementStore for PostgresObjectPlacementStore {
117 async fn get(&self, object: &ActorStorageKey) -> Result<Option<ObjectPlacement>> {
118 let row = self
119 .database
120 .client()
121 .query_opt(
122 "SELECT owner_host_id, owner_epoch, home_region \
123 FROM durable_object_placements WHERE object_id = $1",
124 &[&object.as_str()],
125 )
126 .await
127 .context("load PostgreSQL object placement")?;
128 row.map(|row| placement_from_row(object, &row)).transpose()
129 }
130
131 async fn claim(
132 &self,
133 object: &ActorStorageKey,
134 expected: Option<&ObjectPlacement>,
135 owner: &HostId,
136 home_region: &str,
137 ) -> Result<PlacementClaim> {
138 object.validate()?;
139 validate_region(home_region)?;
140 if expected.is_none() {
141 if let Some(row) = self
142 .database
143 .client()
144 .query_opt(
145 "INSERT INTO durable_object_placements \
146 (object_id, owner_host_id, owner_epoch, home_region) \
147 VALUES ($1, $2, 1, $3) ON CONFLICT DO NOTHING \
148 RETURNING owner_host_id, owner_epoch, home_region",
149 &[&object.as_str(), &owner.as_str(), &home_region],
150 )
151 .await
152 .context("insert PostgreSQL object placement")?
153 {
154 return Ok(PlacementClaim::Acquired(placement_from_row(object, &row)?));
155 }
156 return self.current_claim(object).await;
157 }
158
159 let expected = expected.expect("checked above");
160 ensure!(
161 expected.object == *object && expected.home_region == home_region,
162 "expected object placement does not match the claim"
163 );
164 if &expected.owner == owner {
165 return self.current_claim(object).await;
166 }
167 let expected_epoch = i64::try_from(expected.owner_epoch)
168 .context("object owner epoch exceeds PostgreSQL BIGINT")?;
169 if let Some(row) = self
170 .database
171 .client()
172 .query_opt(
173 "UPDATE durable_object_placements \
174 SET owner_host_id = $2, owner_epoch = owner_epoch + 1, updated_at = clock_timestamp() \
175 WHERE object_id = $1 AND owner_host_id = $3 AND owner_epoch = $4 AND home_region = $5 \
176 RETURNING owner_host_id, owner_epoch, home_region",
177 &[
178 &object.as_str(),
179 &owner.as_str(),
180 &expected.owner.as_str(),
181 &expected_epoch,
182 &home_region,
183 ],
184 )
185 .await
186 .context("claim PostgreSQL object placement")?
187 {
188 return Ok(PlacementClaim::Acquired(placement_from_row(object, &row)?));
189 }
190 self.current_claim(object).await
191 }
192}
193
194impl PostgresObjectPlacementStore {
195 async fn current_claim(&self, object: &ActorStorageKey) -> Result<PlacementClaim> {
196 self.get(object)
197 .await?
198 .map(PlacementClaim::Current)
199 .context("object placement disappeared during claim")
200 }
201}
202
203fn placement_from_row(
204 object: &ActorStorageKey,
205 row: &tokio_postgres::Row,
206) -> Result<ObjectPlacement> {
207 Ok(ObjectPlacement {
208 object: object.clone(),
209 owner: HostId::new(row.get::<_, String>(0)),
210 owner_epoch: u64::try_from(row.get::<_, i64>(1))
211 .context("PostgreSQL object owner epoch is negative")?,
212 home_region: row.get(2),
213 })
214}
215
216pub fn validate_region(region: &str) -> Result<()> {
217 ensure!(
218 !region.is_empty()
219 && region.len() <= 64
220 && region.bytes().all(|byte| {
221 byte.is_ascii_lowercase()
222 || byte.is_ascii_digit()
223 || matches!(byte, b'.' | b'_' | b'-')
224 }),
225 "sandbox region is invalid"
226 );
227 Ok(())
228}
229
230#[cfg(test)]
231mod tests {
232 use super::*;
233
234 #[tokio::test]
235 async fn claims_once_and_increments_epoch_on_transfer() -> Result<()> {
236 let store = LocalObjectPlacementStore::default();
237 let object = ActorStorageKey::new("object.v1.project.Counter.one");
238 let first = match store
239 .claim(&object, None, &HostId::new("host-a"), "us-east")
240 .await?
241 {
242 PlacementClaim::Acquired(placement) => placement,
243 claim => anyhow::bail!("unexpected claim: {claim:?}"),
244 };
245 assert_eq!(first.owner_epoch, 1);
246
247 let second = match store
248 .claim(&object, Some(&first), &HostId::new("host-b"), "us-east")
249 .await?
250 {
251 PlacementClaim::Acquired(placement) => placement,
252 claim => anyhow::bail!("unexpected claim: {claim:?}"),
253 };
254 assert_eq!(second.owner, HostId::new("host-b"));
255 assert_eq!(second.owner_epoch, 2);
256 assert_eq!(second.home_region, "us-east");
257 Ok(())
258 }
259
260 #[tokio::test]
261 async fn stale_claim_observes_the_current_owner() -> Result<()> {
262 let store = LocalObjectPlacementStore::default();
263 let object = ActorStorageKey::new("object.v1.project.Counter.one");
264 let PlacementClaim::Acquired(first) = store
265 .claim(&object, None, &HostId::new("host-a"), "us-east")
266 .await?
267 else {
268 anyhow::bail!("first claim was not acquired")
269 };
270 let PlacementClaim::Acquired(second) = store
271 .claim(&object, Some(&first), &HostId::new("host-b"), "us-east")
272 .await?
273 else {
274 anyhow::bail!("second claim was not acquired")
275 };
276 assert_eq!(
277 store
278 .claim(&object, Some(&first), &HostId::new("host-c"), "us-east",)
279 .await?,
280 PlacementClaim::Current(second)
281 );
282 Ok(())
283 }
284}