use std::time::{SystemTime, UNIX_EPOCH};
use axum::{
extract::{Query, State},
http::{header, StatusCode},
response::{IntoResponse, Response},
routing::post,
Extension, Json, Router,
};
use serde::{Deserialize, Serialize};
use crate::audit::AuditEntry;
use crate::auth::AuthenticatedUser;
use crate::registry_type::RegistryType;
use crate::tokens::Role;
use crate::AppState;
const REINDEX_MIN_INTERVAL_SECS: u64 = 10;
pub fn routes() -> Router<AppState> {
Router::new()
.route("/api/v1/admin/reindex", post(reindex))
.route("/api/v1/admin/tokens", post(admin_create_token))
}
#[derive(Debug, Default, Deserialize)]
struct ReindexQuery {
registry: Option<String>,
}
#[derive(Serialize)]
struct ReindexResponse {
status: &'static str,
scope: String,
}
fn now_epoch_secs() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0)
}
async fn reindex(
State(state): State<AppState>,
Extension(user): Extension<AuthenticatedUser>,
Query(query): Query<ReindexQuery>,
) -> Response {
let target = match query.registry.as_deref() {
Some(name) => match RegistryType::from_str_opt(name) {
Some(rt) => Some(rt),
None => {
return (StatusCode::BAD_REQUEST, format!("unknown registry: {name}"))
.into_response()
}
},
None => None,
};
if let Err(retry_after) = state
.repo_index
.try_accept_reindex(now_epoch_secs(), REINDEX_MIN_INTERVAL_SECS)
{
return (
StatusCode::TOO_MANY_REQUESTS,
[(header::RETRY_AFTER, retry_after.to_string())],
"reindex debounced; retry later",
)
.into_response();
}
let scope = target.map(|rt| rt.as_str().to_string());
match target {
Some(rt) => state.repo_index.invalidate(rt.as_str()),
None => state.repo_index.invalidate_all(),
}
state.audit.log(AuditEntry::new(
"reindex",
&user.0,
"",
scope.as_deref().unwrap_or("all"),
"",
));
tracing::info!(
actor = %user.0,
scope = scope.as_deref().unwrap_or("all"),
"admin reindex triggered"
);
let repo_index = state.repo_index.clone();
let storage = state.storage.clone();
let cancel = state.cancel_token;
let targets: Vec<RegistryType> = match target {
Some(rt) => vec![rt],
None => RegistryType::all().to_vec(),
};
tokio::spawn(async move {
tokio::select! {
_ = cancel.cancelled() => {}
_ = async {
for rt in targets {
let _ = repo_index.get(rt.as_str(), &storage).await;
}
} => {}
}
});
(
StatusCode::ACCEPTED,
Json(ReindexResponse {
status: "reindexing",
scope: scope.unwrap_or_else(|| "all".to_string()),
}),
)
.into_response()
}
#[derive(Deserialize)]
struct AdminCreateTokenRequest {
username: String,
#[serde(default = "default_admin_token_ttl")]
ttl_days: u64,
description: Option<String>,
#[serde(default = "default_admin_token_role")]
role: String,
}
fn default_admin_token_ttl() -> u64 {
30
}
fn default_admin_token_role() -> String {
"read".to_string()
}
#[derive(Serialize)]
struct AdminCreateTokenResponse {
token: String,
expires_in_days: u64,
}
async fn admin_create_token(
State(state): State<AppState>,
Extension(caller): Extension<AuthenticatedUser>,
Json(req): Json<AdminCreateTokenRequest>,
) -> Response {
let role = match req.role.as_str() {
"read" => Role::Read,
"write" => Role::Write,
"admin" => Role::Admin,
_ => {
return (
StatusCode::BAD_REQUEST,
"Invalid role. Use: read, write, admin",
)
.into_response()
}
};
if req.ttl_days == 0 {
return (StatusCode::BAD_REQUEST, "ttl_days must be > 0").into_response();
}
let Some(token_store) = &state.tokens else {
return (
StatusCode::SERVICE_UNAVAILABLE,
"Token storage not configured",
)
.into_response();
};
let role_label = role.to_string();
match token_store.create_token(&req.username, req.ttl_days, req.description, role) {
Ok(token) => {
let detail = format!("role={} ttl_days={}", role_label, req.ttl_days);
state.audit.log(AuditEntry::new(
"admin_token_mint",
&caller.0,
&req.username,
"admin",
&detail,
));
Json(AdminCreateTokenResponse {
token,
expires_in_days: req.ttl_days,
})
.into_response()
}
Err(e) => (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()).into_response(),
}
}
#[cfg(test)]
mod tests {
use crate::test_helpers::{
create_test_context_with_auth, send, send_with_headers, TestContext,
};
use crate::tokens::Role;
use axum::http::{Method, StatusCode};
use base64::{engine::general_purpose::STANDARD, Engine};
const URI: &str = "/api/v1/admin/reindex";
fn mint(ctx: &TestContext, role: Role) -> String {
ctx.state
.tokens
.as_ref()
.expect("token store enabled")
.create_token("tester", 30, None, role)
.expect("create token")
}
async fn post_bearer(ctx: &TestContext, uri: &str, token: &str) -> StatusCode {
let auth = format!("Bearer {token}");
send_with_headers(
&ctx.app,
Method::POST,
uri,
vec![("Authorization", &auth)],
"",
)
.await
.status()
}
async fn post_basic(ctx: &TestContext, uri: &str, cred: &str) -> StatusCode {
let auth = format!("Basic {}", STANDARD.encode(cred));
send_with_headers(
&ctx.app,
Method::POST,
uri,
vec![("Authorization", &auth)],
"",
)
.await
.status()
}
#[tokio::test]
async fn admin_token_accepted() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Admin);
assert_eq!(post_bearer(&ctx, URI, &tok).await, StatusCode::ACCEPTED);
}
#[tokio::test]
async fn write_token_forbidden() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Write);
assert_eq!(post_bearer(&ctx, URI, &tok).await, StatusCode::FORBIDDEN);
}
#[tokio::test]
async fn read_token_forbidden() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Read);
assert_eq!(post_bearer(&ctx, URI, &tok).await, StatusCode::FORBIDDEN);
}
#[tokio::test]
async fn no_auth_unauthorized() {
let ctx = create_test_context_with_auth(&[]);
let status = send(&ctx.app, Method::POST, URI, "").await.status();
assert_eq!(status, StatusCode::UNAUTHORIZED);
}
#[tokio::test]
async fn basic_auth_forbidden_fail_closed() {
let ctx = create_test_context_with_auth(&[("alice", "pw")]);
let cred = STANDARD.encode("alice:pw");
let auth = format!("Basic {cred}");
let status = send_with_headers(
&ctx.app,
Method::POST,
URI,
vec![("Authorization", &auth)],
"",
)
.await
.status();
assert_eq!(status, StatusCode::FORBIDDEN);
}
#[tokio::test]
async fn basic_password_write_token_forbidden() {
let ctx = create_test_context_with_auth(&[("alice", "pw")]);
let tok = mint(&ctx, Role::Write);
assert_eq!(
post_basic(&ctx, URI, &format!("x:{tok}")).await,
StatusCode::FORBIDDEN
);
}
#[tokio::test]
async fn basic_password_admin_token_accepted() {
let ctx = create_test_context_with_auth(&[("alice", "pw")]);
let tok = mint(&ctx, Role::Admin);
assert_eq!(
post_basic(&ctx, URI, &format!("x:{tok}")).await,
StatusCode::ACCEPTED
);
}
#[tokio::test]
async fn unknown_registry_bad_request() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Admin);
let status = post_bearer(&ctx, "/api/v1/admin/reindex?registry=crago", &tok).await;
assert_eq!(status, StatusCode::BAD_REQUEST);
}
#[tokio::test]
async fn scoped_registry_accepted() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Admin);
let status = post_bearer(&ctx, "/api/v1/admin/reindex?registry=cargo", &tok).await;
assert_eq!(status, StatusCode::ACCEPTED);
}
#[tokio::test]
async fn second_call_debounced() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Admin);
assert_eq!(post_bearer(&ctx, URI, &tok).await, StatusCode::ACCEPTED);
assert_eq!(
post_bearer(&ctx, URI, &tok).await,
StatusCode::TOO_MANY_REQUESTS
);
}
#[tokio::test]
async fn raw_reindex_stays_write_gated() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Write);
let status = post_bearer(&ctx, "/raw/-/reindex", &tok).await;
assert_eq!(status, StatusCode::OK);
}
const TOK_URI: &str = "/api/v1/admin/tokens";
async fn post_json_bearer(
ctx: &TestContext,
token: &str,
body: &str,
) -> axum::http::StatusCode {
let auth = format!("Bearer {token}");
send_with_headers(
&ctx.app,
Method::POST,
TOK_URI,
vec![
("Authorization", &auth),
("Content-Type", "application/json"),
],
body.to_string(),
)
.await
.status()
}
#[tokio::test]
async fn admin_mint_admin_token_end_to_end() {
use crate::test_helpers::body_bytes;
let ctx = create_test_context_with_auth(&[]);
let admin = mint(&ctx, Role::Admin);
let auth = format!("Bearer {admin}");
let resp = send_with_headers(
&ctx.app,
Method::POST,
TOK_URI,
vec![
("Authorization", &auth),
("Content-Type", "application/json"),
],
r#"{"username":"svc","role":"admin"}"#.to_string(),
)
.await;
assert_eq!(resp.status(), StatusCode::OK);
let json: serde_json::Value = serde_json::from_slice(&body_bytes(resp).await).unwrap();
let minted = json["token"].as_str().expect("token in response");
assert_eq!(
post_bearer(&ctx, "/api/v1/admin/reindex?registry=cargo", minted).await,
StatusCode::ACCEPTED
);
}
#[tokio::test]
async fn admin_mint_write_token_forbidden() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Write);
assert_eq!(
post_json_bearer(&ctx, &tok, r#"{"username":"svc","role":"admin"}"#).await,
StatusCode::FORBIDDEN
);
}
#[tokio::test]
async fn admin_mint_read_token_forbidden() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Read);
assert_eq!(
post_json_bearer(&ctx, &tok, r#"{"username":"svc","role":"admin"}"#).await,
StatusCode::FORBIDDEN
);
}
#[tokio::test]
async fn admin_mint_basic_no_role_forbidden() {
let ctx = create_test_context_with_auth(&[("alice", "pw")]);
assert_eq!(
post_basic(&ctx, TOK_URI, "alice:pw").await,
StatusCode::FORBIDDEN
);
}
#[tokio::test]
async fn admin_mint_no_auth_unauthorized() {
let ctx = create_test_context_with_auth(&[]);
let status = send(&ctx.app, Method::POST, TOK_URI, "").await.status();
assert_eq!(status, StatusCode::UNAUTHORIZED);
}
#[tokio::test]
async fn admin_mint_invalid_role_bad_request() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Admin);
assert_eq!(
post_json_bearer(&ctx, &tok, r#"{"username":"svc","role":"superuser"}"#).await,
StatusCode::BAD_REQUEST
);
}
#[tokio::test]
async fn admin_mint_zero_ttl_bad_request() {
let ctx = create_test_context_with_auth(&[]);
let tok = mint(&ctx, Role::Admin);
assert_eq!(
post_json_bearer(
&ctx,
&tok,
r#"{"username":"svc","role":"admin","ttl_days":0}"#
)
.await,
StatusCode::BAD_REQUEST
);
}
}