use crate::{
error::{binding_env_var, map_cloud_client_error, ErrorData, Result},
traits::{
ArtifactRegistry, ArtifactRegistryCredentials, ArtifactRegistryPermissions, Binding,
ComputeServiceType, CrossAccountAccess, CrossAccountPermissions, GcpCrossAccountAccess,
RegistryAuthMethod, RepositoryResponse,
},
};
use alien_core::bindings::ArtifactRegistryBinding;
use alien_error::{AlienError, Context};
use alien_gcp_clients::iam::IamPolicy;
use alien_gcp_clients::{
artifactregistry::{ArtifactRegistryApi, ArtifactRegistryClient},
GcpClientConfig, GcpClientConfigExt as _,
};
use async_trait::async_trait;
use chrono;
use tracing::{debug, info, warn};
fn agent_domain(service_type: &ComputeServiceType) -> &'static str {
match service_type {
ComputeServiceType::Worker => "serverless-robot-prod",
ComputeServiceType::Sandbox => "gcp-sa-vertex-sandbox",
}
}
fn is_agent_created_on_first_use(service_type: &ComputeServiceType) -> bool {
match service_type {
ComputeServiceType::Sandbox => true,
ComputeServiceType::Worker => false,
}
}
fn agent_member(service_account: &str) -> Option<(ComputeServiceType, &str)> {
ComputeServiceType::ALL.iter().find_map(|service_type| {
let suffix = format!("@{}.iam.gserviceaccount.com", agent_domain(service_type));
service_account
.strip_prefix("service-")
.and_then(|rest| rest.strip_suffix(&suffix))
.map(|project_number| (service_type.clone(), project_number))
})
}
fn cross_account_members(access: &GcpCrossAccountAccess) -> Vec<String> {
let mut members = Vec::new();
for service_type in &access.allowed_service_types {
let agent_domain = agent_domain(service_type);
for project_number in &access.project_numbers {
members.push(format!(
"serviceAccount:service-{project_number}@{agent_domain}.iam.gserviceaccount.com"
));
}
}
for service_account_email in &access.service_account_emails {
members.push(format!("serviceAccount:{service_account_email}"));
}
members
}
#[derive(Debug)]
pub struct GarArtifactRegistry {
client: ArtifactRegistryClient,
binding_name: String,
project_id: String,
location: String,
repository_name: String,
pull_service_account_email: Option<String>,
push_service_account_email: Option<String>,
gcp_config: GcpClientConfig,
}
impl GarArtifactRegistry {
pub async fn new(
binding_name: String,
binding: ArtifactRegistryBinding,
gcp_config: &GcpClientConfig,
) -> Result<Self> {
info!(
binding_name = %binding_name,
"Initializing GCP Artifact Registry"
);
let client = crate::http_client::create_http_client();
let artifact_registry_client = ArtifactRegistryClient::new(client, gcp_config.clone());
let project_id = gcp_config.project_id.clone();
let location = gcp_config.region.clone();
let config = match binding {
ArtifactRegistryBinding::Gar(config) => config,
_ => {
return Err(AlienError::new(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&binding_name),
binding_name: binding_name.clone(),
reason: "Expected GAR binding, got different service type".to_string(),
}));
}
};
let repository_name = config
.repository_name
.into_value(&binding_name, "repository_name")
.context(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&binding_name),
binding_name: binding_name.clone(),
reason: "Failed to extract repository_name from binding".to_string(),
})?;
let pull_service_account_email = config
.pull_service_account_email
.map(|v| {
v.into_value(&binding_name, "pull_service_account_email")
.context(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&binding_name),
binding_name: binding_name.clone(),
reason: "Failed to extract pull_service_account_email from binding"
.to_string(),
})
})
.transpose()?;
let push_service_account_email = config
.push_service_account_email
.map(|v| {
v.into_value(&binding_name, "push_service_account_email")
.context(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&binding_name),
binding_name: binding_name.clone(),
reason: "Failed to extract push_service_account_email from binding"
.to_string(),
})
})
.transpose()?;
Ok(Self {
client: artifact_registry_client,
binding_name,
project_id,
location,
repository_name,
pull_service_account_email,
push_service_account_email,
gcp_config: gcp_config.clone(),
})
}
fn extract_repo_name(&self, repo_id: &str) -> Result<String> {
if repo_id.is_empty() {
return Ok(self.repository_name.clone());
}
if let Some(name) = repo_id.split('/').last() {
Ok(name.to_string())
} else {
Err(AlienError::new(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&self.binding_name),
binding_name: self.binding_name.clone(),
reason: format!("Invalid repository ID format: {}", repo_id),
}))
}
}
async fn update_policy_members(
&self,
repo_name: &str,
mut current_policy: IamPolicy,
members: Vec<String>,
add_members: bool, ) -> Result<IamPolicy> {
let reader_role = "roles/artifactregistry.reader";
let mut binding_index = None;
for (i, binding) in current_policy.bindings.iter().enumerate() {
if binding.role == reader_role {
binding_index = Some(i);
break;
}
}
if add_members {
if members.is_empty() {
info!(repo_name = %repo_name, "No new members to add");
return Ok(current_policy);
}
match binding_index {
Some(i) => {
let binding = &mut current_policy.bindings[i];
for member in members {
if !binding.members.contains(&member) {
binding.members.push(member);
}
}
}
None => {
current_policy
.bindings
.push(alien_gcp_clients::iam::Binding {
role: reader_role.to_string(),
members,
condition: None,
});
}
}
} else {
if let Some(i) = binding_index {
let binding = &mut current_policy.bindings[i];
binding.members.retain(|member| !members.contains(member));
if binding.members.is_empty() {
current_policy.bindings.remove(i);
}
}
}
let updated = self.client.set_repository_iam_policy(
self.project_id.clone(),
self.location.clone(),
repo_name.to_string(),
current_policy,
).await
.map_err(|e| map_cloud_client_error(
e,
format!("Failed to update cross-account access for GCP Artifact Registry repository '{}'", repo_name),
Some(repo_name.to_string()),
))?;
let action = if add_members { "added" } else { "removed" };
info!(
repo_name = %repo_name,
action = %action,
"GCP Artifact Registry repository cross-account access updated successfully"
);
Ok(updated)
}
}
impl Binding for GarArtifactRegistry {}
#[async_trait]
impl ArtifactRegistry for GarArtifactRegistry {
fn registry_endpoint(&self) -> String {
format!("https://{}-docker.pkg.dev", self.location)
}
fn upstream_repository_prefix(&self) -> String {
format!("{}/{}", self.project_id, self.repository_name)
}
async fn create_repository(&self, repo_name: &str) -> Result<RepositoryResponse> {
let routable_name = format!("{}/{}", self.upstream_repository_prefix(), repo_name);
Ok(RepositoryResponse {
name: routable_name,
uri: None,
created_at: None,
})
}
async fn get_repository(&self, repo_id: &str) -> Result<RepositoryResponse> {
let image_path = self.extract_repo_name(repo_id)?;
let routable_name = format!("{}/{}", self.upstream_repository_prefix(), image_path);
let repository_uri = format!(
"{}-docker.pkg.dev/{}/{}",
self.location, self.project_id, image_path
);
Ok(RepositoryResponse {
name: routable_name,
uri: Some(repository_uri),
created_at: None,
})
}
async fn add_cross_account_access(
&self,
repo_id: &str,
access: CrossAccountAccess,
) -> Result<()> {
let _ = repo_id; let repo_name = self.repository_name.clone();
let gcp_access = match access {
CrossAccountAccess::Gcp(gcp_access) => gcp_access,
_ => {
return Err(AlienError::new(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&self.binding_name),
binding_name: self.binding_name.clone(),
reason: "GCP artifact registry can only accept GCP cross-account access configuration".to_string(),
}));
}
};
info!(
repo_name = %repo_name,
project_numbers = ?gcp_access.project_numbers,
allowed_service_types = ?gcp_access.allowed_service_types,
service_account_emails = ?gcp_access.service_account_emails,
"Adding GCP Artifact Registry repository cross-account access"
);
let current_policy = match self
.client
.get_repository_iam_policy(
self.project_id.clone(),
self.location.clone(),
repo_name.clone(),
)
.await
{
Ok(policy) => policy,
Err(error) if error.http_status_code == Some(404) => IamPolicy {
version: Some(1),
kind: None,
resource_id: None,
bindings: vec![],
etag: None,
},
Err(error) => {
return Err(map_cloud_client_error(
error,
format!(
"Failed to read the IAM policy of GCP Artifact Registry repository '{repo_name}' to grant access"
),
Some(repo_name.to_string()),
))
}
};
let (deferred, present): (Vec<ComputeServiceType>, Vec<ComputeServiceType>) = gcp_access
.allowed_service_types
.iter()
.cloned()
.partition(is_agent_created_on_first_use);
let policy = self
.update_policy_members(
&repo_name,
current_policy,
cross_account_members(&GcpCrossAccountAccess {
allowed_service_types: present,
..gcp_access.clone()
}),
true,
)
.await?;
if deferred.is_empty() {
return Ok(());
}
if policy.etag.is_none() {
return Err(AlienError::new(ErrorData::RemoteResourceConflict {
operation_context: "adding cross-account access".to_string(),
resource_type: "repository".to_string(),
resource_name: repo_name.clone(),
conflict_reason:
"the updated IAM policy carries no etag, so the sandbox agent's member cannot be added without replacing the policy"
.to_string(),
}));
}
self.update_policy_members(
&repo_name,
policy,
cross_account_members(&GcpCrossAccountAccess {
allowed_service_types: deferred,
service_account_emails: Vec::new(),
..gcp_access
}),
true,
)
.await?;
Ok(())
}
async fn remove_cross_account_access(
&self,
repo_id: &str,
access: CrossAccountAccess,
) -> Result<()> {
let _ = repo_id;
let repo_name = self.repository_name.clone();
let gcp_access = match access {
CrossAccountAccess::Gcp(gcp_access) => gcp_access,
_ => {
return Err(AlienError::new(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&self.binding_name),
binding_name: self.binding_name.clone(),
reason: "GCP artifact registry can only accept GCP cross-account access configuration".to_string(),
}));
}
};
info!(
repo_name = %repo_name,
project_numbers = ?gcp_access.project_numbers,
allowed_service_types = ?gcp_access.allowed_service_types,
service_account_emails = ?gcp_access.service_account_emails,
"Removing GCP Artifact Registry repository cross-account access"
);
let current_policy = match self
.client
.get_repository_iam_policy(
self.project_id.clone(),
self.location.clone(),
repo_name.clone(),
)
.await
{
Ok(policy) => policy,
Err(error) if error.http_status_code == Some(404) => {
info!(repo_name = %repo_name, "No existing GCP IAM policy to remove permissions from");
return Ok(());
}
Err(error) => {
return Err(map_cloud_client_error(
error,
format!(
"Failed to read the IAM policy of GCP Artifact Registry repository '{repo_name}' to revoke access"
),
Some(repo_name.to_string()),
))
}
};
let members_to_remove = cross_account_members(&gcp_access);
self.update_policy_members(&repo_name, current_policy, members_to_remove, false)
.await?;
Ok(())
}
async fn get_cross_account_access(&self, repo_id: &str) -> Result<CrossAccountPermissions> {
let _ = repo_id;
let repo_name = self.repository_name.clone();
info!(
repo_name = %repo_name,
"Getting GCP Artifact Registry repository cross-account access"
);
let policy = match self
.client
.get_repository_iam_policy(
self.project_id.clone(),
self.location.clone(),
repo_name.clone(),
)
.await
{
Ok(policy) => policy,
Err(e) => {
warn!(
repo_name = %repo_name,
error = %e,
"Failed to get GCP Artifact Registry repository IAM policy"
);
return Ok(CrossAccountPermissions {
access: CrossAccountAccess::Gcp(GcpCrossAccountAccess {
project_numbers: Vec::new(),
allowed_service_types: Vec::new(),
service_account_emails: Vec::new(),
}),
last_updated: None,
});
}
};
let mut project_numbers = Vec::new();
let mut service_account_emails = Vec::new();
let mut allowed_service_types = Vec::new();
for binding in policy.bindings {
if binding.role.contains("reader") || binding.role.contains("artifactregistry") {
for member in binding.members {
if let Some(service_account) = member.strip_prefix("serviceAccount:") {
match agent_member(service_account) {
Some((service_type, project_number)) => {
project_numbers.push(project_number.to_string());
if !allowed_service_types.contains(&service_type) {
allowed_service_types.push(service_type);
}
}
None => service_account_emails.push(service_account.to_string()),
}
}
}
}
}
project_numbers.sort();
project_numbers.dedup();
service_account_emails.sort();
service_account_emails.dedup();
allowed_service_types.sort_by_key(|rt| format!("{:?}", rt));
allowed_service_types.dedup();
info!(
repo_name = %repo_name,
project_numbers = ?project_numbers,
allowed_service_types = ?allowed_service_types,
service_account_emails = ?service_account_emails,
"Retrieved GCP Artifact Registry repository cross-account access"
);
Ok(CrossAccountPermissions {
access: CrossAccountAccess::Gcp(GcpCrossAccountAccess {
project_numbers,
allowed_service_types,
service_account_emails,
}),
last_updated: None, })
}
async fn generate_credentials(
&self,
repo_id: &str,
permissions: ArtifactRegistryPermissions,
ttl_seconds: Option<u32>,
) -> Result<ArtifactRegistryCredentials> {
info!(
repo_id = %repo_id,
permissions = ?permissions,
ttl_seconds = ?ttl_seconds,
"Generating GCP Artifact Registry credentials by impersonating service account"
);
let _project_id = &self.project_id;
let _location = &self.location;
let service_account_email = match permissions {
ArtifactRegistryPermissions::Pull => {
self.pull_service_account_email.clone()
.ok_or_else(|| AlienError::new(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&self.binding_name),
binding_name: self.binding_name.clone(),
reason: "Pull service account email not available - ensure the artifact registry resource is properly linked".to_string(),
}))?
}
ArtifactRegistryPermissions::PushPull => {
self.push_service_account_email.clone()
.ok_or_else(|| AlienError::new(ErrorData::BindingConfigInvalid {
env_var: binding_env_var(&self.binding_name),
binding_name: self.binding_name.clone(),
reason: "Push service account email not available - ensure the artifact registry resource is properly linked".to_string(),
}))?
}
};
info!(
service_account_email = %service_account_email,
"Using stored service account email for GCP Artifact Registry access"
);
let gcp_config = &self.gcp_config;
let scopes = vec![
"https://www.googleapis.com/auth/cloud-platform".to_string(),
"https://www.googleapis.com/auth/devstorage.read_write".to_string(),
];
let lifetime = ttl_seconds.map(|ttl| format!("{}s", ttl.min(3600)));
let impersonation_config = alien_gcp_clients::GcpImpersonationConfig {
service_account_email: service_account_email.clone(),
scopes,
delegates: None,
lifetime,
target_project_id: None,
target_region: None,
};
let impersonated_config =
gcp_config
.impersonate(impersonation_config)
.await
.map_err(|e| {
map_cloud_client_error(
e,
"Failed to impersonate GCP service account for artifact registry access"
.to_string(),
Some(repo_id.to_string()),
)
})?;
let access_token = impersonated_config
.get_bearer_token("https://www.googleapis.com/")
.await
.map_err(|e| {
map_cloud_client_error(
e,
"Failed to get OAuth token from impersonated service account".to_string(),
Some(repo_id.to_string()),
)
})?;
let expires_at = if let Some(ttl) = ttl_seconds {
Some(
(chrono::Utc::now() + chrono::Duration::seconds(ttl.min(3600) as i64)).to_rfc3339(),
)
} else {
Some((chrono::Utc::now() + chrono::Duration::seconds(3600)).to_rfc3339())
};
info!(
permissions = ?permissions,
service_account = %service_account_email,
"GCP Artifact Registry OAuth token generated successfully with impersonated service account"
);
Ok(ArtifactRegistryCredentials {
auth_method: RegistryAuthMethod::Basic,
username: "oauth2accesstoken".to_string(),
password: access_token,
expires_at,
})
}
async fn delete_repository(&self, repo_id: &str) -> Result<()> {
debug!(
repo_id = %repo_id,
"GCP Artifact Registry delete_repository: no-op (image paths are implicit)"
);
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn each_service_type_resolves_to_its_own_service_agent() {
let members = cross_account_members(&GcpCrossAccountAccess {
project_numbers: vec!["123456789012".to_string()],
allowed_service_types: vec![ComputeServiceType::Worker, ComputeServiceType::Sandbox],
service_account_emails: vec![
"management@test-project.iam.gserviceaccount.com".to_string()
],
});
assert_eq!(
members,
vec![
"serviceAccount:service-123456789012@serverless-robot-prod.iam.gserviceaccount.com",
"serviceAccount:service-123456789012@gcp-sa-vertex-sandbox.iam.gserviceaccount.com",
"serviceAccount:management@test-project.iam.gserviceaccount.com",
]
);
}
#[test]
fn every_service_type_survives_a_write_then_read() {
for service_type in [ComputeServiceType::Worker, ComputeServiceType::Sandbox] {
let written = cross_account_members(&GcpCrossAccountAccess {
project_numbers: vec!["123456789012".to_string()],
allowed_service_types: vec![service_type.clone()],
service_account_emails: Vec::new(),
});
let member = written
.first()
.expect("a project and a type name one member");
let account = member
.strip_prefix("serviceAccount:")
.expect("members carry the IAM principal prefix");
assert_eq!(
agent_member(account),
Some((service_type.clone(), "123456789012")),
"{service_type:?} must decode back to the project that earned it"
);
}
}
#[test]
fn the_worker_and_sandbox_members_are_written_apart() {
let access = GcpCrossAccountAccess {
project_numbers: vec!["123456789012".to_string()],
allowed_service_types: vec![ComputeServiceType::Worker, ComputeServiceType::Sandbox],
service_account_emails: vec![
"management@test-project.iam.gserviceaccount.com".to_string()
],
};
let first = cross_account_members(&GcpCrossAccountAccess {
allowed_service_types: vec![ComputeServiceType::Worker],
..access.clone()
});
let second = cross_account_members(&GcpCrossAccountAccess {
allowed_service_types: vec![ComputeServiceType::Sandbox],
service_account_emails: Vec::new(),
..access.clone()
});
assert!(second
.iter()
.all(|member| member.contains("gcp-sa-vertex-sandbox")));
assert!(first
.iter()
.all(|member| !member.contains("gcp-sa-vertex-sandbox")));
let mut both = [first, second].concat();
both.sort();
let mut whole = cross_account_members(&access);
whole.sort();
assert_eq!(both, whole);
}
#[test]
fn a_service_type_grants_nothing_for_a_project_it_does_not_name() {
let members = cross_account_members(&GcpCrossAccountAccess {
project_numbers: Vec::new(),
allowed_service_types: vec![ComputeServiceType::Sandbox],
service_account_emails: vec![
"management@test-project.iam.gserviceaccount.com".to_string()
],
});
assert_eq!(
members,
vec!["serviceAccount:management@test-project.iam.gserviceaccount.com"]
);
}
}