feat: publish identity events #32

Manually merged
day01 merged 1 commits from feat/0.6-event-delivery into develop 2026-08-30 09:27:02 +00:00
13 changed files with 498 additions and 145 deletions
Showing only changes of commit d17d5217cb - Show all commits
+3
View File
File diff suppressed because it is too large Load Diff
+3
View File
@@ -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
View File
@@ -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");
}
}
}
})
}
+5 -32
View File
@@ -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
}
}
+6 -22
View File
@@ -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
View File
@@ -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
}
}
+56 -4
View File
@@ -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 -87
View File
@@ -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,
&current.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)
}
}
+49
View File
@@ -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,
&current.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)
}
}
+150
View File
@@ -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
View File
@@ -1,0 +1,1 @@
pub mod event_outbox;
+100
View File
@@ -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(())
}