1use 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#[derive(Debug, Clone, PartialEq, Eq)]
18pub struct RepoEntry {
19 pub name: String,
21 pub visibility: RepoVisibility,
23}
24
25#[derive(Debug, Clone, PartialEq, Eq)]
27pub struct RepoPage {
28 pub repos: Vec<RepoEntry>,
30 pub next: Option<Vec<u8>>,
32}
33
34pub(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 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 #[allow(clippy::too_many_lines)] 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)] 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 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}