Skip to main content

little_durable_objects/
placement.rs

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 &current.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}