Skip to main content

mkit_server/pipeline/
list_repos.rs

1//! Coordinator-owned visibility indexes and bounded namespace listing.
2
3use std::collections::VecDeque;
4
5use subtle::ConstantTimeEq;
6
7use super::{Authenticated, Authorizer, HookSet, Pipeline, RepoVisibility, internal, meta_error};
8use crate::error::ServerError;
9use crate::op::{AuthzFacts, CallerView, OpKind};
10use crate::policy::AuthorizerRole;
11use crate::repo::{Addressing, MAX_REPO_NAME_BYTES, RepoName};
12use crate::store::{
13    Batch, Cursor, Key, MultipartBlobStore, NamespaceStore, Partition, Value, codec, keys,
14};
15
16/// A repository name and its effective visibility. The registry has no cheap head/update projection.
17#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct RepoEntry {
19    /// Full name within the requested namespace.
20    pub name: String,
21    /// Explicit visibility or the current deployment default.
22    pub visibility: RepoVisibility,
23}
24
25/// One bounded repository-name page.
26#[derive(Debug, Clone, PartialEq, Eq)]
27pub struct RepoPage {
28    /// Entries in ascending byte order.
29    pub repos: Vec<RepoEntry>,
30    /// Opaque authenticated bytes; the binding uses unpadded base64url.
31    pub next: Option<Vec<u8>>,
32}
33
34/// Append the listing projection to the very apply that writes registration/visibility.
35pub(super) fn index_writes(
36    batch: &mut Batch,
37    repo: &RepoName,
38    registered: bool,
39    stored: Option<&Value>,
40) -> Result<(), ServerError> {
41    let visibility = stored
42        .map(codec::decode_repo_visibility)
43        .transpose()
44        .map_err(meta_error)?;
45    batch
46        .writes
47        .push(crate::store::Write::Delete(keys::repo_listing(repo, true)));
48    batch
49        .writes
50        .push(crate::store::Write::Delete(keys::repo_listing(repo, false)));
51    if registered {
52        let explicit_public = match visibility.map(|row| row.visibility) {
53            None => false,
54            Some(codec::StoredVisibility::Public) => true,
55            Some(codec::StoredVisibility::Private) => return Ok(()),
56        };
57        batch.writes.push(crate::store::Write::Put(
58            keys::repo_listing(repo, explicit_public),
59            Value::default(),
60        ));
61    }
62    Ok(())
63}
64
65struct Source {
66    prefix: Key,
67    start: Key,
68    cursor: Option<Cursor>,
69    last: Option<String>,
70    rows: VecDeque<String>,
71    more: bool,
72    registry: bool,
73}
74
75impl Source {
76    async fn fill<N: NamespaceStore>(
77        &mut self,
78        store: &N,
79        partition: &Partition,
80        calls_left: &mut usize,
81    ) -> Result<(), ServerError> {
82        if !self.rows.is_empty() || !self.more {
83            return Ok(());
84        }
85        // The opaque backend cursor belongs to this fixed range for the entire request.
86        let mut end = self.prefix.as_bytes().to_vec();
87        end.push(0xff);
88        if *calls_left == 0 {
89            return Err(ServerError::resource_exhausted(
90                "repository listing call budget exceeded",
91            ));
92        }
93        *calls_left -= 1;
94        let page = store
95            .scan(
96                partition,
97                &self.start,
98                &Key::new(end),
99                self.cursor.as_ref(),
100                101,
101            )
102            .await
103            .map_err(listing_unavailable)?;
104        if page.entries.len() > 101 {
105            return Err(listing_unavailable("invalid repository listing scan"));
106        }
107        let mut previous = self.last.clone();
108        for (key, value) in page.entries {
109            if !key.as_bytes().starts_with(self.prefix.as_bytes())
110                || key.as_bytes() < self.start.as_bytes()
111            {
112                return Err(listing_unavailable("repository listing key outside range"));
113            }
114            let name = match (self.registry, keys::parse(&key)) {
115                (true, Some(keys::ParsedKey::RepoRecord(repo))) => {
116                    codec::decode_repo_record(&value).map_err(listing_unavailable)?;
117                    repo.as_str().to_owned()
118                }
119                (false, Some(keys::ParsedKey::RepoListing { repo, .. }))
120                    if value.as_bytes().is_empty() =>
121                {
122                    repo.as_str().to_owned()
123                }
124                _ => return Err(listing_unavailable("invalid repository listing row")),
125            };
126            if previous.as_ref().is_some_and(|last| last >= &name) {
127                return Err(listing_unavailable("unordered repository listing"));
128            }
129            previous = Some(name.clone());
130            self.rows.push_back(name);
131        }
132        self.last = previous;
133        self.more = page.next.is_some();
134        self.cursor = page.next;
135        Ok(())
136    }
137}
138
139impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
140    pub(super) async fn plan_listing_visibility(
141        &self,
142        p: &Partition,
143        repo: &crate::repo::RepoId,
144        batch: &mut Batch,
145    ) -> Result<(), ServerError> {
146        let key = keys::repo_record(&repo.name);
147        let record = self.meta.get(p, &key).await.map_err(meta_error)?;
148        if let Some(value) = &record {
149            codec::decode_repo_record(value).map_err(meta_error)?;
150        }
151        batch
152            .preconditions
153            .push(super::lease::observed_guard(key, record.as_ref()));
154        let value = batch
155            .writes
156            .iter()
157            .rev()
158            .find_map(|write| match write {
159                crate::store::Write::Put(key, value)
160                    if *key == keys::repo_visibility(&repo.name) =>
161                {
162                    Some(value.clone())
163                }
164                _ => None,
165            })
166            .ok_or_else(|| internal("visibility projection without visibility write"))?;
167        index_writes(batch, &repo.name, record.is_some(), Some(&value))
168    }
169
170    /// List a namespace without visiting private rows for the public view.
171    /// Signed envelopes use the ordinary X-Repository syntax; its namespace
172    /// must match `namespace`. The name is only a namespace selector for this RPC.
173    /// Grants never extend listing rights. Check and Authority hooks concern the entire namespace.
174    ///
175    /// # Errors
176    /// Invalid namespace/prefix/size/token, mismatched authentication, or unavailable storage/hooks.
177    #[allow(clippy::too_many_lines)] // One bounded merge and its authenticated continuation.
178    pub async fn list_repos_page(
179        &self,
180        a: &Authenticated,
181        namespace: &str,
182        name_prefix: &str,
183        page_size: Option<u32>,
184        token: Option<&[u8]>,
185    ) -> Result<RepoPage, ServerError> {
186        self.observe(
187            a,
188            self.list_repos_inner(a, namespace, name_prefix, page_size, token),
189        )
190        .await
191    }
192
193    #[allow(clippy::too_many_lines)] // One bounded merge and its authenticated continuation.
194    async fn list_repos_inner(
195        &self,
196        a: &Authenticated,
197        namespace: &str,
198        name_prefix: &str,
199        page_size: Option<u32>,
200        token: Option<&[u8]>,
201    ) -> Result<RepoPage, ServerError> {
202        if a.write_grant.is_some() && a.auth.is_none() {
203            return Err(ServerError::unauthenticated(
204                "grant requires signed authentication",
205            ));
206        }
207        if name_prefix.len() > MAX_REPO_NAME_BYTES
208            || !name_prefix.bytes().all(|b| (0x21..=0x7e).contains(&b))
209        {
210            return Err(ServerError::invalid_argument(
211                "invalid repository listing prefix",
212            ));
213        }
214        let mut op = self.identify(
215            a,
216            OpKind::ListRepos {
217                name_prefix: name_prefix.into(),
218            },
219        )?;
220        if namespace != op.repo.namespace.as_str() {
221            return Err(ServerError::invalid_argument(
222                "invalid repository listing namespace or prefix",
223            ));
224        }
225        let size = match page_size.unwrap_or(0) {
226            0 => 100,
227            n @ 1..=100 => n as usize,
228            _ => {
229                return Err(ServerError::invalid_argument(
230                    "repository page size exceeds 100",
231                ));
232            }
233        };
234        if let Addressing::Single { repo } = &self.cfg.addressing {
235            self.hooks
236                .authorizer()
237                .authorize(&op)
238                .await
239                .map_err(ServerError::strip_admission_shape)?;
240            if token.is_some() {
241                return Err(invalid_token());
242            }
243            return Ok(RepoPage {
244                repos: if repo.name.as_str().starts_with(name_prefix) {
245                    vec![RepoEntry {
246                        name: repo.name.as_str().into(),
247                        visibility: RepoVisibility::Public,
248                    }]
249                } else {
250                    Vec::new()
251                },
252                next: None,
253            });
254        }
255        let owner = a.auth.is_some()
256            && a.write_grant.is_none()
257            && matches!(mkit_core::repo_identity::Namespace::parse(namespace),
258                Ok(mkit_core::repo_identity::Namespace::Ed25519(key)) if a.principal.ed25519() == Some(&key));
259        let mut full = owner;
260        let authority = self.cfg.authorizer_role == AuthorizerRole::Authority;
261        if !authority || (a.auth.is_some() && a.write_grant.is_none()) {
262            op.authz = AuthzFacts {
263                owner,
264                caller_view: if owner {
265                    CallerView::Writer
266                } else if a.auth.is_none() {
267                    CallerView::Anonymous
268                } else {
269                    CallerView::Reader
270                },
271                ..AuthzFacts::default()
272            };
273            match self.hooks.authorizer().authorize(&op).await {
274                Ok(facts) => {
275                    full |= authority
276                        && self.cfg.list_repos_authority_full
277                        && facts.caller_view == CallerView::Writer;
278                }
279                Err(error)
280                    if authority
281                        && matches!(
282                            error.code(),
283                            crate::Code::PermissionDenied | crate::Code::NotFound
284                        ) =>
285                {
286                    full = false;
287                }
288                Err(error) => return Err(error.strip_admission_shape()),
289            }
290        }
291        let keyset = self.cfg.ticket_keys.as_ref().ok_or_else(|| {
292            ServerError::failed_precondition("repository listing requires deployment MAC keys")
293        })?;
294        let binding = blake3::hash(
295            &serde_json::to_vec(&(
296                "ListRepos v1",
297                namespace,
298                name_prefix,
299                full,
300                self.cfg.default_repo_visibility == RepoVisibility::Public,
301                a.principal.ed25519(),
302                match &self.cfg.auth {
303                    super::AuthMode::AuthV2(cfg) => cfg.audience(),
304                    _ => "",
305                },
306            ))
307            .map_err(|_| internal("repository listing binding"))?,
308        );
309        let last = token
310            .map(|bytes| decode_token(keyset, bytes, binding.as_bytes()))
311            .transpose()?;
312        if last
313            .as_ref()
314            .is_some_and(|last| !last.starts_with(name_prefix))
315        {
316            return Err(invalid_token());
317        }
318        let prefixes = if full {
319            vec![Key::new(b"rr\0".to_vec())]
320        } else if self.cfg.default_repo_visibility == RepoVisibility::Public {
321            vec![
322                keys::repo_listing_prefix(true),
323                keys::repo_listing_prefix(false),
324            ]
325        } else {
326            vec![keys::repo_listing_prefix(true)]
327        };
328        let mut sources: Vec<_> = prefixes
329            .into_iter()
330            .map(|mut prefix| {
331                let base = prefix.as_bytes().to_vec();
332                let mut bytes = base;
333                bytes.extend_from_slice(name_prefix.as_bytes());
334                prefix = Key::new(bytes);
335                let mut start = prefix.as_bytes().to_vec();
336                if let Some(last) = &last {
337                    start.extend_from_slice(&last.as_bytes()[name_prefix.len()..]);
338                    start.push(0);
339                }
340                Source {
341                    prefix,
342                    start: Key::new(start),
343                    cursor: None,
344                    last: last.clone(),
345                    rows: VecDeque::new(),
346                    more: true,
347                    registry: full,
348                }
349            })
350            .collect();
351        let partition = self.shards.coordinator(&op.repo.namespace);
352        // Reserve one top-level call for the full view's batched visibility read.
353        let mut calls_left = if full { 101 } else { 102 };
354        let mut names = Vec::with_capacity(size);
355        loop {
356            for source in &mut sources {
357                while source.rows.is_empty() && source.more {
358                    source.fill(&self.meta, &partition, &mut calls_left).await?;
359                }
360            }
361            let next = sources
362                .iter()
363                .enumerate()
364                .filter_map(|(i, source)| source.rows.front().map(|name| (i, name)))
365                .min_by(|a, b| a.1.cmp(b.1));
366            let Some((index, _)) = next else {
367                break;
368            };
369            if names.len() == size {
370                break;
371            }
372            let name = sources[index]
373                .rows
374                .pop_front()
375                .ok_or_else(|| internal("repository listing merge"))?;
376            if names.last().is_some_and(|last| last >= &name) {
377                return Err(internal("duplicate repository listing row"));
378            }
379            names.push(name);
380        }
381        let more = sources.iter().any(|source| !source.rows.is_empty());
382        let visibility = if full {
383            let wanted: Vec<_> = names
384                .iter()
385                .map(|name| RepoName::new(name.clone()).map(|repo| keys::repo_visibility(&repo)))
386                .collect::<Result<_, _>>()?;
387            let rows = if wanted.is_empty() {
388                Vec::new()
389            } else {
390                self.meta
391                    .get_many(&partition, &wanted)
392                    .await
393                    .map_err(listing_unavailable)?
394            };
395            if rows.len() != names.len() {
396                return Err(listing_unavailable("short repository visibility read"));
397            }
398            rows.into_iter()
399                .map(|row| {
400                    let row = row
401                        .as_ref()
402                        .map(codec::decode_repo_visibility)
403                        .transpose()
404                        .map_err(listing_unavailable)?;
405                    Ok(
406                        if super::repo_is_private(row.as_ref(), self.cfg.default_repo_visibility) {
407                            RepoVisibility::Private
408                        } else {
409                            RepoVisibility::Public
410                        },
411                    )
412                })
413                .collect::<Result<Vec<_>, ServerError>>()?
414        } else {
415            vec![RepoVisibility::Public; names.len()]
416        };
417        let next = if more {
418            let last = names
419                .last()
420                .ok_or_else(|| internal("empty repository continuation"))?;
421            Some(encode_token(keyset, binding.as_bytes(), last)?)
422        } else {
423            None
424        };
425        Ok(RepoPage {
426            repos: names
427                .into_iter()
428                .zip(visibility)
429                .map(|(name, visibility)| RepoEntry { name, visibility })
430                .collect(),
431            next,
432        })
433    }
434}
435
436fn listing_unavailable(error: impl std::fmt::Display) -> ServerError {
437    tracing::warn!(detail = %error, "repository listing failed");
438    ServerError::unavailable("repository listing unavailable")
439}
440
441fn invalid_token() -> ServerError {
442    ServerError::invalid_argument("invalid repository page token")
443}
444
445fn encode_token(
446    keys: &crate::upload::token::TicketKeys,
447    binding: &[u8; 32],
448    last: &str,
449) -> Result<Vec<u8>, ServerError> {
450    let (id, key) = keys.listing_signing_key();
451    let mut bytes = vec![1, u8::try_from(id.len()).map_err(|_| invalid_token())?];
452    bytes.extend_from_slice(id.as_bytes());
453    bytes.extend_from_slice(binding);
454    bytes.extend_from_slice(last.as_bytes());
455    bytes.extend_from_slice(blake3::keyed_hash(&key, &bytes).as_bytes());
456    Ok(bytes)
457}
458
459fn decode_token(
460    keys: &crate::upload::token::TicketKeys,
461    bytes: &[u8],
462    binding: &[u8; 32],
463) -> Result<String, ServerError> {
464    if bytes.len() > 2 + 32 + 32 + MAX_REPO_NAME_BYTES + 32 || bytes.first() != Some(&1) {
465        return Err(invalid_token());
466    }
467    let id_len = usize::from(*bytes.get(1).ok_or_else(invalid_token)?);
468    let tag_at = bytes.len().checked_sub(32).ok_or_else(invalid_token)?;
469    let message = bytes.get(..tag_at).ok_or_else(invalid_token)?;
470    let id = message.get(2..2 + id_len).ok_or_else(invalid_token)?;
471    let key = keys
472        .listing_verification_key(id)
473        .ok_or_else(invalid_token)?;
474    if !bool::from(
475        blake3::keyed_hash(&key, message)
476            .as_bytes()
477            .as_slice()
478            .ct_eq(&bytes[tag_at..]),
479    ) || message.get(2 + id_len..2 + id_len + 32) != Some(binding.as_slice())
480    {
481        return Err(invalid_token());
482    }
483    let name = std::str::from_utf8(message.get(2 + id_len + 32..).ok_or_else(invalid_token)?)
484        .map_err(|_| invalid_token())?;
485    RepoName::new(name).map_err(|_| invalid_token())?;
486    Ok(name.into())
487}