aegis-api 0.3.0

Independent Aegis fleet control plane HTTP service.
Documentation
use std::sync::Arc;

use aegis_dto::identity::USER_CLIENT_ID;
use async_trait::async_trait;
use phylax_core::{
    AccessTokenGrant, JsonRefreshTokenExchange, JsonRefreshTokenGrant,
    JsonRefreshTokenGrantRequest, JwtIssuer, RefreshTokenValidation, ScopeSet,
};
use phylax_gcp::{FirestoreAuthStore, RefreshTokenGrantRequest, identity::IdentityStore};

#[derive(Clone)]
pub struct UserTokenGrant {
    pub auth: FirestoreAuthStore,
    pub identities: IdentityStore,
    pub issuer: Arc<JwtIssuer>,
    pub audience: String,
}

#[async_trait]
impl JsonRefreshTokenGrant for UserTokenGrant {
    type Extra = super::oauth::OAuthTokenExtra;

    fn issuer(&self) -> Arc<JwtIssuer> {
        self.issuer.clone()
    }

    async fn exchange_refresh_token(
        &self,
        request: JsonRefreshTokenGrantRequest<'_>,
    ) -> anyhow::Result<RefreshTokenValidation<JsonRefreshTokenExchange<Self::Extra>>> {
        let request = RefreshTokenGrantRequest {
            refresh_token: request.refresh_token,
            client_id: USER_CLIENT_ID,
            now_unix: request.now_unix,
        };
        let RefreshTokenValidation::Valid(inspection) =
            self.auth.inspect_refresh_token(request).await?
        else {
            return Ok(RefreshTokenValidation::Invalid);
        };
        let Some(user_id) = inspection.subject.strip_kind("user") else {
            return Ok(RefreshTokenValidation::Invalid);
        };
        let Some(user) = self
            .identities
            .user(user_id)
            .await?
            .filter(|user| !user.disabled)
        else {
            return Ok(RefreshTokenValidation::Invalid);
        };
        let RefreshTokenValidation::Valid(exchange) =
            self.auth.exchange_refresh_token(request).await?
        else {
            return Ok(RefreshTokenValidation::Invalid);
        };
        anyhow::ensure!(
            exchange.subject == inspection.subject,
            "session identity changed during exchange"
        );
        Ok(RefreshTokenValidation::Valid(JsonRefreshTokenExchange {
            access_grant: AccessTokenGrant {
                subject: exchange.subject,
                client_id: USER_CLIENT_ID.into(),
                audience: vec![self.audience.clone()],
                scope: ScopeSet::new(["aegis:read", "aegis:user"])?,
                sid: Some(exchange.sid),
                refresh_expires_unix: exchange.expires_unix,
                extra: super::oauth::OAuthTokenExtra { principal: user.id },
            },
            refresh_token: exchange.refresh_token,
        }))
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use axum::{
        Router,
        body::{Body, to_bytes},
        http::{Request, StatusCode},
    };
    use phylax_core::{AccessClaims, Subject};
    use serde_json::{Value, json};
    use tower::ServiceExt;

    async fn exchange(app: &Router, token: &str) -> (StatusCode, Value) {
        let response = app
            .clone()
            .oneshot(
                Request::post("/v2/auth/token")
                    .header("content-type", "application/json")
                    .body(Body::from(
                        json!({"grant_type":"refresh_token","refresh_token":token}).to_string(),
                    ))
                    .unwrap(),
            )
            .await
            .unwrap();
        (
            response.status(),
            serde_json::from_slice(&to_bytes(response.into_body(), 65536).await.unwrap()).unwrap(),
        )
    }

    #[tokio::test]
    #[ignore = "requires a loopback FIRESTORE_EMULATOR_HOST"]
    async fn emulator_credential_sessions_rotate_revoke_and_work_without_oauth()
    -> anyhow::Result<()> {
        anyhow::ensure!(
            std::env::var("FIRESTORE_EMULATOR_HOST")?.starts_with("127.0.0.1:"),
            "use a local emulator"
        );
        tokio::time::timeout(std::time::Duration::from_secs(60), exercise()).await??;
        Ok(())
    }

    async fn exercise() -> anyhow::Result<()> {
        let project = format!(
            "aegis-keys-{}",
            std::time::SystemTime::now()
                .duration_since(std::time::UNIX_EPOCH)?
                .as_nanos()
        );
        let db = arche_firestore::DatabaseOptions {
            project_id: Some(project),
            database_id: "aegis".into(),
            credentials: arche_firestore::Credentials::AccessToken("owner".into()),
            ..Default::default()
        }
        .validate()?
        .connect()
        .await?;
        let identities = crate::identity::store(&db)?;
        identities
            .bootstrap(
                phylax_gcp::identity::IdentityBootstrapOptions {
                    issuer_url: "https://fleet.example/v2".into(),
                    audience: "test-users".into(),
                    ..Default::default()
                }
                .validate()?,
            )
            .await?;
        arche_firestore::create_typed_at(
            db.inner(),
            identities.parent(),
            "settings",
            "authentication",
            &aegis_dto::identity::AuthenticationConfiguration::default(),
        )
        .await?;
        identities
            .create_user(
                &phylax_gcp::identity::UserRecord {
                    id: "owner".into(),
                    email: "owner@example.com".into(),
                    disabled: false,
                    session_version: 1,
                },
                None,
            )
            .await?;
        let config = identities.load_config().await?;
        let issuer = JwtIssuer::from_config(&config.api_token.jwt_config())?;
        let auth = identities.auth_store(&config, 600)?;
        let now = time::OffsetDateTime::now_utc().unix_timestamp();
        let issue = || phylax_gcp::RefreshTokenIssueRequest {
            subject: Subject::new("user:owner").unwrap(),
            client_id: USER_CLIENT_ID,
            now_unix: now,
        };
        let first = auth.issue_refresh_token(issue()).await?;
        let second = auth.issue_refresh_token(issue()).await?;
        let app = crate::application(db.clone(), None).await?;
        let info = app
            .clone()
            .oneshot(Request::get("/v2/info").body(Body::empty())?)
            .await?;
        let info: Value = serde_json::from_slice(&to_bytes(info.into_body(), 65536).await?)?;
        assert_eq!(
            info["authentication"],
            json!({"credentials":true,"oauth":false})
        );
        assert_eq!(
            app.clone()
                .oneshot(Request::get("/v2/oauth/authorize").body(Body::empty())?)
                .await?
                .status(),
            StatusCode::NOT_FOUND
        );
        assert_eq!(exchange(&app, "invalid").await.0, StatusCode::UNAUTHORIZED);
        let (status, token) = exchange(&app, first.refresh_token()).await;
        assert_eq!(status, StatusCode::OK);
        assert_eq!(token["principal"], "owner");
        let claims: AccessClaims =
            issuer.decode_access(token["access_token"].as_str().unwrap(), "test-users")?;
        assert_eq!(claims.sid.as_deref(), Some(first.session_id()));
        assert!(identities.session_active(&claims, now).await?);
        assert_eq!(
            exchange(&app, first.refresh_token()).await.0,
            StatusCode::UNAUTHORIZED
        );
        let (status, rotated) = exchange(&app, token["refresh_token"].as_str().unwrap()).await;
        assert_eq!(status, StatusCode::OK);
        assert!(
            auth.revoke_refresh_session(
                first.session_id(),
                &Subject::new("user:owner")?,
                USER_CLIENT_ID,
                now
            )
            .await?
        );
        assert!(!identities.session_active(&claims, now).await?);
        assert_eq!(
            exchange(&app, rotated["refresh_token"].as_str().unwrap())
                .await
                .0,
            StatusCode::UNAUTHORIZED
        );
        identities.set_disabled("owner", true).await?;
        assert_eq!(
            exchange(&app, second.refresh_token()).await.0,
            StatusCode::UNAUTHORIZED
        );
        identities.set_disabled("owner", false).await?;
        let (status, second_token) = exchange(&app, second.refresh_token()).await;
        assert_eq!(status, StatusCode::OK);
        let revoked = app
            .clone()
            .oneshot(
                Request::post("/v2/auth/revoke")
                    .header("content-type", "application/x-www-form-urlencoded")
                    .body(Body::from(format!(
                        "token={}",
                        second_token["refresh_token"].as_str().unwrap()
                    )))?,
            )
            .await?;
        assert_eq!(revoked.status(), StatusCode::OK);
        assert_eq!(
            exchange(&app, second_token["refresh_token"].as_str().unwrap())
                .await
                .0,
            StatusCode::UNAUTHORIZED
        );
        assert!(
            crate::application(db, Some("unexpected-secret".into()))
                .await
                .is_err()
        );
        Ok(())
    }
}