From 1d1cdf305084aab9b7f398458e83a9cfed3b1c8c Mon Sep 17 00:00:00 2001 From: Van Tan Minh Date: Sat, 3 Oct 2026 11:20:20 +0700 Subject: [PATCH] feat: connect Knotree Registry account once with owner-verified auto-deploy Co-Authored-By: Claude Opus 5.5 --- .../api/migrations/0022_registry_accounts.sql | 39 + apps/api/src/knotree_registry.rs | 99 ++- apps/api/src/lib.rs | 21 + apps/api/src/registry_accounts.rs | 742 ++++++++++++++++++ apps/api/src/registry_consent.rs | 23 +- docs/intakes/IN-047.md | 20 + docs/stories/US-048.md | 19 + 7 files changed, 941 insertions(+), 22 deletions(-) create mode 100644 apps/api/migrations/0022_registry_accounts.sql create mode 100644 apps/api/src/registry_accounts.rs create mode 100644 docs/intakes/IN-047.md create mode 100644 docs/stories/US-048.md diff --git a/apps/api/migrations/0022_registry_accounts.sql b/apps/api/migrations/0022_registry_accounts.sql new file mode 100644 index 0000000..4718247 --- /dev/null +++ b/apps/api/migrations/0022_registry_accounts.sql @@ -0,0 +1,39 @@ +-- Account-level Knotree Registry connection: one consent per Cloud user grants +-- pull access to that user's whole Registry namespace. Project connections +-- created from it reference the account and always use its live credential. +CREATE TABLE knotree_registry_accounts ( + id UUID PRIMARY KEY, + user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE, + issuer TEXT NOT NULL, + subject TEXT NOT NULL, + registry_username TEXT NOT NULL CHECK (char_length(registry_username) BETWEEN 1 AND 128), + credential_ciphertext TEXT NOT NULL, + delegated_credential_id UUID NOT NULL, + credential_expires_at TIMESTAMPTZ NOT NULL, + revoked_at TIMESTAMPTZ, + created_at TIMESTAMPTZ NOT NULL DEFAULT now(), + updated_at TIMESTAMPTZ NOT NULL DEFAULT now() +); + +CREATE UNIQUE INDEX knotree_registry_accounts_active_user + ON knotree_registry_accounts (user_id) WHERE revoked_at IS NULL; + +CREATE TABLE registry_account_consent_attempts ( + state_hash BYTEA PRIMARY KEY, + session_hash BYTEA NOT NULL, + user_id UUID NOT NULL REFERENCES users(id) ON DELETE CASCADE, + issuer TEXT NOT NULL, + subject TEXT NOT NULL, + verifier_ciphertext TEXT NOT NULL, + return_to TEXT, + expires_at TIMESTAMPTZ NOT NULL +); +CREATE INDEX registry_account_consent_attempts_expiry + ON registry_account_consent_attempts (expires_at); + +ALTER TABLE knotree_registry_connections + ADD COLUMN account_id UUID REFERENCES knotree_registry_accounts(id) ON DELETE CASCADE; + +CREATE UNIQUE INDEX knotree_registry_connections_account_repository + ON knotree_registry_connections (project_id, account_id, repository) + WHERE account_id IS NOT NULL AND revoked_at IS NULL; diff --git a/apps/api/src/knotree_registry.rs b/apps/api/src/knotree_registry.rs index a708bd1..523432c 100644 --- a/apps/api/src/knotree_registry.rs +++ b/apps/api/src/knotree_registry.rs @@ -95,6 +95,26 @@ struct RegistryEventMetadata { registry: Option, tagged_image: Option, is_tag: Option, + #[serde(default)] + owner_issuer: Option, + #[serde(default)] + owner_subject: Option, +} + +/// An account-derived connection auto-deploys only when Registry reports the +/// pushed namespace belongs to the same central identity that connected it. +/// Legacy per-repository connections keep their verified-at-consent scope. +fn owner_matches( + account_identity: Option<(&str, &str)>, + metadata: &RegistryEventMetadata, +) -> bool { + match account_identity { + None => true, + Some((issuer, subject)) => { + metadata.owner_issuer.as_deref() == Some(issuer) + && metadata.owner_subject.as_deref() == Some(subject) + } + } } #[derive(Debug, Clone, PartialEq, Eq)] @@ -219,7 +239,8 @@ pub async fn update_connection( let connection = sqlx::query_as::<_, RegistryConnectionRow>( "SELECT id, registry_username, repository, verified_at FROM knotree_registry_connections - WHERE id = $1 AND project_id = $2 AND revoked_at IS NULL", + WHERE id = $1 AND project_id = $2 AND revoked_at IS NULL + AND account_id IS NULL", ) .bind(connection_id) .bind(project_id) @@ -376,11 +397,17 @@ pub(crate) async fn load_credentials( project_id: Uuid, connection_id: Uuid, ) -> Result, AppError> { + // Connections created from an account always use the account's current + // credential, so reconnecting the account renews every derived project. let row = sqlx::query_as::<_, (String, String)>( - "SELECT registry_username, credential_ciphertext - FROM knotree_registry_connections - WHERE id = $1 AND project_id = $2 AND revoked_at IS NULL - AND (credential_expires_at IS NULL OR credential_expires_at > now())", + "SELECT connection.registry_username, + COALESCE(account.credential_ciphertext, connection.credential_ciphertext) + FROM knotree_registry_connections AS connection + LEFT JOIN knotree_registry_accounts AS account ON account.id = connection.account_id + WHERE connection.id = $1 AND connection.project_id = $2 AND connection.revoked_at IS NULL + AND (connection.credential_expires_at IS NULL OR connection.credential_expires_at > now()) + AND (connection.account_id IS NULL + OR (account.revoked_at IS NULL AND account.credential_expires_at > now()))", ) .bind(connection_id) .bind(project_id) @@ -735,24 +762,43 @@ async fn persist_registry_event( } let image_ref = immutable_image(repository, digest).expect("event repository and digest were validated"); - let services = sqlx::query_as::<_, (Uuid, String)>( - "SELECT id, image - FROM project_app_services - WHERE image_source = 'knotree_registry' - AND auto_deploy_enabled = TRUE - AND status IN ('ready', 'provisioning') - AND registry_connection_id IN ( - SELECT id FROM knotree_registry_connections WHERE revoked_at IS NULL - AND (credential_expires_at IS NULL OR credential_expires_at > now()) - )", + let services = sqlx::query_as::<_, (Uuid, String, String, Option, Option)>( + "SELECT service.id, service.image, connection.repository, + account.issuer, account.subject + FROM project_app_services AS service + JOIN knotree_registry_connections AS connection + ON connection.id = service.registry_connection_id + LEFT JOIN knotree_registry_accounts AS account ON account.id = connection.account_id + WHERE service.image_source = 'knotree_registry' + AND service.auto_deploy_enabled = TRUE + AND service.status IN ('ready', 'provisioning') + AND connection.revoked_at IS NULL + AND (connection.credential_expires_at IS NULL OR connection.credential_expires_at > now()) + AND (connection.account_id IS NULL + OR (account.revoked_at IS NULL AND account.credential_expires_at > now()))", ) .fetch_all(&mut *transaction) .await?; - for (service_id, image) in services { + for (service_id, image, connected_repository, account_issuer, account_subject) in services { let Some(target) = parse_registry_image(&image) else { continue; }; - if target.repository != repository || target.tag != tag { + if target.repository != repository + || target.tag != tag + || connected_repository != repository + { + continue; + } + let account_identity = match (account_issuer.as_deref(), account_subject.as_deref()) { + (Some(issuer), Some(subject)) => Some((issuer, subject)), + _ => None, + }; + if !owner_matches(account_identity, &event.metadata) { + tracing::warn!( + delivery_id = %delivery_id, + app_service_id = %service_id, + "ignoring Knotree Registry push whose owner does not match the connected account" + ); continue; } sqlx::query( @@ -850,6 +896,25 @@ mod tests { )); } + #[test] + fn account_connections_deploy_only_for_matching_owner() { + let metadata = |issuer: Option<&str>, subject: Option<&str>| RegistryEventMetadata { + registry: Some(REGISTRY_HOST.into()), + tagged_image: None, + is_tag: Some(true), + owner_issuer: issuer.map(Into::into), + owner_subject: subject.map(Into::into), + }; + let issuer = "https://accounts.knotree.com"; + let alice = Some((issuer, "alice")); + assert!(owner_matches(alice, &metadata(Some(issuer), Some("alice")))); + assert!(!owner_matches(alice, &metadata(Some(issuer), Some("bob")))); + assert!(!owner_matches(alice, &metadata(Some("https://evil.example"), Some("alice")))); + assert!(!owner_matches(alice, &metadata(None, None))); + // Legacy repository-scoped connections are unaffected. + assert!(owner_matches(None, &metadata(None, None))); + } + #[test] fn creates_registry_specific_docker_auth_without_returning_plaintext() { let config = docker_config_json("service-user", "pull-token").unwrap(); diff --git a/apps/api/src/lib.rs b/apps/api/src/lib.rs index f99e218..a016bb4 100644 --- a/apps/api/src/lib.rs +++ b/apps/api/src/lib.rs @@ -17,6 +17,7 @@ pub mod models; pub mod projects; pub mod public_access; pub mod redis_resources; +pub mod registry_accounts; pub mod registry_consent; pub mod resources; pub mod security; @@ -48,6 +49,26 @@ pub fn router(state: AppState) -> Router { let api = Router::new() .route("/auth/knotree-registry/callback", get(registry_consent::callback)) .route("/workspaces/{workspace_id}/projects/{project_slug}/registry-connections/authorize", post(registry_consent::start)) + .route( + "/workspaces/{workspace_id}/projects/{project_slug}/registry-connections/from-account", + post(registry_accounts::import_into_project), + ) + .route( + "/integrations/knotree-registry", + get(registry_accounts::status).delete(registry_accounts::disconnect), + ) + .route( + "/integrations/knotree-registry/authorize", + post(registry_accounts::start), + ) + .route( + "/integrations/knotree-registry/repositories", + get(registry_accounts::repositories), + ) + .route( + "/integrations/knotree-registry/repositories/{*repository}", + get(registry_accounts::repository_tags), + ) .route("/auth/sso/config", get(sso::configuration)) .route("/auth/sso/start", get(sso::start)) .route("/auth/sso/callback", get(sso::callback)) diff --git a/apps/api/src/registry_accounts.rs b/apps/api/src/registry_accounts.rs new file mode 100644 index 0000000..a6ab88e --- /dev/null +++ b/apps/api/src/registry_accounts.rs @@ -0,0 +1,742 @@ +//! Account-level Knotree Registry connection ("connect once", like installing +//! a GitHub App). One SSO-bound consent gives Cloud a pull-only credential for +//! the user's whole Registry namespace. The picker lists only that namespace, +//! and project connections created from it always use the account credential. + +use axum::{ + Json, + extract::{Path, State}, + http::{HeaderMap, StatusCode}, + response::{IntoResponse, Redirect, Response}, +}; +use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD}; +use serde::{Deserialize, Serialize}; +use sqlx::FromRow; +use time::OffsetDateTime; +use url::Url; +use uuid::Uuid; + +use crate::{ + auth, error::AppError, knotree_registry, projects, registry_consent, security, + state::AppState, +}; + +const REGISTRY: &str = "https://registry.knotree.com"; +const MAX_LIST_BYTES: usize = 1024 * 1024; + +fn invalid() -> AppError { + AppError::BadRequest { + code: "REGISTRY_CONSENT_FAILED", + message: "Registry authorization expired or could not be verified. Start again.", + } +} + +fn not_connected() -> AppError { + AppError::Conflict { + code: "KNOTREE_REGISTRY_ACCOUNT_REQUIRED", + message: "Connect your Knotree Registry account first.", + } +} + +fn no_store(mut response: Response) -> Response { + response + .headers_mut() + .insert("cache-control", "no-store".parse().unwrap()); + response +} + +/// Only same-site relative paths, so the callback can never redirect off Cloud. +fn safe_return_to(value: Option<&str>) -> Option { + let value = value?.trim(); + (value.starts_with('/') + && !value.starts_with("//") + && !value.contains('\\') + && value.len() <= 512 + && !value.chars().any(char::is_control)) + .then(|| value.to_owned()) +} + +async fn sso_subject(state: &AppState, user_id: Uuid, issuer: &str) -> Result, AppError> { + Ok( + sqlx::query_scalar("SELECT subject FROM sso_identities WHERE user_id=$1 AND issuer=$2") + .bind(user_id) + .bind(issuer) + .fetch_optional(&state.db) + .await?, + ) +} + +#[derive(Deserialize, Default)] +#[serde(rename_all = "camelCase")] +pub struct StartInput { + #[serde(default)] + return_to: Option, +} + +#[derive(Deserialize)] +struct Started { + request_id: Uuid, + authorization_url: String, +} + +/// POST /integrations/knotree-registry/authorize +pub async fn start( + State(state): State, + headers: HeaderMap, + input: Option>, +) -> Result { + security::require_csrf(&headers, &state.config)?; + let user = auth::authenticate(&state, &headers).await?; + let config = state.config.sso.as_ref().ok_or(AppError::ServiceUnavailable { + code: "SSO_NOT_CONFIGURED", + message: "Sign in with Knotree Accounts to connect Registry.", + })?; + let subject = sso_subject(&state, user.id, &config.issuer) + .await? + .ok_or(AppError::Forbidden { + code: "ACCOUNTS_SIGN_IN_REQUIRED", + message: "This Cloud account must be linked to Knotree Accounts before connecting Registry.", + })?; + let return_to = safe_return_to(input.as_ref().and_then(|v| v.return_to.as_deref())); + let state_token = security::random_token(); + let verifier = security::random_token(); + let challenge = URL_SAFE_NO_PAD.encode(security::token_hash(&verifier)); + let response = registry_consent::client()? + .post(format!("{REGISTRY}/api/v1/cloud-grants/requests")) + .json(&serde_json::json!({ + "client_id": registry_consent::CLIENT, + "redirect_uri": registry_consent::CALLBACK, + "state": state_token, + "namespace": true, + "code_challenge": challenge, + "code_challenge_method": "S256", + "expected_issuer": config.issuer, + "expected_subject": subject, + })) + .send() + .await + .map_err(|_| invalid())?; + let started: Started = registry_consent::bounded_json(response).await?; + if started.authorization_url != format!("{REGISTRY}/cloud/authorize/{}", started.request_id) { + return Err(invalid()); + } + sqlx::query("DELETE FROM registry_account_consent_attempts WHERE expires_at <= now() OR user_id=$1") + .bind(user.id) + .execute(&state.db) + .await?; + sqlx::query( + "INSERT INTO registry_account_consent_attempts + (state_hash, session_hash, user_id, issuer, subject, verifier_ciphertext, return_to, expires_at) + VALUES ($1, $2, $3, $4, $5, $6, $7, now() + interval '10 minutes')", + ) + .bind(security::token_hash(&state_token)) + .bind(registry_consent::session_hash(&state, &headers)?) + .bind(user.id) + .bind(&config.issuer) + .bind(subject) + .bind(security::encrypt_secret( + &verifier, + &state.config.database_credentials_encryption_key, + )?) + .bind(return_to) + .execute(&state.db) + .await?; + Ok(no_store( + Json(serde_json::json!({"authorizationUrl": started.authorization_url})).into_response(), + )) +} + +#[derive(FromRow)] +struct Attempt { + issuer: String, + subject: String, + verifier_ciphertext: String, + return_to: Option, +} + +#[derive(Deserialize)] +struct Grant { + username: String, + credential: String, + credential_id: Uuid, + namespace: Option, + issuer: String, + subject: String, + expires_at: u64, + actions: Vec, +} + +fn validate_grant(grant: &Grant, attempt: &Attempt) -> Result<(), AppError> { + if grant.issuer != attempt.issuer + || grant.subject != attempt.subject + || grant.namespace.as_deref() != Some(grant.username.as_str()) + || grant.actions != ["pull"] + || grant.credential.is_empty() + || grant.credential.len() > 4096 + || grant.credential.chars().any(char::is_control) + { + return Err(invalid()); + } + knotree_registry::validate_registry_username(&grant.username)?; + let now = OffsetDateTime::now_utc().unix_timestamp() as u64; + if grant.expires_at <= now || grant.expires_at > now + 31 * 86400 { + return Err(invalid()); + } + Ok(()) +} + +/// Completes an account-level consent if `state_token` belongs to one. +/// Returns `None` so the project-level callback can handle other attempts. +pub(crate) async fn complete_callback( + state: &AppState, + headers: &HeaderMap, + user_id: Uuid, + state_token: &str, + code: Option, + error: Option<&str>, +) -> Result, AppError> { + let config = state.config.sso.as_ref().ok_or_else(invalid)?; + // Consumed only for the initiating user and live browser session. + let attempt: Option = sqlx::query_as( + "DELETE FROM registry_account_consent_attempts + WHERE state_hash=$1 AND session_hash=$2 AND user_id=$3 AND issuer=$4 AND expires_at>now() + RETURNING issuer, subject, verifier_ciphertext, return_to", + ) + .bind(security::token_hash(state_token)) + .bind(registry_consent::session_hash(state, headers)?) + .bind(user_id) + .bind(&config.issuer) + .fetch_optional(&state.db) + .await?; + let Some(attempt) = attempt else { + return Ok(None); + }; + if sso_subject(state, user_id, &attempt.issuer).await?.as_deref() != Some(attempt.subject.as_str()) { + return Err(invalid()); + } + let status = if error == Some("access_denied") { + "denied" + } else { + if error.is_some() { + return Err(invalid()); + } + let code = code + .filter(|v| v.len() == 64 && v.bytes().all(|b| b.is_ascii_hexdigit())) + .ok_or_else(invalid)?; + let verifier = security::decrypt_secret( + &attempt.verifier_ciphertext, + &state.config.database_credentials_encryption_key, + )?; + let response = registry_consent::client()? + .post(format!("{REGISTRY}/api/v1/cloud-grants/exchange")) + .json(&serde_json::json!({ + "client_id": registry_consent::CLIENT, + "redirect_uri": registry_consent::CALLBACK, + "code": code, + "code_verifier": verifier, + })) + .send() + .await + .map_err(|_| invalid())?; + let grant: Grant = registry_consent::bounded_json(response).await?; + validate_grant(&grant, &attempt)?; + let encrypted = security::encrypt_secret( + &grant.credential, + &state.config.database_credentials_encryption_key, + )?; + let expiry = + OffsetDateTime::from_unix_timestamp(grant.expires_at as i64).map_err(|_| invalid())?; + store_account(state, user_id, &attempt, &grant, &encrypted, expiry).await?; + "connected" + }; + let mut destination = Url::parse(&config.frontend_url).map_err(|_| invalid())?; + let path = attempt.return_to.as_deref().unwrap_or("/integrations"); + let (path, query) = path.split_once('?').unwrap_or((path, "")); + destination.set_path(path); + destination.set_query((!query.is_empty()).then_some(query)); + destination.query_pairs_mut().append_pair("registry", status); + Ok(Some(no_store(Redirect::to(destination.as_str()).into_response()))) +} + +/// Replaces the user's account connection. Reconnecting (for example after +/// the 30-day grant expires) keeps the same account row, so every project +/// connection derived from it picks up the new credential immediately. +async fn store_account( + state: &AppState, + user_id: Uuid, + attempt: &Attempt, + grant: &Grant, + encrypted: &str, + expiry: OffsetDateTime, +) -> Result<(), AppError> { + let mut transaction = state.db.begin().await?; + let existing: Option<(Uuid, String, String)> = sqlx::query_as( + "SELECT id, subject, registry_username FROM knotree_registry_accounts + WHERE user_id=$1 AND revoked_at IS NULL FOR UPDATE", + ) + .bind(user_id) + .fetch_optional(&mut *transaction) + .await?; + match existing { + Some((id, subject, username)) if subject == attempt.subject && username == grant.username => { + sqlx::query( + "UPDATE knotree_registry_accounts + SET credential_ciphertext=$1, delegated_credential_id=$2, + credential_expires_at=$3, updated_at=now() + WHERE id=$4", + ) + .bind(encrypted) + .bind(grant.credential_id) + .bind(expiry) + .bind(id) + .execute(&mut *transaction) + .await?; + } + existing => { + if let Some((id, _, _)) = existing { + revoke_account_rows(&mut transaction, id, "Knotree Registry account was replaced.").await?; + } + sqlx::query( + "INSERT INTO knotree_registry_accounts + (id, user_id, issuer, subject, registry_username, credential_ciphertext, + delegated_credential_id, credential_expires_at) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8)", + ) + .bind(Uuid::new_v4()) + .bind(user_id) + .bind(&attempt.issuer) + .bind(&attempt.subject) + .bind(&grant.username) + .bind(encrypted) + .bind(grant.credential_id) + .bind(expiry) + .execute(&mut *transaction) + .await?; + } + } + transaction.commit().await?; + Ok(()) +} + +async fn revoke_account_rows( + transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>, + account_id: Uuid, + reason: &str, +) -> Result, AppError> { + sqlx::query("UPDATE knotree_registry_accounts SET revoked_at=now(), updated_at=now() WHERE id=$1") + .bind(account_id) + .execute(&mut **transaction) + .await?; + let connection_ids: Vec = sqlx::query_scalar( + "UPDATE knotree_registry_connections SET revoked_at=now(), updated_at=now() + WHERE account_id=$1 AND revoked_at IS NULL RETURNING id", + ) + .bind(account_id) + .fetch_all(&mut **transaction) + .await?; + let service_ids: Vec = sqlx::query_scalar( + "UPDATE project_app_services + SET auto_deploy_enabled = FALSE, registry_connection_id = NULL, + auto_deploy_error = $2, updated_at = now() + WHERE registry_connection_id = ANY($1) RETURNING id", + ) + .bind(&connection_ids) + .bind(reason) + .fetch_all(&mut **transaction) + .await?; + sqlx::query( + "UPDATE knotree_registry_deploy_jobs + SET status = 'failed', locked_until = NULL, last_error = $2, updated_at = now() + WHERE app_service_id = ANY($1) AND status = 'pending'", + ) + .bind(&service_ids) + .bind(reason) + .execute(&mut **transaction) + .await?; + Ok(service_ids) +} + +#[derive(FromRow)] +struct AccountRow { + id: Uuid, + registry_username: String, + credential_ciphertext: String, + credential_expires_at: OffsetDateTime, +} + +async fn active_account(state: &AppState, user_id: Uuid) -> Result, AppError> { + Ok(sqlx::query_as( + "SELECT id, registry_username, credential_ciphertext, credential_expires_at + FROM knotree_registry_accounts WHERE user_id=$1 AND revoked_at IS NULL", + ) + .bind(user_id) + .fetch_optional(&state.db) + .await?) +} + +#[derive(Serialize)] +#[serde(rename_all = "camelCase")] +struct AccountStatus { + connected: bool, + consent_ready: bool, + auto_deploy_ready: bool, + namespace: Option, + expires_at: Option, + expired: bool, +} + +/// GET /integrations/knotree-registry +pub async fn status( + State(state): State, + headers: HeaderMap, +) -> Result { + let user = auth::authenticate(&state, &headers).await?; + let account = active_account(&state, user.id).await?; + let now = OffsetDateTime::now_utc(); + let body = AccountStatus { + connected: account.is_some(), + consent_ready: state.config.sso.is_some(), + auto_deploy_ready: state.config.knotree_registry_webhook_secret.is_some(), + namespace: account.as_ref().map(|a| a.registry_username.clone()), + expires_at: account.as_ref().and_then(|a| { + a.credential_expires_at + .format(&time::format_description::well_known::Rfc3339) + .ok() + }), + expired: account.as_ref().is_some_and(|a| a.credential_expires_at <= now), + }; + Ok(no_store(Json(body).into_response())) +} + +/// DELETE /integrations/knotree-registry +pub async fn disconnect( + State(state): State, + headers: HeaderMap, +) -> Result { + security::require_csrf(&headers, &state.config)?; + let user = auth::authenticate(&state, &headers).await?; + let account = active_account(&state, user.id).await?.ok_or_else(not_connected)?; + let mut transaction = state.db.begin().await?; + let service_ids = revoke_account_rows( + &mut transaction, + account.id, + "Knotree Registry account was disconnected.", + ) + .await?; + transaction.commit().await?; + if state.config.uses_kubernetes_workloads() { + for service_id in service_ids { + if let Err(error) = + crate::cluster_kubernetes::delete_app_image_pull_secret(&state.config, service_id).await + { + tracing::warn!(app_service_id = %service_id, error = %error, + "could not remove Knotree Registry Kubernetes pull Secret"); + } + } + } + Ok(StatusCode::NO_CONTENT) +} + +/// The signed-in user's own credential. Every picker request uses this and +/// nothing else, so a user can never list another user's images. +async fn own_credentials( + state: &AppState, + user_id: Uuid, +) -> Result<(AccountRow, String), AppError> { + let account = active_account(state, user_id).await?.ok_or_else(not_connected)?; + if account.credential_expires_at <= OffsetDateTime::now_utc() { + return Err(AppError::Conflict { + code: "KNOTREE_REGISTRY_ACCOUNT_EXPIRED", + message: "Your Knotree Registry connection expired. Reconnect to continue.", + }); + } + let secret = security::decrypt_secret( + &account.credential_ciphertext, + &state.config.database_credentials_encryption_key, + )?; + Ok((account, secret)) +} + +fn registry_unavailable() -> AppError { + AppError::ServiceUnavailable { + code: "KNOTREE_REGISTRY_UNAVAILABLE", + message: "Knotree Registry could not be reached. Try again.", + } +} + +async fn registry_get(path: &str, username: &str, secret: &str) -> Result { + let mut response = registry_consent::client()? + .get(format!("{REGISTRY}{path}")) + .basic_auth(username, Some(secret)) + .send() + .await + .map_err(|_| registry_unavailable())?; + match response.status() { + status if status.is_success() => {} + StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN => { + return Err(AppError::Conflict { + code: "KNOTREE_REGISTRY_ACCOUNT_REVOKED", + message: "Knotree Registry rejected the connection. Reconnect your account.", + }); + } + StatusCode::NOT_FOUND => { + return Err(AppError::NotFound { + code: "KNOTREE_REGISTRY_REPOSITORY_NOT_FOUND", + message: "The repository could not be found in your Knotree Registry namespace.", + }); + } + _ => return Err(registry_unavailable()), + } + let mut bytes = Vec::new(); + while let Some(chunk) = response.chunk().await.map_err(|_| registry_unavailable())? { + if bytes.len() + chunk.len() > MAX_LIST_BYTES { + return Err(registry_unavailable()); + } + bytes.extend_from_slice(&chunk); + } + serde_json::from_slice(&bytes).map_err(|_| registry_unavailable()) +} + +fn owned_repository(namespace: &str, repository: &str) -> Result { + let repository = knotree_registry::validate_registry_repository(repository)?; + if !repository.starts_with(&format!("{namespace}/")) { + return Err(AppError::Forbidden { + code: "KNOTREE_REGISTRY_NOT_OWNER", + message: "You can only use repositories in your own Knotree Registry namespace.", + }); + } + Ok(repository) +} + +#[derive(Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +struct RepositorySummary { + #[serde(alias = "name")] + name: String, + #[serde(default, alias = "tag_count")] + tag_count: usize, + #[serde(default, alias = "latest_tag")] + latest_tag: Option, + #[serde(default, alias = "latest_digest")] + latest_digest: Option, + #[serde(default)] + size: u64, + #[serde(default, alias = "updated_at")] + updated_at: Option, +} + +/// GET /integrations/knotree-registry/repositories +pub async fn repositories( + State(state): State, + headers: HeaderMap, +) -> Result { + let user = auth::authenticate(&state, &headers).await?; + let (account, secret) = own_credentials(&state, user.id).await?; + let body = registry_get( + "/api/v1/integrations/cloud/repositories", + &account.registry_username, + &secret, + ) + .await?; + let repositories: Vec = + serde_json::from_value(body["repositories"].clone()).map_err(|_| registry_unavailable())?; + // Defense in depth: never surface anything outside the user's namespace. + let prefix = format!("{}/", account.registry_username); + let repositories: Vec<_> = repositories + .into_iter() + .filter(|repository| repository.name.starts_with(&prefix)) + .collect(); + Ok(no_store( + Json(serde_json::json!({ + "namespace": account.registry_username, + "registryHost": knotree_registry::REGISTRY_HOST, + "repositories": repositories, + })) + .into_response(), + )) +} + +#[derive(Deserialize, Serialize)] +#[serde(rename_all = "camelCase")] +struct TagSummary { + tag: String, + digest: String, + #[serde(default)] + size: u64, + #[serde(default, alias = "created_at")] + created_at: Option, +} + +/// GET /integrations/knotree-registry/repositories/{*repository} +pub async fn repository_tags( + State(state): State, + headers: HeaderMap, + Path(repository): Path, +) -> Result { + let user = auth::authenticate(&state, &headers).await?; + let (account, secret) = own_credentials(&state, user.id).await?; + let repository = owned_repository(&account.registry_username, &repository)?; + let body = registry_get( + &format!("/api/v1/integrations/cloud/repositories/{repository}"), + &account.registry_username, + &secret, + ) + .await?; + let tags: Vec = + serde_json::from_value(body["tags"].clone()).map_err(|_| registry_unavailable())?; + Ok(no_store( + Json(serde_json::json!({"repository": repository, "tags": tags})).into_response(), + )) +} + +#[derive(Deserialize)] +pub struct ImportInput { + repository: String, +} + +/// POST /workspaces/{w}/projects/{p}/registry-connections/from-account +/// +/// Returns a project connection for one repository of the user's own +/// namespace, creating it if needed. The existing app-service flow then takes +/// its `id` as `registryConnectionId`. +pub async fn import_into_project( + State(state): State, + headers: HeaderMap, + Path((workspace_id, project_slug)): Path<(String, String)>, + Json(input): Json, +) -> Result { + security::require_csrf(&headers, &state.config)?; + let user = auth::authenticate(&state, &headers).await?; + let project_id = + projects::accessible_project_id(&state, user.id, &workspace_id, &project_slug).await?; + let (account, secret) = own_credentials(&state, user.id).await?; + let repository = owned_repository(&account.registry_username, &input.repository)?; + if !knotree_registry::verify_pull_access(&account.registry_username, &secret, &repository).await { + return Err(AppError::Conflict { + code: "KNOTREE_REGISTRY_AUTHENTICATION_FAILED", + message: "Knotree Registry could not verify pull access to this repository.", + }); + } + let existing: Option = sqlx::query_scalar( + "SELECT id FROM knotree_registry_connections + WHERE project_id=$1 AND account_id=$2 AND repository=$3 AND revoked_at IS NULL", + ) + .bind(project_id) + .bind(account.id) + .bind(&repository) + .fetch_optional(&state.db) + .await?; + let id = match existing { + Some(id) => id, + None => { + // Membership is re-checked in the INSERT after network I/O. + let id = Uuid::new_v4(); + let inserted = sqlx::query( + "INSERT INTO knotree_registry_connections + (id, project_id, user_id, registry_username, repository, + credential_ciphertext, account_id) + SELECT $1, $2, $3, $4, $5, $6, $7 + WHERE EXISTS (SELECT 1 FROM projects p + JOIN workspace_memberships wm ON wm.workspace_id = p.workspace_id + WHERE p.id = $2 AND wm.user_id = $3) + ON CONFLICT DO NOTHING", + ) + .bind(id) + .bind(project_id) + .bind(user.id) + .bind(&account.registry_username) + .bind(&repository) + .bind(&account.credential_ciphertext) + .bind(account.id) + .execute(&state.db) + .await?; + if inserted.rows_affected() == 1 { + id + } else { + sqlx::query_scalar( + "SELECT id FROM knotree_registry_connections + WHERE project_id=$1 AND account_id=$2 AND repository=$3 AND revoked_at IS NULL", + ) + .bind(project_id) + .bind(account.id) + .bind(&repository) + .fetch_optional(&state.db) + .await? + .ok_or_else(invalid)? + } + } + }; + Ok(no_store( + Json(serde_json::json!({ + "id": id, + "registryHost": knotree_registry::REGISTRY_HOST, + "username": account.registry_username, + "repository": repository, + })) + .into_response(), + )) +} + +#[cfg(test)] +mod tests { + use super::*; + + fn attempt() -> Attempt { + Attempt { + issuer: "https://accounts.knotree.com".into(), + subject: "alice".into(), + verifier_ciphertext: "unused".into(), + return_to: None, + } + } + fn grant() -> Grant { + Grant { + username: "kt-alice".into(), + credential: "secret-test-only".into(), + credential_id: Uuid::new_v4(), + namespace: Some("kt-alice".into()), + issuer: "https://accounts.knotree.com".into(), + subject: "alice".into(), + expires_at: (OffsetDateTime::now_utc().unix_timestamp() + 86400) as u64, + actions: vec!["pull".into()], + } + } + + #[test] + fn namespace_grant_requires_exact_identity_own_namespace_and_pull_only() { + let attempt = attempt(); + assert!(validate_grant(&grant(), &attempt).is_ok()); + let mut wrong = grant(); + wrong.subject = "bob".into(); + assert!(validate_grant(&wrong, &attempt).is_err()); + wrong = grant(); + wrong.namespace = Some("kt-bob".into()); + assert!(validate_grant(&wrong, &attempt).is_err()); + wrong = grant(); + wrong.namespace = None; + assert!(validate_grant(&wrong, &attempt).is_err()); + wrong = grant(); + wrong.actions = vec!["pull".into(), "push".into()]; + assert!(validate_grant(&wrong, &attempt).is_err()); + wrong = grant(); + wrong.expires_at = 0; + assert!(validate_grant(&wrong, &attempt).is_err()); + } + + #[test] + fn picker_and_import_reject_other_namespaces() { + assert!(owned_repository("kt-alice", "kt-alice/api").is_ok()); + assert!(owned_repository("kt-alice", "kt-bob/api").is_err()); + assert!(owned_repository("kt-alice", "kt-alice-evil/api").is_err()); + assert!(owned_repository("kt-alice", "kt-alice/../kt-bob/api").is_err()); + } + + #[test] + fn return_to_stays_on_cloud() { + assert_eq!(safe_return_to(Some("/workspace/w/project/p")).as_deref(), Some("/workspace/w/project/p")); + assert!(safe_return_to(Some("//evil.example")).is_none()); + assert!(safe_return_to(Some("https://evil.example")).is_none()); + assert!(safe_return_to(Some("/\\evil.example")).is_none()); + assert!(safe_return_to(None).is_none()); + } +} diff --git a/apps/api/src/registry_consent.rs b/apps/api/src/registry_consent.rs index b660aaf..4bb82bf 100644 --- a/apps/api/src/registry_consent.rs +++ b/apps/api/src/registry_consent.rs @@ -16,8 +16,8 @@ use uuid::Uuid; use crate::{auth, error::AppError, knotree_registry, projects, security, state::AppState}; const REGISTRY: &str = "https://registry.knotree.com"; -const CLIENT: &str = "knotree-cloud"; -const CALLBACK: &str = "https://cloud.knotree.com/api/v1/auth/knotree-registry/callback"; +pub(crate) const CLIENT: &str = "knotree-cloud"; +pub(crate) const CALLBACK: &str = "https://cloud.knotree.com/api/v1/auth/knotree-registry/callback"; fn invalid() -> AppError { AppError::BadRequest { @@ -25,19 +25,19 @@ fn invalid() -> AppError { message: "Registry authorization expired or could not be verified. Start again.", } } -fn session_hash(state: &AppState, headers: &HeaderMap) -> Result, AppError> { +pub(crate) fn session_hash(state: &AppState, headers: &HeaderMap) -> Result, AppError> { let token = security::get_cookie(headers, state.config.session_cookie_name()).ok_or_else(invalid)?; Ok(security::token_hash(&token)) } -fn client() -> Result { +pub(crate) fn client() -> Result { reqwest::Client::builder() .timeout(Duration::from_secs(10)) .redirect(reqwest::redirect::Policy::none()) .build() .map_err(|_| invalid()) } -async fn bounded_json( +pub(crate) async fn bounded_json( mut response: reqwest::Response, ) -> Result { if !response.status().is_success() { @@ -162,6 +162,19 @@ pub async fn callback( let user = auth::authenticate(&state, &headers).await?; let config = state.config.sso.as_ref().ok_or_else(invalid)?; let token = query.state.filter(|v| v.len() == 43).ok_or_else(invalid)?; + // Account-level ("connect once") consents share this registered callback. + if let Some(response) = crate::registry_accounts::complete_callback( + &state, + &headers, + user.id, + &token, + query.code.clone(), + query.error.as_deref(), + ) + .await? + { + return Ok(response); + } // Consume only for the initiating user and live browser session. Replays, // another signed-in account, and configuration changes fail closed. let attempt: Attempt = sqlx::query_as("DELETE FROM registry_consent_attempts WHERE state_hash=$1 AND session_hash=$2 AND user_id=$3 AND issuer=$4 AND expires_at>now() RETURNING project_id,workspace_id,project_slug,repository,issuer,subject,verifier_ciphertext") diff --git a/docs/intakes/IN-047.md b/docs/intakes/IN-047.md new file mode 100644 index 0000000..29627a4 --- /dev/null +++ b/docs/intakes/IN-047.md @@ -0,0 +1,20 @@ +--- +created_at: "2026-10-03T04:05:25.072688100+00:00" +docs: null +flags: null +id: IN-047 +input_type: change_request +lane: high_risk +links: [] +notes: null +status: pending +stories: [] +story: null +summary: "Account-level Knotree Registry connection: one consent, image picker, owner-verified auto-deploy (needs registry decision 0004)" +type: intake +updated_at: "2026-10-03T04:05:25.072691600+00:00" +--- + +# Intake IN-047 + +Account-level Knotree Registry connection: one consent, image picker, owner-verified auto-deploy (needs registry decision 0004) diff --git a/docs/stories/US-048.md b/docs/stories/US-048.md new file mode 100644 index 0000000..4c79a4e --- /dev/null +++ b/docs/stories/US-048.md @@ -0,0 +1,19 @@ +--- +contract: User connects registry once via SSO-bound consent; picker lists only own namespace; project connections derive from the account; webhook deploys only when owner_subject matches +created_at: "2026-10-03T04:05:25.394976400+00:00" +e2e: 0 +evidence: null +id: US-048 +integration: 0 +lane: high_risk +notes: null +platform: 0 +status: in_progress +title: Connect Knotree Registry account once and import images +type: story +unit: 0 +updated_at: "2026-10-03T04:05:30.864766200+00:00" +verify: cargo test -p knotree-api registry; pnpm --filter web test +--- + +# Connect Knotree Registry account once and import images