Skip to main content

agentplane/store/
postgres_registry.rs

1//! A durable manifest registry on `PostgreSQL`.
2//!
3//! The backend the immutability rule is actually about. On one node almost
4//! anything holds it up; the moment two instances publish concurrently, only
5//! the database can arbitrate — so the rule is enforced by a **unique primary
6//! key plus one transaction that reads and writes**, not by a `SELECT` this
7//! process makes a decision from and then hopes about.
8//!
9//! The decision itself is
10//! [`decide_publish`](crate::manifest::registry::decide_publish), shared with
11//! the embedded backend and with `MemoryRegistry`. Three hand-written copies of
12//! "may this replace what is there" are three chances to get it subtly
13//! different, and the one that is wrong is whichever nobody tested.
14//!
15//! The lock is a **transaction-scoped advisory lock on the key**, taken before
16//! the read. `SELECT … FOR UPDATE` cannot do this job: it locks a row that
17//! exists, and the race that matters is two publishes of different content
18//! both finding *no* row — both decide to insert, and whichever lands second
19//! either fails on the primary key with an error that reads like a backend
20//! fault, or (under an upsert) silently wins. The advisory lock serialises on
21//! the key whether or not the row is there yet, so the second publisher reads
22//! the first one's row and is refused as the immutability rule intends.
23
24use async_trait::async_trait;
25
26use crate::core::{Digest, KeyId, KeySignature, Signer, StoreError, Verifier};
27use crate::manifest::registry::{
28    PublishVerdict, check_signature, decide_publish, reparse, sign_manifest, to_yaml,
29};
30use crate::manifest::{Manifest, Registry, RegistryError};
31
32use super::postgres::PostgresStore;
33
34fn be(e: &tokio_postgres::Error) -> RegistryError {
35    RegistryError::Backend(super::postgres::detail(e))
36}
37
38fn pool_err(e: &impl std::fmt::Display) -> RegistryError {
39    RegistryError::Backend(e.to_string())
40}
41
42/// The stored signature, if the row carries one.
43///
44/// An empty key id is the absence: a signature cannot exist without one, so
45/// there is no sentinel here overlapping a value it might be confused with.
46fn stored_signature(key_id: &str, signature: &str) -> Option<KeySignature> {
47    if key_id.is_empty() {
48        return None;
49    }
50    Some(KeySignature {
51        key_id: key_id.to_owned(),
52        // A signature that does not decode cannot verify, and an empty one
53        // fails exactly as a wrong one does — `BadSignature` either way.
54        signature: hex::decode(signature).unwrap_or_default(),
55    })
56}
57
58impl PostgresStore {
59    /// Apply one publish: lock, decide and write in one transaction.
60    async fn publish_row(
61        &self,
62        manifest: &Manifest,
63        signature: Option<KeySignature>,
64    ) -> Result<Digest, RegistryError> {
65        let name = manifest.metadata.name.clone();
66        let version = manifest.metadata.version.clone();
67        let digest = manifest.digest().map_err(|source| RegistryError::Corrupt {
68            name: name.clone(),
69            version: version.clone(),
70            source,
71        })?;
72        let yaml = to_yaml(manifest)?;
73        let (key_id, signature_hex) = signature
74            .as_ref()
75            .map_or((String::new(), String::new()), |a| {
76                (a.key_id.clone(), hex::encode(&a.signature))
77            });
78
79        let mut client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
80        let tx = client.transaction().await.map_err(|e| be(&e))?;
81        let tenant = self.tenant_name();
82
83        // Serialise on the key before reading, whether or not a row exists.
84        // `hashtext` of the joined key is a 32-bit lock id; a collision between
85        // two unrelated keys only serialises two publishes that did not need
86        // it, never admits one that should have waited.
87        tx.execute(
88            "SELECT pg_advisory_xact_lock(hashtext($1 || E'\\x1f' || $2 || E'\\x1f' || $3))",
89            &[&tenant, &name, &version],
90        )
91        .await
92        .map_err(|e| be(&e))?;
93
94        let existing = tx
95            .query_opt(
96                "SELECT digest, key_id, signature FROM registry_manifests
97                 WHERE tenant = $1 AND name = $2 AND version = $3",
98                &[&tenant, &name, &version],
99            )
100            .await
101            .map_err(|e| be(&e))?;
102
103        let stored = match &existing {
104            Some(row) => {
105                let hex: String = row.get(0);
106                let stored_digest = Digest::from_hex(&hex).map_err(|_| {
107                    RegistryError::Backend(format!(
108                        "registry row '{name}' '{version}' holds '{hex}', which is not a digest"
109                    ))
110                })?;
111                Some((
112                    stored_digest,
113                    stored_signature(&row.get::<_, String>(1), &row.get::<_, String>(2)),
114                ))
115            }
116            None => None,
117        };
118
119        let verdict = decide_publish(
120            &name,
121            &version,
122            digest,
123            signature.as_ref(),
124            stored.as_ref().map(|(d, a)| (*d, a.as_ref())),
125        );
126        // The refusal rolls back rather than committing nothing: an open
127        // transaction returned to the pool is a lock held by whoever gets the
128        // connection next.
129        let verdict = match verdict {
130            Ok(verdict) => verdict,
131            Err(refusal) => {
132                let _ = tx.rollback().await;
133                return Err(refusal);
134            }
135        };
136
137        match verdict {
138            // A plain insert, never an upsert: the advisory lock makes a
139            // conflict here impossible, and an `ON CONFLICT DO UPDATE` would
140            // turn any future hole in that reasoning into a silent overwrite —
141            // the exact outcome the primary key exists to refuse.
142            PublishVerdict::Insert => {
143                tx.execute(
144                    "INSERT INTO registry_manifests (tenant, name, version, digest, yaml, key_id, signature)
145                     VALUES ($1, $2, $3, $4, $5, $6, $7)",
146                    &[
147                        &tenant,
148                        &name,
149                        &version,
150                        &digest.to_hex(),
151                        &yaml,
152                        &key_id,
153                        &signature_hex,
154                    ],
155                )
156                .await
157                .map_err(|e| be(&e))?;
158            }
159            // Only the signature columns move; the content is identical by
160            // construction, which is what `decide_publish` established.
161            PublishVerdict::AdoptSignature => {
162                tx.execute(
163                    "UPDATE registry_manifests SET key_id = $4, signature = $5
164                     WHERE tenant = $1 AND name = $2 AND version = $3",
165                    &[&tenant, &name, &version, &key_id, &signature_hex],
166                )
167                .await
168                .map_err(|e| be(&e))?;
169            }
170            PublishVerdict::Unchanged => {}
171        }
172        tx.commit().await.map_err(|e| be(&e))?;
173        Ok(digest)
174    }
175
176    /// One row, or [`RegistryError::NotFound`].
177    async fn row(
178        &self,
179        name: &str,
180        version: &str,
181    ) -> Result<(String, Option<KeySignature>), RegistryError> {
182        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
183        let row = client
184            .query_opt(
185                "SELECT yaml, key_id, signature FROM registry_manifests
186                 WHERE tenant = $1 AND name = $2 AND version = $3",
187                &[&self.tenant_name(), &name, &version],
188            )
189            .await
190            .map_err(|e| be(&e))?
191            .ok_or_else(|| RegistryError::NotFound {
192                name: name.to_owned(),
193                version: version.to_owned(),
194            })?;
195        Ok((
196            row.get(0),
197            stored_signature(&row.get::<_, String>(1), &row.get::<_, String>(2)),
198        ))
199    }
200}
201
202#[async_trait]
203impl Registry for PostgresStore {
204    async fn publish(&self, manifest: &Manifest) -> Result<Digest, RegistryError> {
205        self.publish_row(manifest, None).await
206    }
207
208    async fn publish_signed(
209        &self,
210        manifest: &Manifest,
211        signer: &dyn Signer,
212    ) -> Result<Digest, RegistryError> {
213        let (_, signature) = sign_manifest(manifest, signer)?;
214        self.publish_row(manifest, Some(signature)).await
215    }
216
217    async fn resolve(&self, name: &str, version: &str) -> Result<Manifest, RegistryError> {
218        let (yaml, _) = self.row(name, version).await?;
219        reparse(name, version, &yaml)
220    }
221
222    async fn resolve_verified(
223        &self,
224        name: &str,
225        version: &str,
226        verifier: &dyn Verifier,
227    ) -> Result<(Manifest, KeyId), RegistryError> {
228        let (yaml, signature) = self.row(name, version).await?;
229        let manifest = reparse(name, version, &yaml)?;
230        let key = check_signature(name, version, &manifest, signature.as_ref(), verifier)?;
231        Ok((manifest, key))
232    }
233
234    async fn versions(&self, name: &str) -> Result<Vec<String>, RegistryError> {
235        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
236        let rows = client
237            .query(
238                "SELECT version FROM registry_manifests
239                 WHERE tenant = $1 AND name = $2
240                 ORDER BY version",
241                &[&self.tenant_name(), &name],
242            )
243            .await
244            .map_err(|e| be(&e))?;
245        Ok(rows.into_iter().map(|row| row.get(0)).collect())
246    }
247
248    async fn names(&self) -> Result<Vec<String>, RegistryError> {
249        let client = self.pool_ref().get().await.map_err(|e| pool_err(&e))?;
250        let rows = client
251            .query(
252                "SELECT DISTINCT name FROM registry_manifests
253                 WHERE tenant = $1
254                 ORDER BY name",
255                &[&self.tenant_name()],
256            )
257            .await
258            .map_err(|e| be(&e))?;
259        Ok(rows.into_iter().map(|row| row.get(0)).collect())
260    }
261}
262
263/// Not a `StoreError`, and the conversion is here so a caller that has one can
264/// still say what it means.
265impl From<StoreError> for RegistryError {
266    fn from(e: StoreError) -> Self {
267        Self::Backend(e.to_string())
268    }
269}