feat: publish identity events #32
File diff suppressed because it is too large
Load Diff
@@ -1,81 +1,84 @@
|
||||
[package]
|
||||
name = "syncode-identity"
|
||||
description = "SynCode identity service"
|
||||
version.workspace = true
|
||||
edition.workspace = true
|
||||
rust-version.workspace = true
|
||||
license.workspace = true
|
||||
repository.workspace = true
|
||||
publish = false
|
||||
|
||||
[workspace]
|
||||
members = [
|
||||
".",
|
||||
"crates/model",
|
||||
"crates/application",
|
||||
"crates/storage",
|
||||
"crates/api-grpc",
|
||||
"crates/api-graphql",
|
||||
"crates/api-oauth",
|
||||
"crates/gitea-client",
|
||||
"crates/mailer",
|
||||
"crates/token-crypto",
|
||||
]
|
||||
resolver = "3"
|
||||
|
||||
[workspace.package]
|
||||
version = "0.5.0"
|
||||
edition = "2024"
|
||||
rust-version = "1.95"
|
||||
license = "MIT"
|
||||
repository = "https://syncode.sh/syncode/identity"
|
||||
|
||||
[workspace.lints.rust]
|
||||
unsafe_code = "forbid"
|
||||
|
||||
[workspace.lints.clippy]
|
||||
expect_used = "deny"
|
||||
panic = "deny"
|
||||
unwrap_used = "deny"
|
||||
|
||||
[dependencies]
|
||||
axum = { version = "0.8.9", default-features = false, features = ["http1", "tokio"] }
|
||||
chrono = { version = "0.4.42", default-features = false, features = ["clock", "serde"] }
|
||||
clap = { version = "4.6.4", features = ["derive", "env"] }
|
||||
cynic = "3.14.0"
|
||||
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] }
|
||||
serde = { version = "1.0", features = ["derive"] }
|
||||
serde_json = "1.0"
|
||||
sha2 = "0.11.0"
|
||||
syncode-identity-api-graphql = { path = "crates/api-graphql" }
|
||||
syncode-identity-api-grpc = { path = "crates/api-grpc" }
|
||||
syncode-identity-api-oauth = { path = "crates/api-oauth" }
|
||||
syncode-identity-application = { path = "crates/application" }
|
||||
syncode-identity-gitea-client = { path = "crates/gitea-client" }
|
||||
syncode-identity-mailer = { path = "crates/mailer" }
|
||||
syncode-identity-model = { path = "crates/model" }
|
||||
syncode-identity-storage = { path = "crates/storage" }
|
||||
syncode-identity-token-crypto = { path = "crates/token-crypto" }
|
||||
thiserror = "2.0.19"
|
||||
tokio = { version = "1.49.0", features = ["macros", "net", "rt-multi-thread", "signal", "time"] }
|
||||
tonic = "0.14.6"
|
||||
tower-http = { version = "0.6.11", default-features = false, features = ["cors"] }
|
||||
tracing = "0.1.44"
|
||||
tracing-subscriber = { version = "0.3.22", default-features = false, features = ["ansi", "fmt"] }
|
||||
uuid = { version = "1.24.0", features = ["v4", "serde"] }
|
||||
|
||||
[dev-dependencies]
|
||||
async-trait = "0.1.89"
|
||||
chrono = { version = "0.4.42", default-features = false, features = ["clock"] }
|
||||
sqlx = { version = "0.9.0", default-features = false, features = ["postgres", "runtime-tokio", "tls-rustls-ring-webpki", "uuid"] }
|
||||
ssh-key = "0.7.0-rc.11"
|
||||
syncode-identity-api-grpc = { path = "crates/api-grpc" }
|
||||
syncode-identity-model = { path = "crates/model" }
|
||||
syncode-identity-storage = { path = "crates/storage" }
|
||||
tokio = { version = "1.49.0", features = ["macros", "rt-multi-thread"] }
|
||||
tonic = "0.14.6"
|
||||
tower = { version = "0.5.3", default-features = false, features = ["util"] }
|
||||
uuid = { version = "1.24.0", features = ["v4"] }
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
[package]
|
||||
name = "syncode-identity"
|
||||
description = "SynCode identity service"
|
||||
version.workspace = true
|
||||
edition.workspace = true
|
||||
rust-version.workspace = true
|
||||
license.workspace = true
|
||||
repository.workspace = true
|
||||
publish = false
|
||||
|
||||
[workspace]
|
||||
members = [
|
||||
".",
|
||||
"crates/model",
|
||||
"crates/application",
|
||||
"crates/storage",
|
||||
"crates/api-grpc",
|
||||
"crates/api-graphql",
|
||||
"crates/api-oauth",
|
||||
"crates/gitea-client",
|
||||
"crates/mailer",
|
||||
"crates/token-crypto",
|
||||
]
|
||||
resolver = "3"
|
||||
|
||||
[workspace.package]
|
||||
version = "0.5.0"
|
||||
edition = "2024"
|
||||
rust-version = "1.95"
|
||||
license = "MIT"
|
||||
repository = "https://syncode.sh/syncode/identity"
|
||||
|
||||
[workspace.lints.rust]
|
||||
unsafe_code = "forbid"
|
||||
|
||||
[workspace.lints.clippy]
|
||||
expect_used = "deny"
|
||||
panic = "deny"
|
||||
unwrap_used = "deny"
|
||||
|
||||
[dependencies]
|
||||
axum = { version = "0.8.9", default-features = false, features = ["http1", "tokio"] }
|
||||
chrono = { version = "0.4.42", default-features = false, features = ["clock", "serde"] }
|
||||
clap = { version = "4.6.4", features = ["derive", "env"] }
|
||||
cynic = "3.14.0"
|
||||
hex = "0.4.3"
|
||||
hmac = "0.13.0"
|
||||
reqwest = { version = "0.12", default-features = false, features = ["json", "rustls-tls"] }
|
||||
serde = { version = "1.0", features = ["derive"] }
|
||||
serde_json = "1.0"
|
||||
sha2 = "0.11.0"
|
||||
syncode-identity-api-graphql = { path = "crates/api-graphql" }
|
||||
syncode-identity-api-grpc = { path = "crates/api-grpc" }
|
||||
syncode-identity-api-oauth = { path = "crates/api-oauth" }
|
||||
syncode-identity-application = { path = "crates/application" }
|
||||
syncode-identity-gitea-client = { path = "crates/gitea-client" }
|
||||
syncode-identity-mailer = { path = "crates/mailer" }
|
||||
syncode-identity-model = { path = "crates/model" }
|
||||
syncode-identity-storage = { path = "crates/storage" }
|
||||
syncode-identity-token-crypto = { path = "crates/token-crypto" }
|
||||
thiserror = "2.0.19"
|
||||
tokio = { version = "1.49.0", features = ["macros", "net", "rt-multi-thread", "signal", "time"] }
|
||||
tonic = "0.14.6"
|
||||
tower-http = { version = "0.6.11", default-features = false, features = ["cors"] }
|
||||
tracing = "0.1.44"
|
||||
tracing-subscriber = { version = "0.3.22", default-features = false, features = ["ansi", "fmt"] }
|
||||
uuid = { version = "1.24.0", features = ["v4", "serde"] }
|
||||
url = "2.5.8"
|
||||
|
||||
[dev-dependencies]
|
||||
async-trait = "0.1.89"
|
||||
chrono = { version = "0.4.42", default-features = false, features = ["clock"] }
|
||||
sqlx = { version = "0.9.0", default-features = false, features = ["postgres", "runtime-tokio", "tls-rustls-ring-webpki", "uuid"] }
|
||||
ssh-key = "0.7.0-rc.11"
|
||||
syncode-identity-api-grpc = { path = "crates/api-grpc" }
|
||||
syncode-identity-model = { path = "crates/model" }
|
||||
syncode-identity-storage = { path = "crates/storage" }
|
||||
tokio = { version = "1.49.0", features = ["macros", "rt-multi-thread"] }
|
||||
tonic = "0.14.6"
|
||||
tower = { version = "0.5.3", default-features = false, features = ["util"] }
|
||||
uuid = { version = "1.24.0", features = ["v4"] }
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
+31
@@ -1,270 +1,301 @@
|
||||
mod sync_codeberg;
|
||||
mod sync_github;
|
||||
mod sync_gitlab;
|
||||
mod sync_native;
|
||||
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::Arc;
|
||||
|
||||
use axum::Router;
|
||||
use axum::http::{Method, header};
|
||||
use clap::Parser;
|
||||
use syncode_identity_storage::Postgres;
|
||||
use tokio::net::TcpListener;
|
||||
use tonic::transport::Server;
|
||||
use tower_http::cors::{AllowOrigin, CorsLayer};
|
||||
|
||||
use sync_codeberg::spawn_codeberg_contribution_sync;
|
||||
use sync_github::spawn_github_contribution_sync;
|
||||
use sync_gitlab::spawn_gitlab_contribution_sync;
|
||||
use sync_native::spawn_native_contribution_sync;
|
||||
|
||||
#[derive(Debug, Parser)]
|
||||
#[command(name = "syncode-identity", version, about)]
|
||||
struct Arguments {
|
||||
/// Where `syn`, `LocalAgent` daemons and the Go bridge open gRPC.
|
||||
#[arg(long, default_value = "127.0.0.1:8190")]
|
||||
listen_grpc: SocketAddr,
|
||||
|
||||
/// Where the Leptos front reaches GraphQL and OAuth start/callback.
|
||||
#[arg(long, default_value = "127.0.0.1:8191")]
|
||||
listen_http: SocketAddr,
|
||||
|
||||
/// Where principals and grants are kept.
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_DATABASE_URL", hide_env_values = true)]
|
||||
database: String,
|
||||
|
||||
#[arg(long, default_value_t = 8)]
|
||||
database_connections: u32,
|
||||
|
||||
/// Shared with the Go bridge; sent as `authorization: Bearer <secret>`.
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SHARED_SECRET", hide_env_values = true)]
|
||||
shared_secret: String,
|
||||
|
||||
/// This service's own externally-reachable origin — see
|
||||
/// `syncode_identity_api_oauth::Config`.
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_PUBLIC_URL")]
|
||||
public_url: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_COOKIE_DOMAIN")]
|
||||
cookie_domain: String,
|
||||
|
||||
#[arg(
|
||||
long,
|
||||
env = "SYNCODE_IDENTITY_ALLOWED_REDIRECT_ORIGINS",
|
||||
value_delimiter = ','
|
||||
)]
|
||||
allowed_redirect_origins: Vec<String>,
|
||||
|
||||
/// Where the email-verification link sends people back to.
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_FRONT_URL")]
|
||||
front_url: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_HOST")]
|
||||
smtp_host: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_PORT", default_value_t = 587)]
|
||||
smtp_port: u16,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_USERNAME")]
|
||||
smtp_username: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_PASSWORD", hide_env_values = true)]
|
||||
smtp_password: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_FROM")]
|
||||
smtp_from: String,
|
||||
|
||||
#[arg(
|
||||
long,
|
||||
env = "SYNCODE_IDENTITY_SMTP_FROM_NAME",
|
||||
default_value = "SynCode"
|
||||
)]
|
||||
smtp_from_name: String,
|
||||
|
||||
/// Gitea's origin inside the compose network — where the internal API
|
||||
/// that executes lock/unlock/hard-delete lives.
|
||||
#[arg(long, env = "SYNCODE_GITEA_INTERNAL_URL")]
|
||||
gitea_internal_url: String,
|
||||
|
||||
/// Shared with `syncode-control`'s `GITEA_INTERNAL_TOKEN`.
|
||||
#[arg(long, env = "SYNCODE_GITEA_INTERNAL_TOKEN", hide_env_values = true)]
|
||||
gitea_internal_token: String,
|
||||
|
||||
/// 32 raw bytes, base64-encoded — encrypts linked GitHub/GitLab OAuth
|
||||
/// tokens at rest (0.4 stage 5, part 2's heatmap sync needs them
|
||||
/// later). Fails fast on startup if missing or malformed, rather than
|
||||
/// silently skipping every token save.
|
||||
#[arg(
|
||||
long,
|
||||
env = "SYNCODE_IDENTITY_TOKEN_ENCRYPTION_KEY",
|
||||
hide_env_values = true
|
||||
)]
|
||||
token_encryption_key: String,
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
tracing_subscriber::fmt().try_init()?;
|
||||
let arguments = Arguments::parse();
|
||||
|
||||
let token_key =
|
||||
syncode_identity_token_crypto::Key::from_base64(&arguments.token_encryption_key)?;
|
||||
|
||||
let store = Postgres::connect(&arguments.database, arguments.database_connections).await?;
|
||||
let mailer = syncode_identity_mailer::Mailer::new(syncode_identity_mailer::Config {
|
||||
host: arguments.smtp_host,
|
||||
port: arguments.smtp_port,
|
||||
username: arguments.smtp_username,
|
||||
password: arguments.smtp_password,
|
||||
from: arguments.smtp_from,
|
||||
from_name: arguments.smtp_from_name,
|
||||
})?;
|
||||
|
||||
let (_health_reporter, health_service) = syncode_identity_api_grpc::health_service();
|
||||
let identity_bridge: Arc<dyn syncode_identity_application::IdentityBridgeUseCases> = Arc::new(
|
||||
syncode_identity_application::IdentityBridgeApplication::new(Arc::new(store.clone())),
|
||||
);
|
||||
let identity_service = syncode_identity_api_grpc::authenticated_identity_service(
|
||||
syncode_identity_api_grpc::IdentityServer::new(identity_bridge),
|
||||
arguments.shared_secret,
|
||||
);
|
||||
let gitea =
|
||||
syncode_identity_gitea_client::GiteaClient::new(syncode_identity_gitea_client::Config {
|
||||
internal_url: arguments.gitea_internal_url,
|
||||
internal_token: arguments.gitea_internal_token,
|
||||
});
|
||||
let application = Arc::new(syncode_identity_application::SettingsApplication::new(
|
||||
Arc::new(store.clone()),
|
||||
Arc::new(gitea.clone()),
|
||||
Arc::new(gitea.clone()),
|
||||
Arc::new(mailer.clone()),
|
||||
arguments.public_url.clone(),
|
||||
));
|
||||
let settings: Arc<dyn syncode_identity_application::SettingsUseCases> = application.clone();
|
||||
let sessions: Arc<dyn syncode_identity_application::SessionUseCases> = application.clone();
|
||||
let oauth: Arc<dyn syncode_identity_application::OAuthUseCases> =
|
||||
Arc::new(syncode_identity_application::OAuthApplication::new(
|
||||
Arc::new(store.clone()),
|
||||
Arc::new(gitea.clone()),
|
||||
Arc::new(gitea.clone()),
|
||||
Arc::new(mailer.clone()),
|
||||
arguments.public_url.clone(),
|
||||
));
|
||||
let local_agents: Arc<dyn syncode_identity_application::LocalAgentUseCases> = Arc::new(
|
||||
syncode_identity_application::LocalAgentApplication::new(Arc::new(store.clone())),
|
||||
);
|
||||
let local_agent_service = syncode_identity_api_grpc::GeneratedLocalAgentServer::new(
|
||||
syncode_identity_api_grpc::LocalAgentServer::new(local_agents, settings.clone()),
|
||||
);
|
||||
let settings_service = syncode_identity_api_grpc::GeneratedSettingsServer::new(
|
||||
syncode_identity_api_grpc::SettingsServer::new(settings.clone()),
|
||||
);
|
||||
let forge_operations = syncode_identity_api_grpc::GeneratedForgeOperationsServer::new(
|
||||
syncode_identity_api_grpc::ForgeOperationsServer::new(Arc::new(
|
||||
syncode_identity_application::ForgeOperationsApplication::new(Arc::new(gitea.clone())),
|
||||
)),
|
||||
);
|
||||
let allowed_origins = arguments.allowed_redirect_origins.clone();
|
||||
let cors = CorsLayer::new()
|
||||
.allow_credentials(true)
|
||||
.allow_methods([Method::GET, Method::POST, Method::OPTIONS])
|
||||
.allow_headers([header::CONTENT_TYPE])
|
||||
.allow_origin(AllowOrigin::predicate(move |origin, _| {
|
||||
let Ok(origin) = origin.to_str() else {
|
||||
return false;
|
||||
};
|
||||
allowed_origins.iter().any(|allowed| allowed == origin)
|
||||
}));
|
||||
|
||||
let deletion_sweep = spawn_deletion_sweep(store.clone(), gitea.clone());
|
||||
let native_sync = spawn_native_contribution_sync(store.clone(), gitea.clone());
|
||||
let github_sync = spawn_github_contribution_sync(store.clone(), token_key.clone());
|
||||
let gitlab_sync = spawn_gitlab_contribution_sync(store.clone(), token_key.clone());
|
||||
let codeberg_sync = spawn_codeberg_contribution_sync(store.clone());
|
||||
|
||||
let http = Router::new()
|
||||
.merge(syncode_identity_api_graphql::router(settings, sessions))
|
||||
.merge(syncode_identity_api_oauth::router(
|
||||
oauth,
|
||||
syncode_identity_api_oauth::Config {
|
||||
public_url: arguments.public_url,
|
||||
cookie_domain: arguments.cookie_domain,
|
||||
allowed_redirect_origins: arguments.allowed_redirect_origins,
|
||||
front_url: arguments.front_url,
|
||||
},
|
||||
token_key,
|
||||
))
|
||||
.layer(cors);
|
||||
let http_listener = TcpListener::bind(arguments.listen_http).await?;
|
||||
|
||||
tokio::select! {
|
||||
served = axum::serve(http_listener, http).into_future() => served?,
|
||||
served = Server::builder()
|
||||
.add_service(health_service)
|
||||
.add_service(identity_service)
|
||||
.add_service(local_agent_service)
|
||||
.add_service(settings_service)
|
||||
.add_service(forge_operations)
|
||||
.serve_with_shutdown(arguments.listen_grpc, shutdown()) => served?,
|
||||
worker = deletion_sweep => return Err(worker_terminated("account deletion sweep", worker)),
|
||||
worker = native_sync => return Err(worker_terminated("native contribution sync", worker)),
|
||||
worker = github_sync => return Err(worker_terminated("GitHub contribution sync", worker)),
|
||||
worker = gitlab_sync => return Err(worker_terminated("GitLab contribution sync", worker)),
|
||||
worker = codeberg_sync => return Err(worker_terminated("Codeberg contribution sync", worker)),
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn shutdown() {
|
||||
if let Err(error) = tokio::signal::ctrl_c().await {
|
||||
tracing::error!(%error, "cannot listen for shutdown signal");
|
||||
}
|
||||
}
|
||||
|
||||
fn worker_terminated(
|
||||
worker: &str,
|
||||
outcome: Result<(), tokio::task::JoinError>,
|
||||
) -> Box<dyn std::error::Error + Send + Sync> {
|
||||
let detail = match outcome {
|
||||
Ok(()) => "terminated without an error".to_owned(),
|
||||
Err(error) => error.to_string(),
|
||||
};
|
||||
std::io::Error::other(format!("{worker} terminated: {detail}")).into()
|
||||
}
|
||||
|
||||
/// Hourly worklist of confirmed deletions whose restore window has closed —
|
||||
/// identity is the one tracking the 30-day timer, so it's the one that has
|
||||
/// to remember to act on it. A failed `hard_delete` just leaves the request
|
||||
/// row in place for the next tick, no separate retry bookkeeping needed.
|
||||
fn spawn_deletion_sweep(
|
||||
store: syncode_identity_storage::Postgres,
|
||||
gitea: syncode_identity_gitea_client::GiteaClient,
|
||||
) -> tokio::task::JoinHandle<()> {
|
||||
tokio::spawn(async move {
|
||||
let mut interval = tokio::time::interval(std::time::Duration::from_secs(3600));
|
||||
loop {
|
||||
interval.tick().await;
|
||||
let expired = match store
|
||||
.find_expired_deletion_requests(chrono::Utc::now())
|
||||
.await
|
||||
{
|
||||
Ok(expired) => expired,
|
||||
Err(error) => {
|
||||
tracing::error!(%error, "cannot load expired account deletion requests");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
for (request_id, user_id) in expired {
|
||||
if let Err(error) = gitea.hard_delete(user_id).await {
|
||||
tracing::warn!(%error, %user_id, %request_id, "cannot hard-delete expired account");
|
||||
continue;
|
||||
}
|
||||
if let Err(error) = store.delete_deletion_request(request_id).await {
|
||||
tracing::error!(%error, %user_id, %request_id, "account was deleted in the forge but its deletion request remains");
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
mod sync_codeberg;
|
||||
mod sync_github;
|
||||
mod sync_gitlab;
|
||||
mod sync_native;
|
||||
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::Arc;
|
||||
|
||||
use axum::Router;
|
||||
use axum::http::{Method, header};
|
||||
use clap::Parser;
|
||||
use syncode_identity::event_outbox::{EventPublisher, EventPublisherConfig};
|
||||
use syncode_identity_storage::Postgres;
|
||||
use tokio::net::TcpListener;
|
||||
use tonic::transport::Server;
|
||||
use tower_http::cors::{AllowOrigin, CorsLayer};
|
||||
|
||||
use sync_codeberg::spawn_codeberg_contribution_sync;
|
||||
use sync_github::spawn_github_contribution_sync;
|
||||
use sync_gitlab::spawn_gitlab_contribution_sync;
|
||||
use sync_native::spawn_native_contribution_sync;
|
||||
|
||||
#[derive(Debug, Parser)]
|
||||
#[command(name = "syncode-identity", version, about)]
|
||||
struct Arguments {
|
||||
/// Where `syn`, `LocalAgent` daemons and the Go bridge open gRPC.
|
||||
#[arg(long, default_value = "127.0.0.1:8190")]
|
||||
listen_grpc: SocketAddr,
|
||||
|
||||
/// Where the Leptos front reaches GraphQL and OAuth start/callback.
|
||||
#[arg(long, default_value = "127.0.0.1:8191")]
|
||||
listen_http: SocketAddr,
|
||||
|
||||
/// Where principals and grants are kept.
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_DATABASE_URL", hide_env_values = true)]
|
||||
database: String,
|
||||
|
||||
#[arg(long, default_value_t = 8)]
|
||||
database_connections: u32,
|
||||
|
||||
/// Shared with the Go bridge; sent as `authorization: Bearer <secret>`.
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SHARED_SECRET", hide_env_values = true)]
|
||||
shared_secret: String,
|
||||
|
||||
/// This service's own externally-reachable origin — see
|
||||
/// `syncode_identity_api_oauth::Config`.
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_PUBLIC_URL")]
|
||||
public_url: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_COOKIE_DOMAIN")]
|
||||
cookie_domain: String,
|
||||
|
||||
#[arg(
|
||||
long,
|
||||
env = "SYNCODE_IDENTITY_ALLOWED_REDIRECT_ORIGINS",
|
||||
value_delimiter = ','
|
||||
)]
|
||||
allowed_redirect_origins: Vec<String>,
|
||||
|
||||
/// Where the email-verification link sends people back to.
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_FRONT_URL")]
|
||||
front_url: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_HOST")]
|
||||
smtp_host: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_PORT", default_value_t = 587)]
|
||||
smtp_port: u16,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_USERNAME")]
|
||||
smtp_username: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_PASSWORD", hide_env_values = true)]
|
||||
smtp_password: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SMTP_FROM")]
|
||||
smtp_from: String,
|
||||
|
||||
#[arg(
|
||||
long,
|
||||
env = "SYNCODE_IDENTITY_SMTP_FROM_NAME",
|
||||
default_value = "SynCode"
|
||||
)]
|
||||
smtp_from_name: String,
|
||||
|
||||
/// Gitea's origin inside the compose network — where the internal API
|
||||
/// that executes lock/unlock/hard-delete lives.
|
||||
#[arg(long, env = "SYNCODE_GITEA_INTERNAL_URL")]
|
||||
gitea_internal_url: String,
|
||||
|
||||
/// Shared with `syncode-control`'s `GITEA_INTERNAL_TOKEN`.
|
||||
#[arg(long, env = "SYNCODE_GITEA_INTERNAL_TOKEN", hide_env_values = true)]
|
||||
gitea_internal_token: String,
|
||||
|
||||
/// 32 raw bytes, base64-encoded — encrypts linked GitHub/GitLab OAuth
|
||||
/// tokens at rest (0.4 stage 5, part 2's heatmap sync needs them
|
||||
/// later). Fails fast on startup if missing or malformed, rather than
|
||||
/// silently skipping every token save.
|
||||
#[arg(
|
||||
long,
|
||||
env = "SYNCODE_IDENTITY_TOKEN_ENCRYPTION_KEY",
|
||||
hide_env_values = true
|
||||
)]
|
||||
token_encryption_key: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_COLLAB_EVENT_URL")]
|
||||
collaboration_event_url: url::Url,
|
||||
|
||||
#[arg(long, env = "SYNCODE_COLLAB_EVENT_SECRET", hide_env_values = true)]
|
||||
collaboration_event_secret: String,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_SOURCE_NODE_ID")]
|
||||
source_node_id: uuid::Uuid,
|
||||
|
||||
#[arg(
|
||||
long,
|
||||
env = "SYNCODE_IDENTITY_EVENT_INTERVAL_SECONDS",
|
||||
default_value_t = 1
|
||||
)]
|
||||
event_interval_seconds: u64,
|
||||
|
||||
#[arg(long, env = "SYNCODE_IDENTITY_EVENT_BATCH", default_value_t = 100)]
|
||||
event_batch: i64,
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
|
||||
tracing_subscriber::fmt().try_init()?;
|
||||
let arguments = Arguments::parse();
|
||||
|
||||
let token_key =
|
||||
syncode_identity_token_crypto::Key::from_base64(&arguments.token_encryption_key)?;
|
||||
|
||||
let store = Postgres::connect(&arguments.database, arguments.database_connections).await?;
|
||||
let event_publisher = EventPublisher::new(
|
||||
store.clone(),
|
||||
EventPublisherConfig {
|
||||
endpoint: arguments.collaboration_event_url.as_str(),
|
||||
secret: arguments.collaboration_event_secret,
|
||||
source_node_id: arguments.source_node_id,
|
||||
interval: std::time::Duration::from_secs(arguments.event_interval_seconds),
|
||||
batch: arguments.event_batch,
|
||||
},
|
||||
)?;
|
||||
let mailer = syncode_identity_mailer::Mailer::new(syncode_identity_mailer::Config {
|
||||
host: arguments.smtp_host,
|
||||
port: arguments.smtp_port,
|
||||
username: arguments.smtp_username,
|
||||
password: arguments.smtp_password,
|
||||
from: arguments.smtp_from,
|
||||
from_name: arguments.smtp_from_name,
|
||||
})?;
|
||||
|
||||
let (_health_reporter, health_service) = syncode_identity_api_grpc::health_service();
|
||||
let identity_bridge: Arc<dyn syncode_identity_application::IdentityBridgeUseCases> = Arc::new(
|
||||
syncode_identity_application::IdentityBridgeApplication::new(Arc::new(store.clone())),
|
||||
);
|
||||
let identity_service = syncode_identity_api_grpc::authenticated_identity_service(
|
||||
syncode_identity_api_grpc::IdentityServer::new(identity_bridge),
|
||||
arguments.shared_secret,
|
||||
);
|
||||
let gitea =
|
||||
syncode_identity_gitea_client::GiteaClient::new(syncode_identity_gitea_client::Config {
|
||||
internal_url: arguments.gitea_internal_url,
|
||||
internal_token: arguments.gitea_internal_token,
|
||||
});
|
||||
let application = Arc::new(syncode_identity_application::SettingsApplication::new(
|
||||
Arc::new(store.clone()),
|
||||
Arc::new(gitea.clone()),
|
||||
Arc::new(gitea.clone()),
|
||||
Arc::new(mailer.clone()),
|
||||
arguments.public_url.clone(),
|
||||
));
|
||||
let settings: Arc<dyn syncode_identity_application::SettingsUseCases> = application.clone();
|
||||
let sessions: Arc<dyn syncode_identity_application::SessionUseCases> = application.clone();
|
||||
let oauth: Arc<dyn syncode_identity_application::OAuthUseCases> =
|
||||
Arc::new(syncode_identity_application::OAuthApplication::new(
|
||||
Arc::new(store.clone()),
|
||||
Arc::new(gitea.clone()),
|
||||
Arc::new(gitea.clone()),
|
||||
Arc::new(mailer.clone()),
|
||||
arguments.public_url.clone(),
|
||||
));
|
||||
let local_agents: Arc<dyn syncode_identity_application::LocalAgentUseCases> = Arc::new(
|
||||
syncode_identity_application::LocalAgentApplication::new(Arc::new(store.clone())),
|
||||
);
|
||||
let local_agent_service = syncode_identity_api_grpc::GeneratedLocalAgentServer::new(
|
||||
syncode_identity_api_grpc::LocalAgentServer::new(local_agents, settings.clone()),
|
||||
);
|
||||
let settings_service = syncode_identity_api_grpc::GeneratedSettingsServer::new(
|
||||
syncode_identity_api_grpc::SettingsServer::new(settings.clone()),
|
||||
);
|
||||
let forge_operations = syncode_identity_api_grpc::GeneratedForgeOperationsServer::new(
|
||||
syncode_identity_api_grpc::ForgeOperationsServer::new(Arc::new(
|
||||
syncode_identity_application::ForgeOperationsApplication::new(Arc::new(gitea.clone())),
|
||||
)),
|
||||
);
|
||||
let allowed_origins = arguments.allowed_redirect_origins.clone();
|
||||
let cors = CorsLayer::new()
|
||||
.allow_credentials(true)
|
||||
.allow_methods([Method::GET, Method::POST, Method::OPTIONS])
|
||||
.allow_headers([header::CONTENT_TYPE])
|
||||
.allow_origin(AllowOrigin::predicate(move |origin, _| {
|
||||
let Ok(origin) = origin.to_str() else {
|
||||
return false;
|
||||
};
|
||||
allowed_origins.iter().any(|allowed| allowed == origin)
|
||||
}));
|
||||
|
||||
let deletion_sweep = spawn_deletion_sweep(store.clone(), gitea.clone());
|
||||
let native_sync = spawn_native_contribution_sync(store.clone(), gitea.clone());
|
||||
let github_sync = spawn_github_contribution_sync(store.clone(), token_key.clone());
|
||||
let gitlab_sync = spawn_gitlab_contribution_sync(store.clone(), token_key.clone());
|
||||
let codeberg_sync = spawn_codeberg_contribution_sync(store.clone());
|
||||
|
||||
let http = Router::new()
|
||||
.merge(syncode_identity_api_graphql::router(settings, sessions))
|
||||
.merge(syncode_identity_api_oauth::router(
|
||||
oauth,
|
||||
syncode_identity_api_oauth::Config {
|
||||
public_url: arguments.public_url,
|
||||
cookie_domain: arguments.cookie_domain,
|
||||
allowed_redirect_origins: arguments.allowed_redirect_origins,
|
||||
front_url: arguments.front_url,
|
||||
},
|
||||
token_key,
|
||||
))
|
||||
.layer(cors);
|
||||
let http_listener = TcpListener::bind(arguments.listen_http).await?;
|
||||
|
||||
tokio::select! {
|
||||
served = axum::serve(http_listener, http).into_future() => served?,
|
||||
served = Server::builder()
|
||||
.add_service(health_service)
|
||||
.add_service(identity_service)
|
||||
.add_service(local_agent_service)
|
||||
.add_service(settings_service)
|
||||
.add_service(forge_operations)
|
||||
.serve_with_shutdown(arguments.listen_grpc, shutdown()) => served?,
|
||||
worker = deletion_sweep => return Err(worker_terminated("account deletion sweep", worker)),
|
||||
worker = native_sync => return Err(worker_terminated("native contribution sync", worker)),
|
||||
worker = github_sync => return Err(worker_terminated("GitHub contribution sync", worker)),
|
||||
worker = gitlab_sync => return Err(worker_terminated("GitLab contribution sync", worker)),
|
||||
worker = codeberg_sync => return Err(worker_terminated("Codeberg contribution sync", worker)),
|
||||
published = event_publisher.run() => published?,
|
||||
}
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn shutdown() {
|
||||
if let Err(error) = tokio::signal::ctrl_c().await {
|
||||
tracing::error!(%error, "cannot listen for shutdown signal");
|
||||
}
|
||||
}
|
||||
|
||||
fn worker_terminated(
|
||||
worker: &str,
|
||||
outcome: Result<(), tokio::task::JoinError>,
|
||||
) -> Box<dyn std::error::Error + Send + Sync> {
|
||||
let detail = match outcome {
|
||||
Ok(()) => "terminated without an error".to_owned(),
|
||||
Err(error) => error.to_string(),
|
||||
};
|
||||
std::io::Error::other(format!("{worker} terminated: {detail}")).into()
|
||||
}
|
||||
|
||||
/// Hourly worklist of confirmed deletions whose restore window has closed —
|
||||
/// identity is the one tracking the 30-day timer, so it's the one that has
|
||||
/// to remember to act on it. A failed `hard_delete` just leaves the request
|
||||
/// row in place for the next tick, no separate retry bookkeeping needed.
|
||||
fn spawn_deletion_sweep(
|
||||
store: syncode_identity_storage::Postgres,
|
||||
gitea: syncode_identity_gitea_client::GiteaClient,
|
||||
) -> tokio::task::JoinHandle<()> {
|
||||
tokio::spawn(async move {
|
||||
let mut interval = tokio::time::interval(std::time::Duration::from_secs(3600));
|
||||
loop {
|
||||
interval.tick().await;
|
||||
let expired = match store
|
||||
.find_expired_deletion_requests(chrono::Utc::now())
|
||||
.await
|
||||
{
|
||||
Ok(expired) => expired,
|
||||
Err(error) => {
|
||||
tracing::error!(%error, "cannot load expired account deletion requests");
|
||||
continue;
|
||||
}
|
||||
};
|
||||
for (request_id, user_id) in expired {
|
||||
if let Err(error) = gitea.hard_delete(user_id).await {
|
||||
tracing::warn!(%error, %user_id, %request_id, "cannot hard-delete expired account");
|
||||
continue;
|
||||
}
|
||||
if let Err(error) = store.delete_deletion_request(request_id).await {
|
||||
tracing::error!(%error, %user_id, %request_id, "account was deleted in the forge but its deletion request remains");
|
||||
}
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -1,271 +1,244 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use syncode_identity_application::IdentityBridgeUseCases;
|
||||
use syncode_identity_model::RepositoryName;
|
||||
use tonic::{Request, Response, Status};
|
||||
|
||||
use self::conversion::{
|
||||
application_error, application_principal, model_repository_owner, model_repository_visibility,
|
||||
parse_resource, parse_resource_kind, parse_uuid, wire_principal_kind, wire_resource,
|
||||
};
|
||||
|
||||
use crate::wire::identity_server::Identity;
|
||||
use crate::wire::{
|
||||
ArchiveRepositoryProjectionRequest, ArchiveRepositoryProjectionResponse,
|
||||
ArchiveRepositoryRequest, ArchiveRepositoryResponse, CheckBranchProtectionBypassRequest,
|
||||
CheckBranchProtectionBypassResponse, CheckCapabilitiesRequest, CheckCapabilitiesResponse,
|
||||
CheckCapabilityRequest, CheckCapabilityResponse, DeleteRepositoryRequest,
|
||||
DeleteRepositoryResponse, ListBranchProtectionRulesRequest, ListBranchProtectionRulesResponse,
|
||||
ListPermittedResourcesRequest, ListPermittedResourcesResponse, ListUserSigningKeysRequest,
|
||||
ListUserSigningKeysResponse, PrincipalKind, ProvisionRepositoryRequest,
|
||||
ProvisionRepositoryResponse, RegisterRepositoryRequest, RegisterRepositoryResponse,
|
||||
ResolveRepositoryRequest, ResolveRepositoryResponse, ResolveSshKeyRequest,
|
||||
ResolveSshKeyResponse, UnarchiveRepositoryRequest, UnarchiveRepositoryResponse,
|
||||
UpdateRepositoryMetadataRequest, UpdateRepositoryMetadataResponse,
|
||||
UpdateRepositoryVisibilityRequest, UpdateRepositoryVisibilityResponse, ValidateSessionRequest,
|
||||
ValidateSessionResponse,
|
||||
};
|
||||
|
||||
mod access;
|
||||
mod conversion;
|
||||
mod repository;
|
||||
mod server;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct IdentityServer {
|
||||
application: Arc<dyn IdentityBridgeUseCases>,
|
||||
}
|
||||
|
||||
#[tonic::async_trait]
|
||||
impl Identity for IdentityServer {
|
||||
async fn resolve_repository(
|
||||
&self,
|
||||
request: Request<ResolveRepositoryRequest>,
|
||||
) -> Result<Response<ResolveRepositoryResponse>, Status> {
|
||||
repository::resolve(self, request).await
|
||||
}
|
||||
|
||||
async fn register_repository(
|
||||
&self,
|
||||
request: Request<RegisterRepositoryRequest>,
|
||||
) -> Result<Response<RegisterRepositoryResponse>, Status> {
|
||||
repository::register(self, request).await
|
||||
}
|
||||
|
||||
async fn validate_session(
|
||||
&self,
|
||||
request: Request<ValidateSessionRequest>,
|
||||
) -> Result<Response<ValidateSessionResponse>, Status> {
|
||||
let token = request.into_inner().session_token;
|
||||
let credential = self
|
||||
.application
|
||||
.validate_credential(&token)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
let (resource_kind, resource_id) = wire_resource(credential.resource);
|
||||
Ok(Response::new(ValidateSessionResponse {
|
||||
principal_id: credential.principal.get().to_string(),
|
||||
principal_kind: wire_principal_kind(credential.principal.kind()) as i32,
|
||||
expires_at_unix: credential.expires_at.timestamp(),
|
||||
owner_user_id: credential.owner_user_id.to_string(),
|
||||
capabilities: credential
|
||||
.capabilities
|
||||
.iter()
|
||||
.map(ToString::to_string)
|
||||
.collect(),
|
||||
audience: credential.audience,
|
||||
resource_kind: resource_kind as i32,
|
||||
resource_id,
|
||||
}))
|
||||
}
|
||||
|
||||
async fn check_capability(
|
||||
&self,
|
||||
request: Request<CheckCapabilityRequest>,
|
||||
) -> Result<Response<CheckCapabilityResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let resource = parse_resource(request.resource_kind(), &request.resource_id)?;
|
||||
|
||||
let principal = application_principal(request.principal_kind(), &request.principal_id)?;
|
||||
let capability = request
|
||||
.capability
|
||||
.parse()
|
||||
.map_err(|_| Status::invalid_argument("unsupported capability"))?;
|
||||
let allowed = self
|
||||
.application
|
||||
.check_capability(principal, resource, capability)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
|
||||
Ok(Response::new(CheckCapabilityResponse { allowed }))
|
||||
}
|
||||
|
||||
async fn check_capabilities(
|
||||
&self,
|
||||
request: Request<CheckCapabilitiesRequest>,
|
||||
) -> Result<Response<CheckCapabilitiesResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let principal = application_principal(request.principal_kind(), &request.principal_id)?;
|
||||
let mut pairs = Vec::with_capacity(request.queries.len());
|
||||
for query in &request.queries {
|
||||
let resource = parse_resource(query.resource_kind(), &query.resource_id)?;
|
||||
let capability = query
|
||||
.capability
|
||||
.parse()
|
||||
.map_err(|_| Status::invalid_argument("unsupported capability"))?;
|
||||
pairs.push((resource, capability));
|
||||
}
|
||||
let allowed = self
|
||||
.application
|
||||
.check_capabilities(principal, &pairs)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(CheckCapabilitiesResponse { allowed }))
|
||||
}
|
||||
|
||||
async fn list_permitted_resources(
|
||||
&self,
|
||||
request: Request<ListPermittedResourcesRequest>,
|
||||
) -> Result<Response<ListPermittedResourcesResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let principal = application_principal(request.principal_kind(), &request.principal_id)?;
|
||||
let resource_kind = parse_resource_kind(request.resource_kind())?;
|
||||
let capability = request
|
||||
.capability
|
||||
.parse()
|
||||
.map_err(|_| Status::invalid_argument("unsupported capability"))?;
|
||||
let permitted = self
|
||||
.application
|
||||
.list_permitted_resources(principal, resource_kind, capability)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ListPermittedResourcesResponse {
|
||||
resource_ids: permitted
|
||||
.resource_ids
|
||||
.iter()
|
||||
.map(ToString::to_string)
|
||||
.collect(),
|
||||
includes_public: permitted.includes_public,
|
||||
}))
|
||||
}
|
||||
|
||||
async fn resolve_ssh_key(
|
||||
&self,
|
||||
request: Request<ResolveSshKeyRequest>,
|
||||
) -> Result<Response<ResolveSshKeyResponse>, Status> {
|
||||
let fingerprint = request.into_inner().fingerprint;
|
||||
if fingerprint.is_empty() {
|
||||
return Err(Status::invalid_argument("fingerprint is required"));
|
||||
}
|
||||
let owner = self
|
||||
.application
|
||||
.resolve_ssh_key(&fingerprint)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
let principal_kind = match owner.owner {
|
||||
syncode_identity_application::SshKeyOwnerId::User(_) => PrincipalKind::User,
|
||||
syncode_identity_application::SshKeyOwnerId::LocalAgent(_) => PrincipalKind::LocalAgent,
|
||||
};
|
||||
Ok(Response::new(ResolveSshKeyResponse {
|
||||
principal_id: owner.owner.get().to_string(),
|
||||
principal_kind: principal_kind as i32,
|
||||
owner_user_id: owner.owner_user_id.to_string(),
|
||||
}))
|
||||
}
|
||||
|
||||
async fn check_branch_protection_bypass(
|
||||
&self,
|
||||
request: Request<CheckBranchProtectionBypassRequest>,
|
||||
) -> Result<Response<CheckBranchProtectionBypassResponse>, Status> {
|
||||
access::check_branch_protection_bypass(self, request).await
|
||||
}
|
||||
|
||||
async fn list_branch_protection_rules(
|
||||
&self,
|
||||
request: Request<ListBranchProtectionRulesRequest>,
|
||||
) -> Result<Response<ListBranchProtectionRulesResponse>, Status> {
|
||||
access::list_branch_protection_rules(self, request).await
|
||||
}
|
||||
|
||||
async fn list_user_signing_keys(
|
||||
&self,
|
||||
request: Request<ListUserSigningKeysRequest>,
|
||||
) -> Result<Response<ListUserSigningKeysResponse>, Status> {
|
||||
access::list_user_signing_keys(self, request).await
|
||||
}
|
||||
|
||||
async fn provision_repository(
|
||||
&self,
|
||||
request: Request<ProvisionRepositoryRequest>,
|
||||
) -> Result<Response<ProvisionRepositoryResponse>, Status> {
|
||||
repository::provision(self, request).await
|
||||
}
|
||||
|
||||
async fn update_repository_visibility(
|
||||
&self,
|
||||
request: Request<UpdateRepositoryVisibilityRequest>,
|
||||
) -> Result<Response<UpdateRepositoryVisibilityResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
self.application
|
||||
.update_repository_visibility(
|
||||
parse_uuid(&request.repository_id)?.into(),
|
||||
model_repository_visibility(request.visibility())?,
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(UpdateRepositoryVisibilityResponse {}))
|
||||
}
|
||||
|
||||
async fn update_repository_metadata(
|
||||
&self,
|
||||
request: Request<UpdateRepositoryMetadataRequest>,
|
||||
) -> Result<Response<UpdateRepositoryMetadataResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
self.application
|
||||
.update_repository_metadata(
|
||||
parse_uuid(&request.repository_id)?.into(),
|
||||
model_repository_owner(request.owner_kind(), &request.owner_id)?,
|
||||
request
|
||||
.name
|
||||
.parse::<RepositoryName>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?,
|
||||
model_repository_visibility(request.visibility())?,
|
||||
parse_uuid(&request.actor_id)?.into(),
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(UpdateRepositoryMetadataResponse {}))
|
||||
}
|
||||
|
||||
async fn archive_repository(
|
||||
&self,
|
||||
request: Request<ArchiveRepositoryRequest>,
|
||||
) -> Result<Response<ArchiveRepositoryResponse>, Status> {
|
||||
self.application
|
||||
.archive_repository(parse_uuid(&request.into_inner().repository_id)?.into())
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ArchiveRepositoryResponse {}))
|
||||
}
|
||||
|
||||
async fn archive_repository_projection(
|
||||
&self,
|
||||
request: Request<ArchiveRepositoryProjectionRequest>,
|
||||
) -> Result<Response<ArchiveRepositoryProjectionResponse>, Status> {
|
||||
repository::archive(self, request).await
|
||||
}
|
||||
|
||||
async fn unarchive_repository(
|
||||
&self,
|
||||
request: Request<UnarchiveRepositoryRequest>,
|
||||
) -> Result<Response<UnarchiveRepositoryResponse>, Status> {
|
||||
repository::unarchive(self, request).await
|
||||
}
|
||||
|
||||
async fn delete_repository(
|
||||
&self,
|
||||
request: Request<DeleteRepositoryRequest>,
|
||||
) -> Result<Response<DeleteRepositoryResponse>, Status> {
|
||||
repository::delete(self, request).await
|
||||
}
|
||||
}
|
||||
use std::sync::Arc;
|
||||
|
||||
use syncode_identity_application::IdentityBridgeUseCases;
|
||||
use tonic::{Request, Response, Status};
|
||||
|
||||
use self::conversion::{
|
||||
application_error, application_principal, parse_resource, parse_resource_kind,
|
||||
wire_principal_kind, wire_resource,
|
||||
};
|
||||
|
||||
use crate::wire::identity_server::Identity;
|
||||
use crate::wire::{
|
||||
ArchiveRepositoryProjectionRequest, ArchiveRepositoryProjectionResponse,
|
||||
ArchiveRepositoryRequest, ArchiveRepositoryResponse, CheckBranchProtectionBypassRequest,
|
||||
CheckBranchProtectionBypassResponse, CheckCapabilitiesRequest, CheckCapabilitiesResponse,
|
||||
CheckCapabilityRequest, CheckCapabilityResponse, DeleteRepositoryRequest,
|
||||
DeleteRepositoryResponse, ListBranchProtectionRulesRequest, ListBranchProtectionRulesResponse,
|
||||
ListPermittedResourcesRequest, ListPermittedResourcesResponse, ListUserSigningKeysRequest,
|
||||
ListUserSigningKeysResponse, PrincipalKind, ProvisionRepositoryRequest,
|
||||
ProvisionRepositoryResponse, RegisterRepositoryRequest, RegisterRepositoryResponse,
|
||||
ResolveRepositoryRequest, ResolveRepositoryResponse, ResolveSshKeyRequest,
|
||||
ResolveSshKeyResponse, UnarchiveRepositoryRequest, UnarchiveRepositoryResponse,
|
||||
UpdateRepositoryMetadataRequest, UpdateRepositoryMetadataResponse,
|
||||
UpdateRepositoryVisibilityRequest, UpdateRepositoryVisibilityResponse, ValidateSessionRequest,
|
||||
ValidateSessionResponse,
|
||||
};
|
||||
|
||||
mod access;
|
||||
mod conversion;
|
||||
mod repository;
|
||||
mod server;
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct IdentityServer {
|
||||
application: Arc<dyn IdentityBridgeUseCases>,
|
||||
}
|
||||
|
||||
#[tonic::async_trait]
|
||||
impl Identity for IdentityServer {
|
||||
async fn resolve_repository(
|
||||
&self,
|
||||
request: Request<ResolveRepositoryRequest>,
|
||||
) -> Result<Response<ResolveRepositoryResponse>, Status> {
|
||||
repository::resolve(self, request).await
|
||||
}
|
||||
|
||||
async fn register_repository(
|
||||
&self,
|
||||
request: Request<RegisterRepositoryRequest>,
|
||||
) -> Result<Response<RegisterRepositoryResponse>, Status> {
|
||||
repository::register(self, request).await
|
||||
}
|
||||
|
||||
async fn validate_session(
|
||||
&self,
|
||||
request: Request<ValidateSessionRequest>,
|
||||
) -> Result<Response<ValidateSessionResponse>, Status> {
|
||||
let token = request.into_inner().session_token;
|
||||
let credential = self
|
||||
.application
|
||||
.validate_credential(&token)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
let (resource_kind, resource_id) = wire_resource(credential.resource);
|
||||
Ok(Response::new(ValidateSessionResponse {
|
||||
principal_id: credential.principal.get().to_string(),
|
||||
principal_kind: wire_principal_kind(credential.principal.kind()) as i32,
|
||||
expires_at_unix: credential.expires_at.timestamp(),
|
||||
owner_user_id: credential.owner_user_id.to_string(),
|
||||
capabilities: credential
|
||||
.capabilities
|
||||
.iter()
|
||||
.map(ToString::to_string)
|
||||
.collect(),
|
||||
audience: credential.audience,
|
||||
resource_kind: resource_kind as i32,
|
||||
resource_id,
|
||||
}))
|
||||
}
|
||||
|
||||
async fn check_capability(
|
||||
&self,
|
||||
request: Request<CheckCapabilityRequest>,
|
||||
) -> Result<Response<CheckCapabilityResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let resource = parse_resource(request.resource_kind(), &request.resource_id)?;
|
||||
|
||||
let principal = application_principal(request.principal_kind(), &request.principal_id)?;
|
||||
let capability = request
|
||||
.capability
|
||||
.parse()
|
||||
.map_err(|_| Status::invalid_argument("unsupported capability"))?;
|
||||
let allowed = self
|
||||
.application
|
||||
.check_capability(principal, resource, capability)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
|
||||
Ok(Response::new(CheckCapabilityResponse { allowed }))
|
||||
}
|
||||
|
||||
async fn check_capabilities(
|
||||
&self,
|
||||
request: Request<CheckCapabilitiesRequest>,
|
||||
) -> Result<Response<CheckCapabilitiesResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let principal = application_principal(request.principal_kind(), &request.principal_id)?;
|
||||
let mut pairs = Vec::with_capacity(request.queries.len());
|
||||
for query in &request.queries {
|
||||
let resource = parse_resource(query.resource_kind(), &query.resource_id)?;
|
||||
let capability = query
|
||||
.capability
|
||||
.parse()
|
||||
.map_err(|_| Status::invalid_argument("unsupported capability"))?;
|
||||
pairs.push((resource, capability));
|
||||
}
|
||||
let allowed = self
|
||||
.application
|
||||
.check_capabilities(principal, &pairs)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(CheckCapabilitiesResponse { allowed }))
|
||||
}
|
||||
|
||||
async fn list_permitted_resources(
|
||||
&self,
|
||||
request: Request<ListPermittedResourcesRequest>,
|
||||
) -> Result<Response<ListPermittedResourcesResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let principal = application_principal(request.principal_kind(), &request.principal_id)?;
|
||||
let resource_kind = parse_resource_kind(request.resource_kind())?;
|
||||
let capability = request
|
||||
.capability
|
||||
.parse()
|
||||
.map_err(|_| Status::invalid_argument("unsupported capability"))?;
|
||||
let permitted = self
|
||||
.application
|
||||
.list_permitted_resources(principal, resource_kind, capability)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ListPermittedResourcesResponse {
|
||||
resource_ids: permitted
|
||||
.resource_ids
|
||||
.iter()
|
||||
.map(ToString::to_string)
|
||||
.collect(),
|
||||
includes_public: permitted.includes_public,
|
||||
}))
|
||||
}
|
||||
|
||||
async fn resolve_ssh_key(
|
||||
&self,
|
||||
request: Request<ResolveSshKeyRequest>,
|
||||
) -> Result<Response<ResolveSshKeyResponse>, Status> {
|
||||
let fingerprint = request.into_inner().fingerprint;
|
||||
if fingerprint.is_empty() {
|
||||
return Err(Status::invalid_argument("fingerprint is required"));
|
||||
}
|
||||
let owner = self
|
||||
.application
|
||||
.resolve_ssh_key(&fingerprint)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
let principal_kind = match owner.owner {
|
||||
syncode_identity_application::SshKeyOwnerId::User(_) => PrincipalKind::User,
|
||||
syncode_identity_application::SshKeyOwnerId::LocalAgent(_) => PrincipalKind::LocalAgent,
|
||||
};
|
||||
Ok(Response::new(ResolveSshKeyResponse {
|
||||
principal_id: owner.owner.get().to_string(),
|
||||
principal_kind: principal_kind as i32,
|
||||
owner_user_id: owner.owner_user_id.to_string(),
|
||||
}))
|
||||
}
|
||||
|
||||
async fn check_branch_protection_bypass(
|
||||
&self,
|
||||
request: Request<CheckBranchProtectionBypassRequest>,
|
||||
) -> Result<Response<CheckBranchProtectionBypassResponse>, Status> {
|
||||
access::check_branch_protection_bypass(self, request).await
|
||||
}
|
||||
|
||||
async fn list_branch_protection_rules(
|
||||
&self,
|
||||
request: Request<ListBranchProtectionRulesRequest>,
|
||||
) -> Result<Response<ListBranchProtectionRulesResponse>, Status> {
|
||||
access::list_branch_protection_rules(self, request).await
|
||||
}
|
||||
|
||||
async fn list_user_signing_keys(
|
||||
&self,
|
||||
request: Request<ListUserSigningKeysRequest>,
|
||||
) -> Result<Response<ListUserSigningKeysResponse>, Status> {
|
||||
access::list_user_signing_keys(self, request).await
|
||||
}
|
||||
|
||||
async fn provision_repository(
|
||||
&self,
|
||||
request: Request<ProvisionRepositoryRequest>,
|
||||
) -> Result<Response<ProvisionRepositoryResponse>, Status> {
|
||||
repository::provision(self, request).await
|
||||
}
|
||||
|
||||
async fn update_repository_visibility(
|
||||
&self,
|
||||
request: Request<UpdateRepositoryVisibilityRequest>,
|
||||
) -> Result<Response<UpdateRepositoryVisibilityResponse>, Status> {
|
||||
repository::update_visibility(self, request).await
|
||||
}
|
||||
|
||||
async fn update_repository_metadata(
|
||||
&self,
|
||||
request: Request<UpdateRepositoryMetadataRequest>,
|
||||
) -> Result<Response<UpdateRepositoryMetadataResponse>, Status> {
|
||||
repository::update_metadata(self, request).await
|
||||
}
|
||||
|
||||
async fn archive_repository(
|
||||
&self,
|
||||
request: Request<ArchiveRepositoryRequest>,
|
||||
) -> Result<Response<ArchiveRepositoryResponse>, Status> {
|
||||
repository::archive_repository(self, request).await
|
||||
}
|
||||
|
||||
async fn archive_repository_projection(
|
||||
&self,
|
||||
request: Request<ArchiveRepositoryProjectionRequest>,
|
||||
) -> Result<Response<ArchiveRepositoryProjectionResponse>, Status> {
|
||||
repository::archive(self, request).await
|
||||
}
|
||||
|
||||
async fn unarchive_repository(
|
||||
&self,
|
||||
request: Request<UnarchiveRepositoryRequest>,
|
||||
) -> Result<Response<UnarchiveRepositoryResponse>, Status> {
|
||||
repository::unarchive(self, request).await
|
||||
}
|
||||
|
||||
async fn delete_repository(
|
||||
&self,
|
||||
request: Request<DeleteRepositoryRequest>,
|
||||
) -> Result<Response<DeleteRepositoryResponse>, Status> {
|
||||
repository::delete(self, request).await
|
||||
}
|
||||
}
|
||||
@@ -1,263 +1,247 @@
|
||||
use async_trait::async_trait;
|
||||
use syncode_identity_application::{
|
||||
AccessTokenIdentity, ActiveBranchProtectionRule, ApplicationError, IdentityBridgeRepository,
|
||||
LocalAgentIdentity, SigningKeyIdentity, SshKeyIdentity,
|
||||
};
|
||||
use syncode_identity_model::{
|
||||
AccessTokenId, Grant, GrantPrincipal, LocalAgentId, PlatformAgentId, RepositoryId,
|
||||
RepositoryName, RepositoryOwner, RepositoryOwnerId, RepositoryVisibility, Resource, UserId,
|
||||
UserPrincipal,
|
||||
};
|
||||
|
||||
use crate::{Postgres, application_persistence as persistence};
|
||||
|
||||
mod repository;
|
||||
mod repository_grants;
|
||||
mod repository_registration;
|
||||
mod token;
|
||||
|
||||
use token::{access_token, active_branch_protection_rule, signing_key, ssh_key};
|
||||
|
||||
#[async_trait]
|
||||
impl IdentityBridgeRepository for Postgres {
|
||||
async fn validate_web_session(
|
||||
&self,
|
||||
token: &str,
|
||||
) -> Result<Option<(UserId, chrono::DateTime<chrono::Utc>)>, ApplicationError> {
|
||||
self.validate_session(token)
|
||||
.await
|
||||
.map(|session| session.map(|(id, expires_at)| (id.into(), expires_at)))
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn validate_access_token(
|
||||
&self,
|
||||
token: &str,
|
||||
) -> Result<Option<AccessTokenIdentity>, ApplicationError> {
|
||||
Postgres::validate_access_token(self, token)
|
||||
.await
|
||||
.map(|token| token.map(access_token))
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn load_access_token(
|
||||
&self,
|
||||
token_id: AccessTokenId,
|
||||
) -> Result<Option<AccessTokenIdentity>, ApplicationError> {
|
||||
Postgres::load_access_token(self, token_id.get())
|
||||
.await
|
||||
.map(|token| token.map(access_token))
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn load_user_principal(
|
||||
&self,
|
||||
user_id: UserId,
|
||||
) -> Result<UserPrincipal, ApplicationError> {
|
||||
Postgres::load_user_principal(self, user_id.get())
|
||||
.await
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn load_granted_resources(
|
||||
&self,
|
||||
principals: &[(syncode_identity_model::GrantPrincipalKind, uuid::Uuid)],
|
||||
resource_kind: syncode_identity_model::ResourceKind,
|
||||
) -> Result<
|
||||
Vec<(
|
||||
uuid::Uuid,
|
||||
std::collections::BTreeSet<syncode_identity_model::Capability>,
|
||||
)>,
|
||||
ApplicationError,
|
||||
> {
|
||||
Postgres::load_granted_resources(self, principals, resource_kind)
|
||||
.await
|
||||
.map_err(|error| ApplicationError::Persistence(error.to_string()))
|
||||
}
|
||||
|
||||
async fn load_active_grants(&self, resource: Resource) -> Result<Vec<Grant>, ApplicationError> {
|
||||
Postgres::load_active_grants(self, resource)
|
||||
.await
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn load_active_repository_visibility(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
) -> Result<Option<RepositoryVisibility>, ApplicationError> {
|
||||
self.bridge_repository_visibility(repository_id.get()).await
|
||||
}
|
||||
|
||||
async fn load_active_repository_owner(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
) -> Result<Option<RepositoryOwnerId>, ApplicationError> {
|
||||
self.bridge_repository_owner(repository_id.get()).await
|
||||
}
|
||||
|
||||
async fn platform_agent_is_active(
|
||||
&self,
|
||||
agent_id: PlatformAgentId,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
Postgres::platform_agent_is_active(self, agent_id.get())
|
||||
.await
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn local_agent(
|
||||
&self,
|
||||
agent_id: LocalAgentId,
|
||||
) -> Result<Option<LocalAgentIdentity>, ApplicationError> {
|
||||
self.load_local_agent(agent_id.get())
|
||||
.await
|
||||
.map(|agent| {
|
||||
agent.map(|agent| LocalAgentIdentity {
|
||||
owner_user_id: agent.owner_user_id.into(),
|
||||
restriction: agent.restriction,
|
||||
})
|
||||
})
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn resolve_ssh_key(
|
||||
&self,
|
||||
fingerprint: &str,
|
||||
) -> Result<Option<SshKeyIdentity>, ApplicationError> {
|
||||
Postgres::resolve_ssh_key(self, fingerprint)
|
||||
.await
|
||||
.map_err(persistence)?
|
||||
.map(ssh_key)
|
||||
.transpose()
|
||||
}
|
||||
|
||||
async fn branch_protection_bypass_allowed(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
pattern: &str,
|
||||
principal: GrantPrincipal,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
Postgres::branch_protection_bypass_allowed(
|
||||
self,
|
||||
repository_id.get(),
|
||||
pattern,
|
||||
principal.kind(),
|
||||
principal.get(),
|
||||
)
|
||||
.await
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn active_branch_protection_rules(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
) -> Result<Vec<ActiveBranchProtectionRule>, ApplicationError> {
|
||||
self.find_branch_protection_rules(repository_id.get())
|
||||
.await
|
||||
.map(|rules| {
|
||||
rules
|
||||
.into_iter()
|
||||
.map(active_branch_protection_rule)
|
||||
.collect()
|
||||
})
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn user_signing_keys(
|
||||
&self,
|
||||
user_id: UserId,
|
||||
) -> Result<Vec<SigningKeyIdentity>, ApplicationError> {
|
||||
self.find_user_signing_keys(user_id.get())
|
||||
.await
|
||||
.map(|keys| keys.into_iter().map(signing_key).collect())
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn resolve_repository(
|
||||
&self,
|
||||
owner: &RepositoryOwner,
|
||||
name: &RepositoryName,
|
||||
) -> Result<Option<RepositoryId>, ApplicationError> {
|
||||
self.bridge_resolve_repository(owner, name).await
|
||||
}
|
||||
|
||||
async fn provision_repository(
|
||||
&self,
|
||||
owner: RepositoryOwnerId,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
creator_id: UserId,
|
||||
) -> Result<RepositoryId, ApplicationError> {
|
||||
self.bridge_provision_repository(
|
||||
owner.kind(),
|
||||
owner.get(),
|
||||
name,
|
||||
visibility,
|
||||
creator_id.get(),
|
||||
)
|
||||
.await
|
||||
.map(RepositoryId::from)
|
||||
}
|
||||
|
||||
async fn register_repository(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
owner: RepositoryOwnerId,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
creator_id: UserId,
|
||||
) -> Result<(), ApplicationError> {
|
||||
self.bridge_register_repository(
|
||||
repository_id.get(),
|
||||
owner.kind(),
|
||||
owner.get(),
|
||||
name,
|
||||
visibility,
|
||||
creator_id.get(),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn update_repository_visibility(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
visibility: RepositoryVisibility,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
self.bridge_update_repository_visibility(repository_id.get(), visibility)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn update_repository_metadata(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
owner: RepositoryOwnerId,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
actor_id: UserId,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
self.bridge_update_repository_metadata(
|
||||
repository_id.get(),
|
||||
owner.kind(),
|
||||
owner.get(),
|
||||
name,
|
||||
visibility,
|
||||
actor_id.get(),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn archive_repository(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
self.bridge_archive_repository(repository_id.get()).await
|
||||
}
|
||||
|
||||
async fn set_repository_archived(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
archived: bool,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
self.bridge_set_repository_archived(repository_id.get(), archived)
|
||||
.await
|
||||
}
|
||||
}
|
||||
use async_trait::async_trait;
|
||||
use syncode_identity_application::{
|
||||
AccessTokenIdentity, ActiveBranchProtectionRule, ApplicationError, IdentityBridgeRepository,
|
||||
LocalAgentIdentity, SigningKeyIdentity, SshKeyIdentity,
|
||||
};
|
||||
use syncode_identity_model::{
|
||||
AccessTokenId, Grant, GrantPrincipal, LocalAgentId, PlatformAgentId, RepositoryId,
|
||||
RepositoryName, RepositoryOwner, RepositoryOwnerId, RepositoryVisibility, Resource, UserId,
|
||||
UserPrincipal,
|
||||
};
|
||||
|
||||
use crate::{Postgres, application_persistence as persistence};
|
||||
|
||||
mod identities;
|
||||
mod repository;
|
||||
mod repository_grants;
|
||||
mod repository_metadata;
|
||||
mod repository_registration;
|
||||
mod token;
|
||||
|
||||
use token::{access_token, ssh_key};
|
||||
|
||||
#[async_trait]
|
||||
impl IdentityBridgeRepository for Postgres {
|
||||
async fn validate_web_session(
|
||||
&self,
|
||||
token: &str,
|
||||
) -> Result<Option<(UserId, chrono::DateTime<chrono::Utc>)>, ApplicationError> {
|
||||
self.validate_session(token)
|
||||
.await
|
||||
.map(|session| session.map(|(id, expires_at)| (id.into(), expires_at)))
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn validate_access_token(
|
||||
&self,
|
||||
token: &str,
|
||||
) -> Result<Option<AccessTokenIdentity>, ApplicationError> {
|
||||
Postgres::validate_access_token(self, token)
|
||||
.await
|
||||
.map(|token| token.map(access_token))
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn load_access_token(
|
||||
&self,
|
||||
token_id: AccessTokenId,
|
||||
) -> Result<Option<AccessTokenIdentity>, ApplicationError> {
|
||||
Postgres::load_access_token(self, token_id.get())
|
||||
.await
|
||||
.map(|token| token.map(access_token))
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn load_user_principal(
|
||||
&self,
|
||||
user_id: UserId,
|
||||
) -> Result<UserPrincipal, ApplicationError> {
|
||||
Postgres::load_user_principal(self, user_id.get())
|
||||
.await
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn load_granted_resources(
|
||||
&self,
|
||||
principals: &[(syncode_identity_model::GrantPrincipalKind, uuid::Uuid)],
|
||||
resource_kind: syncode_identity_model::ResourceKind,
|
||||
) -> Result<
|
||||
Vec<(
|
||||
uuid::Uuid,
|
||||
std::collections::BTreeSet<syncode_identity_model::Capability>,
|
||||
)>,
|
||||
ApplicationError,
|
||||
> {
|
||||
Postgres::load_granted_resources(self, principals, resource_kind)
|
||||
.await
|
||||
.map_err(|error| ApplicationError::Persistence(error.to_string()))
|
||||
}
|
||||
|
||||
async fn load_active_grants(&self, resource: Resource) -> Result<Vec<Grant>, ApplicationError> {
|
||||
Postgres::load_active_grants(self, resource)
|
||||
.await
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn load_active_repository_visibility(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
) -> Result<Option<RepositoryVisibility>, ApplicationError> {
|
||||
self.bridge_repository_visibility(repository_id.get()).await
|
||||
}
|
||||
|
||||
async fn load_active_repository_owner(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
) -> Result<Option<RepositoryOwnerId>, ApplicationError> {
|
||||
self.bridge_repository_owner(repository_id.get()).await
|
||||
}
|
||||
|
||||
async fn platform_agent_is_active(
|
||||
&self,
|
||||
agent_id: PlatformAgentId,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
Postgres::platform_agent_is_active(self, agent_id.get())
|
||||
.await
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn local_agent(
|
||||
&self,
|
||||
agent_id: LocalAgentId,
|
||||
) -> Result<Option<LocalAgentIdentity>, ApplicationError> {
|
||||
self.bridge_local_agent(agent_id.get()).await
|
||||
}
|
||||
|
||||
async fn resolve_ssh_key(
|
||||
&self,
|
||||
fingerprint: &str,
|
||||
) -> Result<Option<SshKeyIdentity>, ApplicationError> {
|
||||
Postgres::resolve_ssh_key(self, fingerprint)
|
||||
.await
|
||||
.map_err(persistence)?
|
||||
.map(ssh_key)
|
||||
.transpose()
|
||||
}
|
||||
|
||||
async fn branch_protection_bypass_allowed(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
pattern: &str,
|
||||
principal: GrantPrincipal,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
Postgres::branch_protection_bypass_allowed(
|
||||
self,
|
||||
repository_id.get(),
|
||||
pattern,
|
||||
principal.kind(),
|
||||
principal.get(),
|
||||
)
|
||||
.await
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
async fn active_branch_protection_rules(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
) -> Result<Vec<ActiveBranchProtectionRule>, ApplicationError> {
|
||||
self.bridge_active_branch_protection_rules(repository_id.get())
|
||||
.await
|
||||
}
|
||||
|
||||
async fn user_signing_keys(
|
||||
&self,
|
||||
user_id: UserId,
|
||||
) -> Result<Vec<SigningKeyIdentity>, ApplicationError> {
|
||||
self.bridge_user_signing_keys(user_id.get()).await
|
||||
}
|
||||
|
||||
async fn resolve_repository(
|
||||
&self,
|
||||
owner: &RepositoryOwner,
|
||||
name: &RepositoryName,
|
||||
) -> Result<Option<RepositoryId>, ApplicationError> {
|
||||
self.bridge_resolve_repository(owner, name).await
|
||||
}
|
||||
|
||||
async fn provision_repository(
|
||||
&self,
|
||||
owner: RepositoryOwnerId,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
creator_id: UserId,
|
||||
) -> Result<RepositoryId, ApplicationError> {
|
||||
self.bridge_provision_repository(
|
||||
owner.kind(),
|
||||
owner.get(),
|
||||
name,
|
||||
visibility,
|
||||
creator_id.get(),
|
||||
)
|
||||
.await
|
||||
.map(RepositoryId::from)
|
||||
}
|
||||
|
||||
async fn register_repository(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
owner: RepositoryOwnerId,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
creator_id: UserId,
|
||||
) -> Result<(), ApplicationError> {
|
||||
self.bridge_register_repository(
|
||||
repository_id.get(),
|
||||
owner.kind(),
|
||||
owner.get(),
|
||||
name,
|
||||
visibility,
|
||||
creator_id.get(),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn update_repository_visibility(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
visibility: RepositoryVisibility,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
self.bridge_update_repository_visibility(repository_id.get(), visibility)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn update_repository_metadata(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
owner: RepositoryOwnerId,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
actor_id: UserId,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
self.bridge_update_repository_metadata(
|
||||
repository_id.get(),
|
||||
owner.kind(),
|
||||
owner.get(),
|
||||
name,
|
||||
visibility,
|
||||
actor_id.get(),
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn archive_repository(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
self.bridge_archive_repository(repository_id.get()).await
|
||||
}
|
||||
|
||||
async fn set_repository_archived(
|
||||
&self,
|
||||
repository_id: RepositoryId,
|
||||
archived: bool,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
self.bridge_set_repository_archived(repository_id.get(), archived)
|
||||
.await
|
||||
}
|
||||
}
|
||||
@@ -1,131 +1,132 @@
|
||||
mod access_tokens;
|
||||
mod account;
|
||||
mod account_deletion;
|
||||
mod application;
|
||||
mod branch_protection;
|
||||
mod bridge;
|
||||
mod cli_login;
|
||||
mod contributions;
|
||||
mod email_verification;
|
||||
mod emails;
|
||||
mod events;
|
||||
mod grants;
|
||||
mod heatmap_query;
|
||||
mod linked_identity_tokens;
|
||||
mod local_agent_application;
|
||||
mod local_agent_sessions;
|
||||
mod local_agents;
|
||||
mod oauth;
|
||||
mod oauth_application;
|
||||
mod oauth_user_projection;
|
||||
mod personal_tokens;
|
||||
mod platform_agents;
|
||||
mod principals;
|
||||
mod repository_identity;
|
||||
mod seed;
|
||||
mod sessions;
|
||||
mod settings;
|
||||
mod ssh_keys;
|
||||
mod teams;
|
||||
|
||||
pub use access_tokens::AccessTokenRecord;
|
||||
pub use account::{LinkedIdentityRecord, UnlinkError, UserRecord};
|
||||
pub use branch_protection::{BranchProtectionBypassRecord, BranchProtectionRuleRecord};
|
||||
pub use contributions::{LinkedIdentityTokenRecord, LinkedUsername};
|
||||
pub use email_verification::PendingEmailVerification;
|
||||
pub use emails::{RemoveEmailError, SetPrimaryEmailError, UserEmailRecord};
|
||||
pub use grants::GrantRecord;
|
||||
pub use heatmap_query::{HeatmapPrivacySetting, NativeContributionDay, ProviderContributionDay};
|
||||
pub use local_agent_sessions::{
|
||||
EnrollLocalAgentSession, EnrolledLocalAgentSession, HeartbeatLocalAgentSession,
|
||||
LocalAgentHeartbeat,
|
||||
};
|
||||
pub use local_agents::{CreateLocalAgentError, LocalAgentRecord};
|
||||
pub use oauth::{AuthProvider, LinkError, OAuthClientConfig};
|
||||
pub use platform_agents::{AgentDefinitionRecord, PlatformAgentGrantRecord, PlatformAgentRecord};
|
||||
pub use repository_identity::RepositoryIdentityRecord;
|
||||
pub use seed::MIGRATION_SYSTEM_USERNAME;
|
||||
pub use settings::{
|
||||
AddKeyError, CreateOrganizationError, PersonalTokenRecord, SigningKeyRecord, UserSshKeyRecord,
|
||||
};
|
||||
pub use ssh_keys::{AddSshKeyError, SshKeyOwner};
|
||||
pub use syncode_identity_application::SESSION_COOKIE_NAME;
|
||||
pub use teams::{OrganizationRecord, TeamGrantRecord, TeamRecord};
|
||||
|
||||
use sqlx::migrate::MigrateError;
|
||||
use sqlx::postgres::{PgPool, PgPoolOptions};
|
||||
use syncode_identity_model::{Capability, ForgeProjection, ForgeProjectionStatus};
|
||||
use thiserror::Error;
|
||||
|
||||
pub(crate) fn application_persistence(
|
||||
error: impl std::fmt::Display,
|
||||
) -> syncode_identity_application::ApplicationError {
|
||||
syncode_identity_application::ApplicationError::Persistence(error.to_string())
|
||||
}
|
||||
|
||||
#[derive(Debug, Error)]
|
||||
pub enum StoreError {
|
||||
#[error("the store cannot be reached: {0}")]
|
||||
Unreachable(#[from] sqlx::Error),
|
||||
|
||||
#[error("the schema migrations did not apply: {0}")]
|
||||
Migration(#[from] MigrateError),
|
||||
|
||||
#[error("a stored grant or resource kind is not one this build recognizes: {0}")]
|
||||
UnknownKind(String),
|
||||
|
||||
#[error("a provider's oauth_client_config does not match the expected shape: {0}")]
|
||||
MalformedConfig(#[source] serde_json::Error),
|
||||
|
||||
#[error("a platform agent policy cannot be encoded: {0}")]
|
||||
MalformedPolicy(#[source] serde_json::Error),
|
||||
|
||||
#[error("imported identity data conflicts with existing state: {0}")]
|
||||
ImportConflict(String),
|
||||
}
|
||||
|
||||
pub(crate) fn decode_capabilities(
|
||||
values: Vec<String>,
|
||||
) -> Result<std::collections::BTreeSet<Capability>, StoreError> {
|
||||
values
|
||||
.into_iter()
|
||||
.map(|value| {
|
||||
value
|
||||
.parse()
|
||||
.map_err(|_| StoreError::UnknownKind(format!("capability:{value}")))
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
pub(crate) fn decode_forge_projection(
|
||||
status: String,
|
||||
error: Option<String>,
|
||||
) -> Result<ForgeProjection, StoreError> {
|
||||
Ok(ForgeProjection {
|
||||
status: status
|
||||
.parse::<ForgeProjectionStatus>()
|
||||
.map_err(|_| StoreError::UnknownKind(format!("forge_projection_status:{status}")))?,
|
||||
error,
|
||||
})
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Postgres {
|
||||
pool: PgPool,
|
||||
}
|
||||
|
||||
impl Postgres {
|
||||
pub async fn connect(url: &str, connections: u32) -> Result<Self, StoreError> {
|
||||
let pool = PgPoolOptions::new()
|
||||
.max_connections(connections)
|
||||
.connect(url)
|
||||
.await?;
|
||||
sqlx::migrate!("./migrations").run(&pool).await?;
|
||||
Ok(Self { pool })
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub const fn pool(&self) -> &PgPool {
|
||||
&self.pool
|
||||
}
|
||||
}
|
||||
mod access_tokens;
|
||||
mod account;
|
||||
mod account_deletion;
|
||||
mod application;
|
||||
mod branch_protection;
|
||||
mod bridge;
|
||||
mod cli_login;
|
||||
mod contributions;
|
||||
mod email_verification;
|
||||
mod emails;
|
||||
mod events;
|
||||
mod grants;
|
||||
mod heatmap_query;
|
||||
mod linked_identity_tokens;
|
||||
mod local_agent_application;
|
||||
mod local_agent_sessions;
|
||||
mod local_agents;
|
||||
mod oauth;
|
||||
mod oauth_application;
|
||||
mod oauth_user_projection;
|
||||
mod personal_tokens;
|
||||
mod platform_agents;
|
||||
mod principals;
|
||||
mod repository_identity;
|
||||
mod seed;
|
||||
mod sessions;
|
||||
mod settings;
|
||||
mod ssh_keys;
|
||||
mod teams;
|
||||
|
||||
pub use access_tokens::AccessTokenRecord;
|
||||
pub use account::{LinkedIdentityRecord, UnlinkError, UserRecord};
|
||||
pub use branch_protection::{BranchProtectionBypassRecord, BranchProtectionRuleRecord};
|
||||
pub use contributions::{LinkedIdentityTokenRecord, LinkedUsername};
|
||||
pub use email_verification::PendingEmailVerification;
|
||||
pub use emails::{RemoveEmailError, SetPrimaryEmailError, UserEmailRecord};
|
||||
pub use events::PendingEvent;
|
||||
pub use grants::GrantRecord;
|
||||
pub use heatmap_query::{HeatmapPrivacySetting, NativeContributionDay, ProviderContributionDay};
|
||||
pub use local_agent_sessions::{
|
||||
EnrollLocalAgentSession, EnrolledLocalAgentSession, HeartbeatLocalAgentSession,
|
||||
LocalAgentHeartbeat,
|
||||
};
|
||||
pub use local_agents::{CreateLocalAgentError, LocalAgentRecord};
|
||||
pub use oauth::{AuthProvider, LinkError, OAuthClientConfig};
|
||||
pub use platform_agents::{AgentDefinitionRecord, PlatformAgentGrantRecord, PlatformAgentRecord};
|
||||
pub use repository_identity::RepositoryIdentityRecord;
|
||||
pub use seed::MIGRATION_SYSTEM_USERNAME;
|
||||
pub use settings::{
|
||||
AddKeyError, CreateOrganizationError, PersonalTokenRecord, SigningKeyRecord, UserSshKeyRecord,
|
||||
};
|
||||
pub use ssh_keys::{AddSshKeyError, SshKeyOwner};
|
||||
pub use syncode_identity_application::SESSION_COOKIE_NAME;
|
||||
pub use teams::{OrganizationRecord, TeamGrantRecord, TeamRecord};
|
||||
|
||||
use sqlx::migrate::MigrateError;
|
||||
use sqlx::postgres::{PgPool, PgPoolOptions};
|
||||
use syncode_identity_model::{Capability, ForgeProjection, ForgeProjectionStatus};
|
||||
use thiserror::Error;
|
||||
|
||||
pub(crate) fn application_persistence(
|
||||
error: impl std::fmt::Display,
|
||||
) -> syncode_identity_application::ApplicationError {
|
||||
syncode_identity_application::ApplicationError::Persistence(error.to_string())
|
||||
}
|
||||
|
||||
#[derive(Debug, Error)]
|
||||
pub enum StoreError {
|
||||
#[error("the store cannot be reached: {0}")]
|
||||
Unreachable(#[from] sqlx::Error),
|
||||
|
||||
#[error("the schema migrations did not apply: {0}")]
|
||||
Migration(#[from] MigrateError),
|
||||
|
||||
#[error("a stored grant or resource kind is not one this build recognizes: {0}")]
|
||||
UnknownKind(String),
|
||||
|
||||
#[error("a provider's oauth_client_config does not match the expected shape: {0}")]
|
||||
MalformedConfig(#[source] serde_json::Error),
|
||||
|
||||
#[error("a platform agent policy cannot be encoded: {0}")]
|
||||
MalformedPolicy(#[source] serde_json::Error),
|
||||
|
||||
#[error("imported identity data conflicts with existing state: {0}")]
|
||||
ImportConflict(String),
|
||||
}
|
||||
|
||||
pub(crate) fn decode_capabilities(
|
||||
values: Vec<String>,
|
||||
) -> Result<std::collections::BTreeSet<Capability>, StoreError> {
|
||||
values
|
||||
.into_iter()
|
||||
.map(|value| {
|
||||
value
|
||||
.parse()
|
||||
.map_err(|_| StoreError::UnknownKind(format!("capability:{value}")))
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
pub(crate) fn decode_forge_projection(
|
||||
status: String,
|
||||
error: Option<String>,
|
||||
) -> Result<ForgeProjection, StoreError> {
|
||||
Ok(ForgeProjection {
|
||||
status: status
|
||||
.parse::<ForgeProjectionStatus>()
|
||||
.map_err(|_| StoreError::UnknownKind(format!("forge_projection_status:{status}")))?,
|
||||
error,
|
||||
})
|
||||
}
|
||||
|
||||
#[derive(Clone)]
|
||||
pub struct Postgres {
|
||||
pool: PgPool,
|
||||
}
|
||||
|
||||
impl Postgres {
|
||||
pub async fn connect(url: &str, connections: u32) -> Result<Self, StoreError> {
|
||||
let pool = PgPoolOptions::new()
|
||||
.max_connections(connections)
|
||||
.connect(url)
|
||||
.await?;
|
||||
sqlx::migrate!("./migrations").run(&pool).await?;
|
||||
Ok(Self { pool })
|
||||
}
|
||||
|
||||
#[must_use]
|
||||
pub const fn pool(&self) -> &PgPool {
|
||||
&self.pool
|
||||
}
|
||||
}
|
||||
@@ -1,127 +1,179 @@
|
||||
use syncode_identity_model::{RepositoryName, RepositoryOwner};
|
||||
use tonic::{Request, Response, Status};
|
||||
|
||||
use super::IdentityServer;
|
||||
use super::conversion::{
|
||||
application_error, model_repository_owner, model_repository_visibility, parse_uuid,
|
||||
};
|
||||
use crate::wire::{
|
||||
ArchiveRepositoryProjectionRequest, ArchiveRepositoryProjectionResponse,
|
||||
DeleteRepositoryRequest, DeleteRepositoryResponse, ProvisionRepositoryRequest,
|
||||
ProvisionRepositoryResponse, RegisterRepositoryRequest, RegisterRepositoryResponse,
|
||||
ResolveRepositoryRequest, ResolveRepositoryResponse, UnarchiveRepositoryRequest,
|
||||
UnarchiveRepositoryResponse,
|
||||
};
|
||||
|
||||
pub(super) async fn resolve(
|
||||
server: &IdentityServer,
|
||||
request: Request<ResolveRepositoryRequest>,
|
||||
) -> Result<Response<ResolveRepositoryResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let owner = request
|
||||
.owner
|
||||
.parse::<RepositoryOwner>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?;
|
||||
let name = request
|
||||
.name
|
||||
.parse::<RepositoryName>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?;
|
||||
let repository_id = server
|
||||
.application
|
||||
.resolve_repository(owner, name)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ResolveRepositoryResponse {
|
||||
repository_id: repository_id.to_string(),
|
||||
}))
|
||||
}
|
||||
|
||||
pub(super) async fn provision(
|
||||
server: &IdentityServer,
|
||||
request: Request<ProvisionRepositoryRequest>,
|
||||
) -> Result<Response<ProvisionRepositoryResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let owner = model_repository_owner(request.owner_kind(), &request.owner_id)?;
|
||||
let visibility = model_repository_visibility(request.visibility())?;
|
||||
let name = request
|
||||
.name
|
||||
.parse::<RepositoryName>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?;
|
||||
let repository_id = server
|
||||
.application
|
||||
.provision_repository(
|
||||
owner,
|
||||
name,
|
||||
visibility,
|
||||
parse_uuid(&request.creator_id)?.into(),
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ProvisionRepositoryResponse {
|
||||
repository_id: repository_id.to_string(),
|
||||
}))
|
||||
}
|
||||
|
||||
pub(super) async fn register(
|
||||
server: &IdentityServer,
|
||||
request: Request<RegisterRepositoryRequest>,
|
||||
) -> Result<Response<RegisterRepositoryResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
server
|
||||
.application
|
||||
.register_repository(
|
||||
parse_uuid(&request.repository_id)?.into(),
|
||||
model_repository_owner(request.owner_kind(), &request.owner_id)?,
|
||||
request
|
||||
.name
|
||||
.parse::<RepositoryName>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?,
|
||||
model_repository_visibility(request.visibility())?,
|
||||
parse_uuid(&request.creator_id)?.into(),
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(RegisterRepositoryResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn archive(
|
||||
server: &IdentityServer,
|
||||
request: Request<ArchiveRepositoryProjectionRequest>,
|
||||
) -> Result<Response<ArchiveRepositoryProjectionResponse>, Status> {
|
||||
server
|
||||
.application
|
||||
.set_repository_archived(
|
||||
parse_uuid(&request.into_inner().repository_id)?.into(),
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ArchiveRepositoryProjectionResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn unarchive(
|
||||
server: &IdentityServer,
|
||||
request: Request<UnarchiveRepositoryRequest>,
|
||||
) -> Result<Response<UnarchiveRepositoryResponse>, Status> {
|
||||
server
|
||||
.application
|
||||
.set_repository_archived(
|
||||
parse_uuid(&request.into_inner().repository_id)?.into(),
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(UnarchiveRepositoryResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn delete(
|
||||
server: &IdentityServer,
|
||||
request: Request<DeleteRepositoryRequest>,
|
||||
) -> Result<Response<DeleteRepositoryResponse>, Status> {
|
||||
server
|
||||
.application
|
||||
.delete_repository(parse_uuid(&request.into_inner().repository_id)?.into())
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(DeleteRepositoryResponse {}))
|
||||
}
|
||||
use syncode_identity_model::{RepositoryName, RepositoryOwner};
|
||||
use tonic::{Request, Response, Status};
|
||||
|
||||
use super::IdentityServer;
|
||||
use super::conversion::{
|
||||
application_error, model_repository_owner, model_repository_visibility, parse_uuid,
|
||||
};
|
||||
use crate::wire::{
|
||||
ArchiveRepositoryProjectionRequest, ArchiveRepositoryProjectionResponse,
|
||||
ArchiveRepositoryRequest, ArchiveRepositoryResponse, DeleteRepositoryRequest,
|
||||
DeleteRepositoryResponse, ProvisionRepositoryRequest, ProvisionRepositoryResponse,
|
||||
RegisterRepositoryRequest, RegisterRepositoryResponse, ResolveRepositoryRequest,
|
||||
ResolveRepositoryResponse, UnarchiveRepositoryRequest, UnarchiveRepositoryResponse,
|
||||
UpdateRepositoryMetadataRequest, UpdateRepositoryMetadataResponse,
|
||||
UpdateRepositoryVisibilityRequest, UpdateRepositoryVisibilityResponse,
|
||||
};
|
||||
|
||||
pub(super) async fn resolve(
|
||||
server: &IdentityServer,
|
||||
request: Request<ResolveRepositoryRequest>,
|
||||
) -> Result<Response<ResolveRepositoryResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let owner = request
|
||||
.owner
|
||||
.parse::<RepositoryOwner>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?;
|
||||
let name = request
|
||||
.name
|
||||
.parse::<RepositoryName>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?;
|
||||
let repository_id = server
|
||||
.application
|
||||
.resolve_repository(owner, name)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ResolveRepositoryResponse {
|
||||
repository_id: repository_id.to_string(),
|
||||
}))
|
||||
}
|
||||
|
||||
pub(super) async fn provision(
|
||||
server: &IdentityServer,
|
||||
request: Request<ProvisionRepositoryRequest>,
|
||||
) -> Result<Response<ProvisionRepositoryResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
let owner = model_repository_owner(request.owner_kind(), &request.owner_id)?;
|
||||
let visibility = model_repository_visibility(request.visibility())?;
|
||||
let name = request
|
||||
.name
|
||||
.parse::<RepositoryName>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?;
|
||||
let repository_id = server
|
||||
.application
|
||||
.provision_repository(
|
||||
owner,
|
||||
name,
|
||||
visibility,
|
||||
parse_uuid(&request.creator_id)?.into(),
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ProvisionRepositoryResponse {
|
||||
repository_id: repository_id.to_string(),
|
||||
}))
|
||||
}
|
||||
|
||||
pub(super) async fn register(
|
||||
server: &IdentityServer,
|
||||
request: Request<RegisterRepositoryRequest>,
|
||||
) -> Result<Response<RegisterRepositoryResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
server
|
||||
.application
|
||||
.register_repository(
|
||||
parse_uuid(&request.repository_id)?.into(),
|
||||
model_repository_owner(request.owner_kind(), &request.owner_id)?,
|
||||
request
|
||||
.name
|
||||
.parse::<RepositoryName>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?,
|
||||
model_repository_visibility(request.visibility())?,
|
||||
parse_uuid(&request.creator_id)?.into(),
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(RegisterRepositoryResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn update_visibility(
|
||||
server: &IdentityServer,
|
||||
request: Request<UpdateRepositoryVisibilityRequest>,
|
||||
) -> Result<Response<UpdateRepositoryVisibilityResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
server
|
||||
.application
|
||||
.update_repository_visibility(
|
||||
parse_uuid(&request.repository_id)?.into(),
|
||||
model_repository_visibility(request.visibility())?,
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(UpdateRepositoryVisibilityResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn update_metadata(
|
||||
server: &IdentityServer,
|
||||
request: Request<UpdateRepositoryMetadataRequest>,
|
||||
) -> Result<Response<UpdateRepositoryMetadataResponse>, Status> {
|
||||
let request = request.into_inner();
|
||||
server
|
||||
.application
|
||||
.update_repository_metadata(
|
||||
parse_uuid(&request.repository_id)?.into(),
|
||||
model_repository_owner(request.owner_kind(), &request.owner_id)?,
|
||||
request
|
||||
.name
|
||||
.parse::<RepositoryName>()
|
||||
.map_err(|error| Status::invalid_argument(error.to_string()))?,
|
||||
model_repository_visibility(request.visibility())?,
|
||||
parse_uuid(&request.actor_id)?.into(),
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(UpdateRepositoryMetadataResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn archive_repository(
|
||||
server: &IdentityServer,
|
||||
request: Request<ArchiveRepositoryRequest>,
|
||||
) -> Result<Response<ArchiveRepositoryResponse>, Status> {
|
||||
server
|
||||
.application
|
||||
.archive_repository(parse_uuid(&request.into_inner().repository_id)?.into())
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ArchiveRepositoryResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn archive(
|
||||
server: &IdentityServer,
|
||||
request: Request<ArchiveRepositoryProjectionRequest>,
|
||||
) -> Result<Response<ArchiveRepositoryProjectionResponse>, Status> {
|
||||
server
|
||||
.application
|
||||
.set_repository_archived(
|
||||
parse_uuid(&request.into_inner().repository_id)?.into(),
|
||||
true,
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(ArchiveRepositoryProjectionResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn unarchive(
|
||||
server: &IdentityServer,
|
||||
request: Request<UnarchiveRepositoryRequest>,
|
||||
) -> Result<Response<UnarchiveRepositoryResponse>, Status> {
|
||||
server
|
||||
.application
|
||||
.set_repository_archived(
|
||||
parse_uuid(&request.into_inner().repository_id)?.into(),
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(UnarchiveRepositoryResponse {}))
|
||||
}
|
||||
|
||||
pub(super) async fn delete(
|
||||
server: &IdentityServer,
|
||||
request: Request<DeleteRepositoryRequest>,
|
||||
) -> Result<Response<DeleteRepositoryResponse>, Status> {
|
||||
server
|
||||
.application
|
||||
.delete_repository(parse_uuid(&request.into_inner().repository_id)?.into())
|
||||
.await
|
||||
.map_err(application_error)?;
|
||||
Ok(Response::new(DeleteRepositoryResponse {}))
|
||||
}
|
||||
@@ -1,263 +1,177 @@
|
||||
use syncode_identity_application::ApplicationError;
|
||||
use syncode_identity_model::{
|
||||
RepositoryId, RepositoryName, RepositoryOwner, RepositoryOwnerId, RepositoryOwnerKind,
|
||||
RepositoryVisibility,
|
||||
};
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::repository_grants::{
|
||||
ensure_repository_admin_grants, owner_admin_principals, revoke_previous_owner_grants,
|
||||
};
|
||||
use crate::{Postgres, application_persistence as persistence};
|
||||
|
||||
impl Postgres {
|
||||
pub(super) async fn bridge_resolve_repository(
|
||||
&self,
|
||||
owner: &RepositoryOwner,
|
||||
name: &RepositoryName,
|
||||
) -> Result<Option<RepositoryId>, ApplicationError> {
|
||||
Postgres::resolve_repository(self, owner.as_str(), name.as_str())
|
||||
.await
|
||||
.map(|repository| repository.map(RepositoryId::from))
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_repository_visibility(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
) -> Result<Option<RepositoryVisibility>, ApplicationError> {
|
||||
sqlx::query_scalar!(
|
||||
r#"SELECT visibility::text AS "visibility!" FROM repository
|
||||
WHERE id = $1 AND status = 'active' AND deleted_at IS NULL"#,
|
||||
repository_id
|
||||
)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(persistence)?
|
||||
.map(|visibility| match visibility.as_str() {
|
||||
"public" => Ok(RepositoryVisibility::Public),
|
||||
"private" => Ok(RepositoryVisibility::Private),
|
||||
"limited" => Ok(RepositoryVisibility::Limited),
|
||||
other => Err(ApplicationError::Persistence(format!(
|
||||
"unknown repository visibility {other}"
|
||||
))),
|
||||
})
|
||||
.transpose()
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_repository_owner(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
) -> Result<Option<RepositoryOwnerId>, ApplicationError> {
|
||||
sqlx::query!(
|
||||
r#"SELECT owner_kind::text AS "owner_kind!", owner_id FROM repository
|
||||
WHERE id = $1 AND status = 'active' AND deleted_at IS NULL"#,
|
||||
repository_id
|
||||
)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(persistence)?
|
||||
.map(|owner| {
|
||||
let owner_id = match owner.owner_kind.as_str() {
|
||||
"user" => RepositoryOwnerId::User(owner.owner_id.into()),
|
||||
"organization" => RepositoryOwnerId::Organization(owner.owner_id.into()),
|
||||
other => {
|
||||
return Err(ApplicationError::Persistence(format!(
|
||||
"unknown repository owner kind {other}"
|
||||
)));
|
||||
}
|
||||
};
|
||||
Ok(owner_id)
|
||||
})
|
||||
.transpose()
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_provision_repository(
|
||||
&self,
|
||||
owner_kind: RepositoryOwnerKind,
|
||||
owner_id: Uuid,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
creator_id: Uuid,
|
||||
) -> Result<Uuid, ApplicationError> {
|
||||
let mut transaction = self.pool.begin().await.map_err(persistence)?;
|
||||
let repository_id = Uuid::new_v4();
|
||||
let repository_id = sqlx::query_scalar!(
|
||||
"INSERT INTO repository (id, owner_kind, owner_id, name, visibility)
|
||||
VALUES ($1, $2::text::owner_kind, $3, $4, $5::text::repository_visibility)
|
||||
ON CONFLICT (owner_kind, owner_id, name) WHERE deleted_at IS NULL DO UPDATE
|
||||
SET visibility = EXCLUDED.visibility
|
||||
RETURNING id",
|
||||
repository_id,
|
||||
owner_kind.as_str(),
|
||||
owner_id,
|
||||
name.as_str(),
|
||||
visibility.as_str()
|
||||
)
|
||||
.fetch_one(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
let principals = owner_admin_principals(&mut transaction, owner_kind, owner_id).await?;
|
||||
ensure_repository_admin_grants(&mut transaction, repository_id, &principals, creator_id)
|
||||
.await?;
|
||||
transaction.commit().await.map_err(persistence)?;
|
||||
Ok(repository_id)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_update_repository_visibility(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
visibility: RepositoryVisibility,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
sqlx::query_scalar!(
|
||||
"UPDATE repository SET visibility = $2::text::repository_visibility
|
||||
WHERE id = $1 AND status = 'active' AND deleted_at IS NULL
|
||||
RETURNING id",
|
||||
repository_id,
|
||||
visibility.as_str()
|
||||
)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map(|updated| updated.is_some())
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_update_repository_metadata(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
owner_kind: RepositoryOwnerKind,
|
||||
owner_id: Uuid,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
actor_id: Uuid,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
let mut transaction = self.pool.begin().await.map_err(persistence)?;
|
||||
let current = sqlx::query!(
|
||||
r#"SELECT owner_kind::text AS "owner_kind!", owner_id, name FROM repository
|
||||
WHERE id = $1 AND status = 'active' AND deleted_at IS NULL
|
||||
FOR UPDATE"#,
|
||||
repository_id
|
||||
)
|
||||
.fetch_optional(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
let Some(current) = current else {
|
||||
return Ok(false);
|
||||
};
|
||||
let principals = owner_admin_principals(&mut transaction, owner_kind, owner_id).await?;
|
||||
if current.owner_kind != owner_kind.as_str() || current.owner_id != owner_id {
|
||||
revoke_previous_owner_grants(
|
||||
&mut transaction,
|
||||
repository_id,
|
||||
¤t.owner_kind,
|
||||
current.owner_id,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
sqlx::query!(
|
||||
"UPDATE repository SET owner_kind = $2::text::owner_kind, owner_id = $3,
|
||||
name = $4, visibility = $5::text::repository_visibility
|
||||
WHERE id = $1",
|
||||
repository_id,
|
||||
owner_kind.as_str(),
|
||||
owner_id,
|
||||
name.as_str(),
|
||||
visibility.as_str()
|
||||
)
|
||||
.execute(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
ensure_repository_admin_grants(&mut transaction, repository_id, &principals, actor_id)
|
||||
.await?;
|
||||
// Coordinates are identity's fact and nobody else's. Without these two
|
||||
// events a rename or a transfer is invisible to every other plane.
|
||||
if current.name != name.as_str() {
|
||||
crate::events::record_repository_event(
|
||||
&mut transaction,
|
||||
repository_id,
|
||||
crate::events::REPOSITORY_RENAMED,
|
||||
serde_json::json!({
|
||||
"repository_id": repository_id,
|
||||
"previous_name": current.name,
|
||||
"name": name.as_str(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
}
|
||||
if current.owner_kind != owner_kind.as_str() || current.owner_id != owner_id {
|
||||
crate::events::record_repository_event(
|
||||
&mut transaction,
|
||||
repository_id,
|
||||
crate::events::REPOSITORY_TRANSFERRED,
|
||||
serde_json::json!({
|
||||
"repository_id": repository_id,
|
||||
"previous_owner_kind": current.owner_kind,
|
||||
"previous_owner_id": current.owner_id,
|
||||
"owner_kind": owner_kind.as_str(),
|
||||
"owner_id": owner_id,
|
||||
"name": name.as_str(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
}
|
||||
transaction.commit().await.map_err(persistence)?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_archive_repository(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
let mut transaction = self.pool.begin().await.map_err(persistence)?;
|
||||
let exists = sqlx::query_scalar!(
|
||||
"SELECT id FROM repository WHERE id = $1 FOR UPDATE",
|
||||
repository_id
|
||||
)
|
||||
.fetch_optional(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?
|
||||
.is_some();
|
||||
if !exists {
|
||||
return Ok(false);
|
||||
}
|
||||
sqlx::query!(
|
||||
"UPDATE repository SET status = 'deleted', deleted_at = COALESCE(deleted_at, now())
|
||||
WHERE id = $1",
|
||||
repository_id
|
||||
)
|
||||
.execute(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
sqlx::query!(
|
||||
r#"UPDATE "grant" SET revoked_at = now()
|
||||
WHERE resource_kind = 'repository' AND resource_id = $1 AND revoked_at IS NULL"#,
|
||||
repository_id
|
||||
)
|
||||
.execute(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
transaction.commit().await.map_err(persistence)?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_set_repository_archived(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
archived: bool,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
let status = if archived { "archived" } else { "active" };
|
||||
sqlx::query_scalar!(
|
||||
"UPDATE repository SET status = $2::text::repository_status
|
||||
WHERE id = $1 AND status IN ('active', 'archived') AND deleted_at IS NULL
|
||||
RETURNING id",
|
||||
repository_id,
|
||||
status
|
||||
)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map(|updated| updated.is_some())
|
||||
.map_err(persistence)
|
||||
}
|
||||
}
|
||||
use syncode_identity_application::ApplicationError;
|
||||
use syncode_identity_model::{
|
||||
RepositoryId, RepositoryName, RepositoryOwner, RepositoryOwnerId, RepositoryOwnerKind,
|
||||
RepositoryVisibility,
|
||||
};
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::repository_grants::{ensure_repository_admin_grants, owner_admin_principals};
|
||||
use crate::{Postgres, application_persistence as persistence};
|
||||
|
||||
impl Postgres {
|
||||
pub(super) async fn bridge_resolve_repository(
|
||||
&self,
|
||||
owner: &RepositoryOwner,
|
||||
name: &RepositoryName,
|
||||
) -> Result<Option<RepositoryId>, ApplicationError> {
|
||||
Postgres::resolve_repository(self, owner.as_str(), name.as_str())
|
||||
.await
|
||||
.map(|repository| repository.map(RepositoryId::from))
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_repository_visibility(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
) -> Result<Option<RepositoryVisibility>, ApplicationError> {
|
||||
sqlx::query_scalar!(
|
||||
r#"SELECT visibility::text AS "visibility!" FROM repository
|
||||
WHERE id = $1 AND status = 'active' AND deleted_at IS NULL"#,
|
||||
repository_id
|
||||
)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(persistence)?
|
||||
.map(|visibility| match visibility.as_str() {
|
||||
"public" => Ok(RepositoryVisibility::Public),
|
||||
"private" => Ok(RepositoryVisibility::Private),
|
||||
"limited" => Ok(RepositoryVisibility::Limited),
|
||||
other => Err(ApplicationError::Persistence(format!(
|
||||
"unknown repository visibility {other}"
|
||||
))),
|
||||
})
|
||||
.transpose()
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_repository_owner(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
) -> Result<Option<RepositoryOwnerId>, ApplicationError> {
|
||||
sqlx::query!(
|
||||
r#"SELECT owner_kind::text AS "owner_kind!", owner_id FROM repository
|
||||
WHERE id = $1 AND status = 'active' AND deleted_at IS NULL"#,
|
||||
repository_id
|
||||
)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map_err(persistence)?
|
||||
.map(|owner| {
|
||||
let owner_id = match owner.owner_kind.as_str() {
|
||||
"user" => RepositoryOwnerId::User(owner.owner_id.into()),
|
||||
"organization" => RepositoryOwnerId::Organization(owner.owner_id.into()),
|
||||
other => {
|
||||
return Err(ApplicationError::Persistence(format!(
|
||||
"unknown repository owner kind {other}"
|
||||
)));
|
||||
}
|
||||
};
|
||||
Ok(owner_id)
|
||||
})
|
||||
.transpose()
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_provision_repository(
|
||||
&self,
|
||||
owner_kind: RepositoryOwnerKind,
|
||||
owner_id: Uuid,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
creator_id: Uuid,
|
||||
) -> Result<Uuid, ApplicationError> {
|
||||
let mut transaction = self.pool.begin().await.map_err(persistence)?;
|
||||
let repository_id = Uuid::new_v4();
|
||||
let repository_id = sqlx::query_scalar!(
|
||||
"INSERT INTO repository (id, owner_kind, owner_id, name, visibility)
|
||||
VALUES ($1, $2::text::owner_kind, $3, $4, $5::text::repository_visibility)
|
||||
ON CONFLICT (owner_kind, owner_id, name) WHERE deleted_at IS NULL DO UPDATE
|
||||
SET visibility = EXCLUDED.visibility
|
||||
RETURNING id",
|
||||
repository_id,
|
||||
owner_kind.as_str(),
|
||||
owner_id,
|
||||
name.as_str(),
|
||||
visibility.as_str()
|
||||
)
|
||||
.fetch_one(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
let principals = owner_admin_principals(&mut transaction, owner_kind, owner_id).await?;
|
||||
ensure_repository_admin_grants(&mut transaction, repository_id, &principals, creator_id)
|
||||
.await?;
|
||||
transaction.commit().await.map_err(persistence)?;
|
||||
Ok(repository_id)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_update_repository_visibility(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
visibility: RepositoryVisibility,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
sqlx::query_scalar!(
|
||||
"UPDATE repository SET visibility = $2::text::repository_visibility
|
||||
WHERE id = $1 AND status = 'active' AND deleted_at IS NULL
|
||||
RETURNING id",
|
||||
repository_id,
|
||||
visibility.as_str()
|
||||
)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map(|updated| updated.is_some())
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_archive_repository(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
let mut transaction = self.pool.begin().await.map_err(persistence)?;
|
||||
let exists = sqlx::query_scalar!(
|
||||
"SELECT id FROM repository WHERE id = $1 FOR UPDATE",
|
||||
repository_id
|
||||
)
|
||||
.fetch_optional(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?
|
||||
.is_some();
|
||||
if !exists {
|
||||
return Ok(false);
|
||||
}
|
||||
sqlx::query!(
|
||||
"UPDATE repository SET status = 'deleted', deleted_at = COALESCE(deleted_at, now())
|
||||
WHERE id = $1",
|
||||
repository_id
|
||||
)
|
||||
.execute(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
sqlx::query!(
|
||||
r#"UPDATE "grant" SET revoked_at = now()
|
||||
WHERE resource_kind = 'repository' AND resource_id = $1 AND revoked_at IS NULL"#,
|
||||
repository_id
|
||||
)
|
||||
.execute(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
transaction.commit().await.map_err(persistence)?;
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_set_repository_archived(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
archived: bool,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
let status = if archived { "archived" } else { "active" };
|
||||
sqlx::query_scalar!(
|
||||
"UPDATE repository SET status = $2::text::repository_status
|
||||
WHERE id = $1 AND status IN ('active', 'archived') AND deleted_at IS NULL
|
||||
RETURNING id",
|
||||
repository_id,
|
||||
status
|
||||
)
|
||||
.fetch_optional(&self.pool)
|
||||
.await
|
||||
.map(|updated| updated.is_some())
|
||||
.map_err(persistence)
|
||||
}
|
||||
}
|
||||
@@ -1,0 +1,49 @@
|
||||
use syncode_identity_application::{
|
||||
ActiveBranchProtectionRule, ApplicationError, LocalAgentIdentity, SigningKeyIdentity,
|
||||
};
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::token::{active_branch_protection_rule, signing_key};
|
||||
use crate::{Postgres, application_persistence as persistence};
|
||||
|
||||
impl Postgres {
|
||||
pub(super) async fn bridge_local_agent(
|
||||
&self,
|
||||
agent_id: Uuid,
|
||||
) -> Result<Option<LocalAgentIdentity>, ApplicationError> {
|
||||
self.load_local_agent(agent_id)
|
||||
.await
|
||||
.map(|agent| {
|
||||
agent.map(|agent| LocalAgentIdentity {
|
||||
owner_user_id: agent.owner_user_id.into(),
|
||||
restriction: agent.restriction,
|
||||
})
|
||||
})
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_active_branch_protection_rules(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
) -> Result<Vec<ActiveBranchProtectionRule>, ApplicationError> {
|
||||
self.find_branch_protection_rules(repository_id)
|
||||
.await
|
||||
.map(|rules| {
|
||||
rules
|
||||
.into_iter()
|
||||
.map(active_branch_protection_rule)
|
||||
.collect()
|
||||
})
|
||||
.map_err(persistence)
|
||||
}
|
||||
|
||||
pub(super) async fn bridge_user_signing_keys(
|
||||
&self,
|
||||
user_id: Uuid,
|
||||
) -> Result<Vec<SigningKeyIdentity>, ApplicationError> {
|
||||
self.find_user_signing_keys(user_id)
|
||||
.await
|
||||
.map(|keys| keys.into_iter().map(signing_key).collect())
|
||||
.map_err(persistence)
|
||||
}
|
||||
}
|
||||
@@ -1,0 +1,92 @@
|
||||
use syncode_identity_application::ApplicationError;
|
||||
use syncode_identity_model::{RepositoryName, RepositoryOwnerKind, RepositoryVisibility};
|
||||
use uuid::Uuid;
|
||||
|
||||
use super::repository_grants::{
|
||||
ensure_repository_admin_grants, owner_admin_principals, revoke_previous_owner_grants,
|
||||
};
|
||||
use crate::{Postgres, application_persistence as persistence};
|
||||
|
||||
impl Postgres {
|
||||
pub(super) async fn bridge_update_repository_metadata(
|
||||
&self,
|
||||
repository_id: Uuid,
|
||||
owner_kind: RepositoryOwnerKind,
|
||||
owner_id: Uuid,
|
||||
name: &RepositoryName,
|
||||
visibility: RepositoryVisibility,
|
||||
actor_id: Uuid,
|
||||
) -> Result<bool, ApplicationError> {
|
||||
let mut transaction = self.pool.begin().await.map_err(persistence)?;
|
||||
let current = sqlx::query!(
|
||||
r#"SELECT owner_kind::text AS "owner_kind!", owner_id, name FROM repository
|
||||
WHERE id = $1 AND status = 'active' AND deleted_at IS NULL
|
||||
FOR UPDATE"#,
|
||||
repository_id
|
||||
)
|
||||
.fetch_optional(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
let Some(current) = current else {
|
||||
return Ok(false);
|
||||
};
|
||||
let principals = owner_admin_principals(&mut transaction, owner_kind, owner_id).await?;
|
||||
if current.owner_kind != owner_kind.as_str() || current.owner_id != owner_id {
|
||||
revoke_previous_owner_grants(
|
||||
&mut transaction,
|
||||
repository_id,
|
||||
¤t.owner_kind,
|
||||
current.owner_id,
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
sqlx::query!(
|
||||
"UPDATE repository SET owner_kind = $2::text::owner_kind, owner_id = $3,
|
||||
name = $4, visibility = $5::text::repository_visibility
|
||||
WHERE id = $1",
|
||||
repository_id,
|
||||
owner_kind.as_str(),
|
||||
owner_id,
|
||||
name.as_str(),
|
||||
visibility.as_str()
|
||||
)
|
||||
.execute(&mut *transaction)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
ensure_repository_admin_grants(&mut transaction, repository_id, &principals, actor_id)
|
||||
.await?;
|
||||
if current.name != name.as_str() {
|
||||
crate::events::record_repository_event(
|
||||
&mut transaction,
|
||||
repository_id,
|
||||
crate::events::REPOSITORY_RENAMED,
|
||||
serde_json::json!({
|
||||
"repository_id": repository_id,
|
||||
"previous_name": current.name,
|
||||
"name": name.as_str(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
}
|
||||
if current.owner_kind != owner_kind.as_str() || current.owner_id != owner_id {
|
||||
crate::events::record_repository_event(
|
||||
&mut transaction,
|
||||
repository_id,
|
||||
crate::events::REPOSITORY_TRANSFERRED,
|
||||
serde_json::json!({
|
||||
"repository_id": repository_id,
|
||||
"previous_owner_kind": current.owner_kind,
|
||||
"previous_owner_id": current.owner_id,
|
||||
"owner_kind": owner_kind.as_str(),
|
||||
"owner_id": owner_id,
|
||||
"name": name.as_str(),
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.map_err(persistence)?;
|
||||
}
|
||||
transaction.commit().await.map_err(persistence)?;
|
||||
Ok(true)
|
||||
}
|
||||
}
|
||||
@@ -1,0 +1,150 @@
|
||||
use std::time::Duration;
|
||||
|
||||
use hmac::{Hmac, KeyInit, Mac};
|
||||
use serde::Serialize;
|
||||
use sha2::Sha256;
|
||||
use syncode_identity_storage::{PendingEvent, Postgres};
|
||||
use thiserror::Error;
|
||||
use uuid::Uuid;
|
||||
|
||||
pub struct EventPublisher {
|
||||
store: Postgres,
|
||||
endpoint: reqwest::Url,
|
||||
secret: String,
|
||||
source_node_id: Uuid,
|
||||
client: reqwest::Client,
|
||||
interval: Duration,
|
||||
batch: i64,
|
||||
}
|
||||
|
||||
pub struct EventPublisherConfig<'a> {
|
||||
pub endpoint: &'a str,
|
||||
pub secret: String,
|
||||
pub source_node_id: Uuid,
|
||||
pub interval: Duration,
|
||||
pub batch: i64,
|
||||
}
|
||||
|
||||
#[derive(Debug, Error)]
|
||||
pub enum EventPublisherError {
|
||||
#[error("identity event interval, batch, and signing secret must be non-empty")]
|
||||
InvalidConfiguration,
|
||||
#[error("identity event endpoint is invalid: {0}")]
|
||||
InvalidEndpoint(#[from] url::ParseError),
|
||||
#[error("identity event HTTP client cannot be created: {0}")]
|
||||
HttpClient(#[from] reqwest::Error),
|
||||
#[error("identity event outbox failed: {0}")]
|
||||
Store(#[from] syncode_identity_storage::StoreError),
|
||||
}
|
||||
|
||||
#[derive(Serialize)]
|
||||
struct Envelope<'a> {
|
||||
message_id: Uuid,
|
||||
protocol_version: &'static str,
|
||||
schema_version: i32,
|
||||
source_node_id: Uuid,
|
||||
repository_id: Uuid,
|
||||
message_type: &'a str,
|
||||
sequence: i64,
|
||||
term: i64,
|
||||
payload: &'a serde_json::Value,
|
||||
}
|
||||
|
||||
impl EventPublisher {
|
||||
pub fn new(
|
||||
store: Postgres,
|
||||
config: EventPublisherConfig<'_>,
|
||||
) -> Result<Self, EventPublisherError> {
|
||||
if config.interval.is_zero() || config.batch <= 0 || config.secret.is_empty() {
|
||||
return Err(EventPublisherError::InvalidConfiguration);
|
||||
}
|
||||
Ok(Self {
|
||||
store,
|
||||
endpoint: config.endpoint.parse()?,
|
||||
secret: config.secret,
|
||||
source_node_id: config.source_node_id,
|
||||
client: reqwest::Client::builder()
|
||||
.timeout(Duration::from_secs(10))
|
||||
.build()?,
|
||||
interval: config.interval,
|
||||
batch: config.batch,
|
||||
})
|
||||
}
|
||||
|
||||
pub async fn run(self) -> Result<(), EventPublisherError> {
|
||||
let mut ticker = tokio::time::interval(self.interval);
|
||||
loop {
|
||||
ticker.tick().await;
|
||||
self.publish_once().await?;
|
||||
}
|
||||
}
|
||||
|
||||
pub async fn publish_once(&self) -> Result<(), EventPublisherError> {
|
||||
let lease = i64::try_from(self.interval.as_secs().max(30)).unwrap_or(i64::MAX);
|
||||
let pending = self.store.claim_identity_events(self.batch, lease).await?;
|
||||
for event in pending {
|
||||
match self.publish(&event).await {
|
||||
Ok(()) => {
|
||||
self.store
|
||||
.mark_identity_event_delivered(event.message_id)
|
||||
.await?;
|
||||
}
|
||||
Err(error) => {
|
||||
self.store
|
||||
.mark_identity_event_failed(
|
||||
event.message_id,
|
||||
&error,
|
||||
retry_delay(event.attempts),
|
||||
)
|
||||
.await?;
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn publish(&self, event: &PendingEvent) -> Result<(), String> {
|
||||
let body = serde_json::to_vec(&Envelope {
|
||||
message_id: event.message_id,
|
||||
protocol_version: "1.0",
|
||||
schema_version: event.schema_version,
|
||||
source_node_id: self.source_node_id,
|
||||
repository_id: event.resource_id,
|
||||
message_type: &event.event_type,
|
||||
sequence: event.sequence,
|
||||
term: 1,
|
||||
payload: &event.payload,
|
||||
})
|
||||
.map_err(|error| error.to_string())?;
|
||||
let mut mac = Hmac::<Sha256>::new_from_slice(self.secret.as_bytes())
|
||||
.map_err(|error| error.to_string())?;
|
||||
mac.update(&body);
|
||||
let signature = hex::encode(mac.finalize().into_bytes());
|
||||
let response = self
|
||||
.client
|
||||
.post(self.endpoint.clone())
|
||||
.header("x-syncode-event", &event.event_type)
|
||||
.header("x-syncode-delivery", event.message_id.to_string())
|
||||
.header("x-syncode-signature", signature)
|
||||
.header("content-type", "application/json")
|
||||
.body(body)
|
||||
.send()
|
||||
.await
|
||||
.map_err(|error| error.to_string())?;
|
||||
if response.status().is_success() {
|
||||
Ok(())
|
||||
} else {
|
||||
Err(format!(
|
||||
"identity event endpoint returned HTTP {}",
|
||||
response.status().as_u16()
|
||||
))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn retry_delay(attempts: i32) -> i64 {
|
||||
let exponent = u32::try_from(attempts.saturating_sub(1))
|
||||
.unwrap_or(u32::MAX)
|
||||
.min(8);
|
||||
1_i64.checked_shl(exponent).unwrap_or(300).min(300)
|
||||
}
|
||||
@@ -1,0 +1,1 @@
|
||||
pub mod event_outbox;
|
||||
@@ -1,0 +1,100 @@
|
||||
#![allow(clippy::expect_used, clippy::panic, clippy::unwrap_used)]
|
||||
|
||||
use std::error::Error;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
use axum::body::Bytes;
|
||||
use axum::extract::State;
|
||||
use axum::http::{HeaderMap, StatusCode};
|
||||
use axum::routing::post;
|
||||
use hmac::{Hmac, KeyInit, Mac};
|
||||
use sha2::Sha256;
|
||||
use syncode_identity::event_outbox::{EventPublisher, EventPublisherConfig};
|
||||
use syncode_identity_storage::Postgres;
|
||||
use uuid::Uuid;
|
||||
|
||||
type Recorded = Arc<Mutex<Option<(HeaderMap, Vec<u8>)>>>;
|
||||
type TestResult<T = ()> = Result<T, Box<dyn Error + Send + Sync>>;
|
||||
|
||||
async fn record(State(recorded): State<Recorded>, headers: HeaderMap, body: Bytes) -> StatusCode {
|
||||
recorded
|
||||
.lock()
|
||||
.map(|mut value| *value = Some((headers, body.to_vec())))
|
||||
.map_or(StatusCode::INTERNAL_SERVER_ERROR, |()| {
|
||||
StatusCode::NO_CONTENT
|
||||
})
|
||||
}
|
||||
|
||||
fn database_url() -> String {
|
||||
std::env::var("SYNCODE_IDENTITY_TEST_DATABASE_URL")
|
||||
.unwrap_or_else(|_| "postgres:///syncode_identity_test".to_owned())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn identity_event_is_delivered_and_acknowledged() -> TestResult {
|
||||
let store = Postgres::connect(&database_url(), 2).await?;
|
||||
let message_id = Uuid::new_v4();
|
||||
let repository_id = Uuid::new_v4();
|
||||
sqlx::query(
|
||||
"INSERT INTO identity_event \
|
||||
(message_id, resource_kind, resource_id, event_type, payload) \
|
||||
VALUES ($1, 'repository', $2, 'identity.repository.renamed', $3)",
|
||||
)
|
||||
.bind(message_id)
|
||||
.bind(repository_id)
|
||||
.bind(serde_json::json!({"name": "renamed"}))
|
||||
.execute(store.pool())
|
||||
.await?;
|
||||
sqlx::query(
|
||||
"INSERT INTO identity_outbox (message_id, topic) \
|
||||
VALUES ($1, 'identity.repository.renamed')",
|
||||
)
|
||||
.bind(message_id)
|
||||
.execute(store.pool())
|
||||
.await?;
|
||||
|
||||
let recorded = Recorded::default();
|
||||
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
|
||||
let endpoint = format!("http://{}/events", listener.local_addr()?);
|
||||
let app = axum::Router::new()
|
||||
.route("/events", post(record))
|
||||
.with_state(Arc::clone(&recorded));
|
||||
let server = tokio::spawn(async move { axum::serve(listener, app).await });
|
||||
|
||||
EventPublisher::new(
|
||||
store.clone(),
|
||||
EventPublisherConfig {
|
||||
endpoint: &endpoint,
|
||||
secret: "identity-secret".to_owned(),
|
||||
source_node_id: Uuid::new_v4(),
|
||||
interval: std::time::Duration::from_secs(1),
|
||||
batch: 10,
|
||||
},
|
||||
)?
|
||||
.publish_once()
|
||||
.await?;
|
||||
|
||||
let (headers, body) = recorded
|
||||
.lock()
|
||||
.map_err(|_| "recorded event lock was poisoned")?
|
||||
.clone()
|
||||
.ok_or("identity event was not delivered")?;
|
||||
let envelope: serde_json::Value = serde_json::from_slice(&body)?;
|
||||
assert_eq!(headers["x-syncode-event"], "identity.repository.renamed");
|
||||
let signature = hex::decode(headers["x-syncode-signature"].to_str()?)?;
|
||||
let mut mac = Hmac::<Sha256>::new_from_slice(b"identity-secret")?;
|
||||
mac.update(&body);
|
||||
mac.verify_slice(&signature)?;
|
||||
assert_eq!(envelope["protocol_version"], "1.0");
|
||||
assert_eq!(envelope["repository_id"], repository_id.to_string());
|
||||
assert_eq!(envelope["message_type"], "identity.repository.renamed");
|
||||
let delivered = sqlx::query_scalar::<_, bool>(
|
||||
"SELECT delivered_at IS NOT NULL FROM identity_outbox WHERE message_id = $1",
|
||||
)
|
||||
.bind(message_id)
|
||||
.fetch_one(store.pool())
|
||||
.await?;
|
||||
assert!(delivered);
|
||||
server.abort();
|
||||
Ok(())
|
||||
}
|
||||
Reference in New Issue
Block a user