1use crate::{
2 error::{binding_env_var, map_cloud_client_error, ErrorData, Result},
3 traits::{
4 ArtifactRegistry, ArtifactRegistryCredentials, ArtifactRegistryPermissions, Binding,
5 ComputeServiceType, CrossAccountAccess, CrossAccountPermissions, GcpCrossAccountAccess,
6 RegistryAuthMethod, RepositoryResponse,
7 },
8};
9use alien_core::bindings::ArtifactRegistryBinding;
10use alien_error::{AlienError, Context};
11use alien_gcp_clients::iam::IamPolicy;
12use alien_gcp_clients::{
13 artifactregistry::{ArtifactRegistryApi, ArtifactRegistryClient},
14 GcpClientConfig, GcpClientConfigExt as _,
15};
16use async_trait::async_trait;
17use chrono;
18use tracing::{debug, info, warn};
19
20fn agent_domain(service_type: &ComputeServiceType) -> &'static str {
25 match service_type {
26 ComputeServiceType::Worker => "serverless-robot-prod",
27 ComputeServiceType::Sandbox => "gcp-sa-vertex-sandbox",
28 }
29}
30
31fn is_agent_created_on_first_use(service_type: &ComputeServiceType) -> bool {
34 match service_type {
35 ComputeServiceType::Sandbox => true,
36 ComputeServiceType::Worker => false,
37 }
38}
39
40fn agent_member(service_account: &str) -> Option<(ComputeServiceType, &str)> {
45 ComputeServiceType::ALL.iter().find_map(|service_type| {
46 let suffix = format!("@{}.iam.gserviceaccount.com", agent_domain(service_type));
47 service_account
48 .strip_prefix("service-")
49 .and_then(|rest| rest.strip_suffix(&suffix))
50 .map(|project_number| (service_type.clone(), project_number))
51 })
52}
53
54fn cross_account_members(access: &GcpCrossAccountAccess) -> Vec<String> {
59 let mut members = Vec::new();
60 for service_type in &access.allowed_service_types {
61 let agent_domain = agent_domain(service_type);
62 for project_number in &access.project_numbers {
63 members.push(format!(
64 "serviceAccount:service-{project_number}@{agent_domain}.iam.gserviceaccount.com"
65 ));
66 }
67 }
68 for service_account_email in &access.service_account_emails {
69 members.push(format!("serviceAccount:{service_account_email}"));
70 }
71 members
72}
73
74#[derive(Debug)]
76pub struct GarArtifactRegistry {
77 client: ArtifactRegistryClient,
78 binding_name: String,
79 project_id: String,
80 location: String,
81 repository_name: String,
82 pull_service_account_email: Option<String>,
83 push_service_account_email: Option<String>,
84 gcp_config: GcpClientConfig,
85}
86
87impl GarArtifactRegistry {
88 pub async fn new(
90 binding_name: String,
91 binding: ArtifactRegistryBinding,
92 gcp_config: &GcpClientConfig,
93 ) -> Result<Self> {
94 info!(
95 binding_name = %binding_name,
96 "Initializing GCP Artifact Registry"
97 );
98
99 let client = crate::http_client::create_http_client();
100 let artifact_registry_client = ArtifactRegistryClient::new(client, gcp_config.clone());
101
102 let project_id = gcp_config.project_id.clone();
104 let location = gcp_config.region.clone();
105
106 let config = match binding {
108 ArtifactRegistryBinding::Gar(config) => config,
109 _ => {
110 return Err(AlienError::new(ErrorData::BindingConfigInvalid {
111 env_var: binding_env_var(&binding_name),
112 binding_name: binding_name.clone(),
113 reason: "Expected GAR binding, got different service type".to_string(),
114 }));
115 }
116 };
117
118 let repository_name = config
119 .repository_name
120 .into_value(&binding_name, "repository_name")
121 .context(ErrorData::BindingConfigInvalid {
122 env_var: binding_env_var(&binding_name),
123 binding_name: binding_name.clone(),
124 reason: "Failed to extract repository_name from binding".to_string(),
125 })?;
126
127 let pull_service_account_email = config
128 .pull_service_account_email
129 .map(|v| {
130 v.into_value(&binding_name, "pull_service_account_email")
131 .context(ErrorData::BindingConfigInvalid {
132 env_var: binding_env_var(&binding_name),
133 binding_name: binding_name.clone(),
134 reason: "Failed to extract pull_service_account_email from binding"
135 .to_string(),
136 })
137 })
138 .transpose()?;
139
140 let push_service_account_email = config
141 .push_service_account_email
142 .map(|v| {
143 v.into_value(&binding_name, "push_service_account_email")
144 .context(ErrorData::BindingConfigInvalid {
145 env_var: binding_env_var(&binding_name),
146 binding_name: binding_name.clone(),
147 reason: "Failed to extract push_service_account_email from binding"
148 .to_string(),
149 })
150 })
151 .transpose()?;
152
153 Ok(Self {
154 client: artifact_registry_client,
155 binding_name,
156 project_id,
157 location,
158 repository_name,
159 pull_service_account_email,
160 push_service_account_email,
161 gcp_config: gcp_config.clone(),
162 })
163 }
164
165 fn extract_repo_name(&self, repo_id: &str) -> Result<String> {
169 if repo_id.is_empty() {
170 return Ok(self.repository_name.clone());
171 }
172 if let Some(name) = repo_id.split('/').last() {
173 Ok(name.to_string())
174 } else {
175 Err(AlienError::new(ErrorData::BindingConfigInvalid {
176 env_var: binding_env_var(&self.binding_name),
177 binding_name: self.binding_name.clone(),
178 reason: format!("Invalid repository ID format: {}", repo_id),
179 }))
180 }
181 }
182
183 async fn update_policy_members(
185 &self,
186 repo_name: &str,
187 mut current_policy: IamPolicy,
188 members: Vec<String>,
189 add_members: bool, ) -> Result<IamPolicy> {
191 let reader_role = "roles/artifactregistry.reader";
192
193 let mut binding_index = None;
195 for (i, binding) in current_policy.bindings.iter().enumerate() {
196 if binding.role == reader_role {
197 binding_index = Some(i);
198 break;
199 }
200 }
201
202 if add_members {
203 if members.is_empty() {
205 info!(repo_name = %repo_name, "No new members to add");
206 return Ok(current_policy);
207 }
208
209 match binding_index {
210 Some(i) => {
211 let binding = &mut current_policy.bindings[i];
213 for member in members {
214 if !binding.members.contains(&member) {
215 binding.members.push(member);
216 }
217 }
218 }
219 None => {
220 current_policy
222 .bindings
223 .push(alien_gcp_clients::iam::Binding {
224 role: reader_role.to_string(),
225 members,
226 condition: None,
227 });
228 }
229 }
230 } else {
231 if let Some(i) = binding_index {
233 let binding = &mut current_policy.bindings[i];
234 binding.members.retain(|member| !members.contains(member));
235
236 if binding.members.is_empty() {
238 current_policy.bindings.remove(i);
239 }
240 }
241 }
243
244 let updated = self.client.set_repository_iam_policy(
246 self.project_id.clone(),
247 self.location.clone(),
248 repo_name.to_string(),
249 current_policy,
250 ).await
251 .map_err(|e| map_cloud_client_error(
252 e,
253 format!("Failed to update cross-account access for GCP Artifact Registry repository '{}'", repo_name),
254 Some(repo_name.to_string()),
255 ))?;
256
257 let action = if add_members { "added" } else { "removed" };
258 info!(
259 repo_name = %repo_name,
260 action = %action,
261 "GCP Artifact Registry repository cross-account access updated successfully"
262 );
263 Ok(updated)
264 }
265}
266
267impl Binding for GarArtifactRegistry {}
268
269#[async_trait]
270impl ArtifactRegistry for GarArtifactRegistry {
271 fn registry_endpoint(&self) -> String {
272 format!("https://{}-docker.pkg.dev", self.location)
273 }
274
275 fn upstream_repository_prefix(&self) -> String {
276 format!("{}/{}", self.project_id, self.repository_name)
277 }
278
279 async fn create_repository(&self, repo_name: &str) -> Result<RepositoryResponse> {
280 let routable_name = format!("{}/{}", self.upstream_repository_prefix(), repo_name);
285 Ok(RepositoryResponse {
286 name: routable_name,
287 uri: None,
288 created_at: None,
289 })
290 }
291
292 async fn get_repository(&self, repo_id: &str) -> Result<RepositoryResponse> {
293 let image_path = self.extract_repo_name(repo_id)?;
296 let routable_name = format!("{}/{}", self.upstream_repository_prefix(), image_path);
297 let repository_uri = format!(
298 "{}-docker.pkg.dev/{}/{}",
299 self.location, self.project_id, image_path
300 );
301
302 Ok(RepositoryResponse {
303 name: routable_name,
304 uri: Some(repository_uri),
305 created_at: None,
306 })
307 }
308
309 async fn add_cross_account_access(
310 &self,
311 repo_id: &str,
312 access: CrossAccountAccess,
313 ) -> Result<()> {
314 let _ = repo_id; let repo_name = self.repository_name.clone();
321
322 let gcp_access = match access {
323 CrossAccountAccess::Gcp(gcp_access) => gcp_access,
324 _ => {
325 return Err(AlienError::new(ErrorData::BindingConfigInvalid {
326 env_var: binding_env_var(&self.binding_name),
327 binding_name: self.binding_name.clone(),
328 reason: "GCP artifact registry can only accept GCP cross-account access configuration".to_string(),
329 }));
330 }
331 };
332
333 info!(
334 repo_name = %repo_name,
335 project_numbers = ?gcp_access.project_numbers,
336 allowed_service_types = ?gcp_access.allowed_service_types,
337 service_account_emails = ?gcp_access.service_account_emails,
338 "Adding GCP Artifact Registry repository cross-account access"
339 );
340
341 let current_policy = match self
345 .client
346 .get_repository_iam_policy(
347 self.project_id.clone(),
348 self.location.clone(),
349 repo_name.clone(),
350 )
351 .await
352 {
353 Ok(policy) => policy,
354 Err(error) if error.http_status_code == Some(404) => IamPolicy {
355 version: Some(1),
356 kind: None,
357 resource_id: None,
358 bindings: vec![],
359 etag: None,
360 },
361 Err(error) => {
362 return Err(map_cloud_client_error(
363 error,
364 format!(
365 "Failed to read the IAM policy of GCP Artifact Registry repository '{repo_name}' to grant access"
366 ),
367 Some(repo_name.to_string()),
368 ))
369 }
370 };
371
372 let (deferred, present): (Vec<ComputeServiceType>, Vec<ComputeServiceType>) = gcp_access
376 .allowed_service_types
377 .iter()
378 .cloned()
379 .partition(is_agent_created_on_first_use);
380
381 let policy = self
382 .update_policy_members(
383 &repo_name,
384 current_policy,
385 cross_account_members(&GcpCrossAccountAccess {
386 allowed_service_types: present,
387 ..gcp_access.clone()
388 }),
389 true,
390 )
391 .await?;
392
393 if deferred.is_empty() {
394 return Ok(());
395 }
396 if policy.etag.is_none() {
400 return Err(AlienError::new(ErrorData::RemoteResourceConflict {
401 operation_context: "adding cross-account access".to_string(),
402 resource_type: "repository".to_string(),
403 resource_name: repo_name.clone(),
404 conflict_reason:
405 "the updated IAM policy carries no etag, so the sandbox agent's member cannot be added without replacing the policy"
406 .to_string(),
407 }));
408 }
409 self.update_policy_members(
410 &repo_name,
411 policy,
412 cross_account_members(&GcpCrossAccountAccess {
413 allowed_service_types: deferred,
414 service_account_emails: Vec::new(),
415 ..gcp_access
416 }),
417 true,
418 )
419 .await?;
420 Ok(())
421 }
422
423 async fn remove_cross_account_access(
424 &self,
425 repo_id: &str,
426 access: CrossAccountAccess,
427 ) -> Result<()> {
428 let _ = repo_id;
430 let repo_name = self.repository_name.clone();
431
432 let gcp_access = match access {
433 CrossAccountAccess::Gcp(gcp_access) => gcp_access,
434 _ => {
435 return Err(AlienError::new(ErrorData::BindingConfigInvalid {
436 env_var: binding_env_var(&self.binding_name),
437 binding_name: self.binding_name.clone(),
438 reason: "GCP artifact registry can only accept GCP cross-account access configuration".to_string(),
439 }));
440 }
441 };
442
443 info!(
444 repo_name = %repo_name,
445 project_numbers = ?gcp_access.project_numbers,
446 allowed_service_types = ?gcp_access.allowed_service_types,
447 service_account_emails = ?gcp_access.service_account_emails,
448 "Removing GCP Artifact Registry repository cross-account access"
449 );
450
451 let current_policy = match self
453 .client
454 .get_repository_iam_policy(
455 self.project_id.clone(),
456 self.location.clone(),
457 repo_name.clone(),
458 )
459 .await
460 {
461 Ok(policy) => policy,
462 Err(error) if error.http_status_code == Some(404) => {
466 info!(repo_name = %repo_name, "No existing GCP IAM policy to remove permissions from");
467 return Ok(());
468 }
469 Err(error) => {
470 return Err(map_cloud_client_error(
471 error,
472 format!(
473 "Failed to read the IAM policy of GCP Artifact Registry repository '{repo_name}' to revoke access"
474 ),
475 Some(repo_name.to_string()),
476 ))
477 }
478 };
479
480 let members_to_remove = cross_account_members(&gcp_access);
481 self.update_policy_members(&repo_name, current_policy, members_to_remove, false)
482 .await?;
483 Ok(())
484 }
485
486 async fn get_cross_account_access(&self, repo_id: &str) -> Result<CrossAccountPermissions> {
487 let _ = repo_id;
489 let repo_name = self.repository_name.clone();
490
491 info!(
492 repo_name = %repo_name,
493 "Getting GCP Artifact Registry repository cross-account access"
494 );
495
496 let policy = match self
497 .client
498 .get_repository_iam_policy(
499 self.project_id.clone(),
500 self.location.clone(),
501 repo_name.clone(),
502 )
503 .await
504 {
505 Ok(policy) => policy,
506 Err(e) => {
507 warn!(
508 repo_name = %repo_name,
509 error = %e,
510 "Failed to get GCP Artifact Registry repository IAM policy"
511 );
512 return Ok(CrossAccountPermissions {
514 access: CrossAccountAccess::Gcp(GcpCrossAccountAccess {
515 project_numbers: Vec::new(),
516 allowed_service_types: Vec::new(),
517 service_account_emails: Vec::new(),
518 }),
519 last_updated: None,
520 });
521 }
522 };
523
524 let mut project_numbers = Vec::new();
525 let mut service_account_emails = Vec::new();
526 let mut allowed_service_types = Vec::new();
527
528 for binding in policy.bindings {
529 if binding.role.contains("reader") || binding.role.contains("artifactregistry") {
531 for member in binding.members {
532 if let Some(service_account) = member.strip_prefix("serviceAccount:") {
534 match agent_member(service_account) {
535 Some((service_type, project_number)) => {
536 project_numbers.push(project_number.to_string());
537 if !allowed_service_types.contains(&service_type) {
538 allowed_service_types.push(service_type);
539 }
540 }
541 None => service_account_emails.push(service_account.to_string()),
542 }
543 }
544 }
545 }
546 }
547
548 project_numbers.sort();
550 project_numbers.dedup();
551 service_account_emails.sort();
552 service_account_emails.dedup();
553 allowed_service_types.sort_by_key(|rt| format!("{:?}", rt));
554 allowed_service_types.dedup();
555
556 info!(
557 repo_name = %repo_name,
558 project_numbers = ?project_numbers,
559 allowed_service_types = ?allowed_service_types,
560 service_account_emails = ?service_account_emails,
561 "Retrieved GCP Artifact Registry repository cross-account access"
562 );
563
564 Ok(CrossAccountPermissions {
565 access: CrossAccountAccess::Gcp(GcpCrossAccountAccess {
566 project_numbers,
567 allowed_service_types,
568 service_account_emails,
569 }),
570 last_updated: None, })
572 }
573
574 async fn generate_credentials(
575 &self,
576 repo_id: &str,
577 permissions: ArtifactRegistryPermissions,
578 ttl_seconds: Option<u32>,
579 ) -> Result<ArtifactRegistryCredentials> {
580 info!(
581 repo_id = %repo_id,
582 permissions = ?permissions,
583 ttl_seconds = ?ttl_seconds,
584 "Generating GCP Artifact Registry credentials by impersonating service account"
585 );
586
587 let _project_id = &self.project_id;
590 let _location = &self.location;
591
592 let service_account_email = match permissions {
594 ArtifactRegistryPermissions::Pull => {
595 self.pull_service_account_email.clone()
596 .ok_or_else(|| AlienError::new(ErrorData::BindingConfigInvalid {
597 env_var: binding_env_var(&self.binding_name),
598 binding_name: self.binding_name.clone(),
599 reason: "Pull service account email not available - ensure the artifact registry resource is properly linked".to_string(),
600 }))?
601 }
602 ArtifactRegistryPermissions::PushPull => {
603 self.push_service_account_email.clone()
604 .ok_or_else(|| AlienError::new(ErrorData::BindingConfigInvalid {
605 env_var: binding_env_var(&self.binding_name),
606 binding_name: self.binding_name.clone(),
607 reason: "Push service account email not available - ensure the artifact registry resource is properly linked".to_string(),
608 }))?
609 }
610 };
611
612 info!(
613 service_account_email = %service_account_email,
614 "Using stored service account email for GCP Artifact Registry access"
615 );
616
617 let gcp_config = &self.gcp_config;
619
620 let scopes = vec![
621 "https://www.googleapis.com/auth/cloud-platform".to_string(),
622 "https://www.googleapis.com/auth/devstorage.read_write".to_string(),
623 ];
624
625 let lifetime = ttl_seconds.map(|ttl| format!("{}s", ttl.min(3600))); let impersonation_config = alien_gcp_clients::GcpImpersonationConfig {
628 service_account_email: service_account_email.clone(),
629 scopes,
630 delegates: None,
631 lifetime,
632 target_project_id: None,
633 target_region: None,
634 };
635
636 let impersonated_config =
638 gcp_config
639 .impersonate(impersonation_config)
640 .await
641 .map_err(|e| {
642 map_cloud_client_error(
643 e,
644 "Failed to impersonate GCP service account for artifact registry access"
645 .to_string(),
646 Some(repo_id.to_string()),
647 )
648 })?;
649
650 let access_token = impersonated_config
652 .get_bearer_token("https://www.googleapis.com/")
653 .await
654 .map_err(|e| {
655 map_cloud_client_error(
656 e,
657 "Failed to get OAuth token from impersonated service account".to_string(),
658 Some(repo_id.to_string()),
659 )
660 })?;
661
662 let expires_at = if let Some(ttl) = ttl_seconds {
664 Some(
665 (chrono::Utc::now() + chrono::Duration::seconds(ttl.min(3600) as i64)).to_rfc3339(),
666 )
667 } else {
668 Some((chrono::Utc::now() + chrono::Duration::seconds(3600)).to_rfc3339())
669 };
671
672 info!(
673 permissions = ?permissions,
674 service_account = %service_account_email,
675 "GCP Artifact Registry OAuth token generated successfully with impersonated service account"
676 );
677
678 Ok(ArtifactRegistryCredentials {
680 auth_method: RegistryAuthMethod::Basic,
681 username: "oauth2accesstoken".to_string(),
682 password: access_token,
683 expires_at,
684 })
685 }
686
687 async fn delete_repository(&self, repo_id: &str) -> Result<()> {
688 debug!(
702 repo_id = %repo_id,
703 "GCP Artifact Registry delete_repository: no-op (image paths are implicit)"
704 );
705 Ok(())
706 }
707}
708
709#[cfg(test)]
710mod tests {
711 use super::*;
712
713 #[test]
717 fn each_service_type_resolves_to_its_own_service_agent() {
718 let members = cross_account_members(&GcpCrossAccountAccess {
719 project_numbers: vec!["123456789012".to_string()],
720 allowed_service_types: vec![ComputeServiceType::Worker, ComputeServiceType::Sandbox],
721 service_account_emails: vec![
722 "management@test-project.iam.gserviceaccount.com".to_string()
723 ],
724 });
725
726 assert_eq!(
727 members,
728 vec![
729 "serviceAccount:service-123456789012@serverless-robot-prod.iam.gserviceaccount.com",
730 "serviceAccount:service-123456789012@gcp-sa-vertex-sandbox.iam.gserviceaccount.com",
731 "serviceAccount:management@test-project.iam.gserviceaccount.com",
732 ]
733 );
734 }
735
736 #[test]
741 fn every_service_type_survives_a_write_then_read() {
742 for service_type in [ComputeServiceType::Worker, ComputeServiceType::Sandbox] {
743 let written = cross_account_members(&GcpCrossAccountAccess {
744 project_numbers: vec!["123456789012".to_string()],
745 allowed_service_types: vec![service_type.clone()],
746 service_account_emails: Vec::new(),
747 });
748 let member = written
749 .first()
750 .expect("a project and a type name one member");
751 let account = member
752 .strip_prefix("serviceAccount:")
753 .expect("members carry the IAM principal prefix");
754 assert_eq!(
756 agent_member(account),
757 Some((service_type.clone(), "123456789012")),
758 "{service_type:?} must decode back to the project that earned it"
759 );
760 }
761 }
762
763 #[test]
767 fn the_worker_and_sandbox_members_are_written_apart() {
768 let access = GcpCrossAccountAccess {
769 project_numbers: vec!["123456789012".to_string()],
770 allowed_service_types: vec![ComputeServiceType::Worker, ComputeServiceType::Sandbox],
771 service_account_emails: vec![
772 "management@test-project.iam.gserviceaccount.com".to_string()
773 ],
774 };
775 let first = cross_account_members(&GcpCrossAccountAccess {
776 allowed_service_types: vec![ComputeServiceType::Worker],
777 ..access.clone()
778 });
779 let second = cross_account_members(&GcpCrossAccountAccess {
780 allowed_service_types: vec![ComputeServiceType::Sandbox],
781 service_account_emails: Vec::new(),
782 ..access.clone()
783 });
784
785 assert!(second
786 .iter()
787 .all(|member| member.contains("gcp-sa-vertex-sandbox")));
788 assert!(first
789 .iter()
790 .all(|member| !member.contains("gcp-sa-vertex-sandbox")));
791 let mut both = [first, second].concat();
793 both.sort();
794 let mut whole = cross_account_members(&access);
795 whole.sort();
796 assert_eq!(both, whole);
797 }
798
799 #[test]
801 fn a_service_type_grants_nothing_for_a_project_it_does_not_name() {
802 let members = cross_account_members(&GcpCrossAccountAccess {
805 project_numbers: Vec::new(),
806 allowed_service_types: vec![ComputeServiceType::Sandbox],
807 service_account_emails: vec![
808 "management@test-project.iam.gserviceaccount.com".to_string()
809 ],
810 });
811
812 assert_eq!(
813 members,
814 vec!["serviceAccount:management@test-project.iam.gserviceaccount.com"]
815 );
816 }
817}