fix: preauthorize repository access without persisting secrets

This commit is contained in:
2026-08-26 17:23:43 +09:00
parent df34533765
commit b644971d45
3 changed files with 313 additions and 44 deletions
+75 -2
View File
@@ -585,12 +585,13 @@ impl Runtime {
let worker_id = request.worker_id; let worker_id = request.worker_id;
let worker_ref = WorkerRef::new(worker_id); let worker_ref = WorkerRef::new(worker_id);
let durable_request = durable_create_worker_request(&request);
let record = WorkerRecord { let record = WorkerRecord {
worker_ref: worker_ref.clone(), worker_ref: worker_ref.clone(),
worker_id: worker_id.clone(), worker_id: worker_id.clone(),
status: WorkerStatus::Stopped, status: WorkerStatus::Stopped,
workspace_id: scope.map(|scope| scope.workspace_id.clone()), workspace_id: scope.map(|scope| scope.workspace_id.clone()),
request: request.clone(), request: durable_request,
run_generation: 1, run_generation: 1,
working_directory: None, working_directory: None,
execution_handle: None, execution_handle: None,
@@ -2733,6 +2734,16 @@ fn worker_status_from_run_state(run_state: WorkerExecutionRunState) -> WorkerSta
} }
} }
fn durable_create_worker_request(request: &CreateWorkerRequest) -> CreateWorkerRequest {
let mut durable = request.clone();
if let Some(working_directory) = durable.working_directory_request.as_mut()
&& let Some(materialization) = working_directory.materialization.as_mut()
{
materialization.ssh = None;
}
durable
}
fn requested_primary_workdir_id(request: &CreateWorkerRequest) -> Option<&str> { fn requested_primary_workdir_id(request: &CreateWorkerRequest) -> Option<&str> {
request request
.working_directory .working_directory
@@ -2903,7 +2914,9 @@ fn subscription_worker_state(status: WorkerStatus) -> SubscriptionWorkerState {
mod tests { mod tests {
use super::*; use super::*;
use crate::catalog::{ use crate::catalog::{
ConfigBundleRef, ProfileSelector, WorkingDirectoryClaim, WorkspaceApiRef, ConfigBundleRef, MaterializerKind, ProfileSelector, RepositoryMaterializationContext,
RepositorySshMaterializationAccess, SensitiveString, WorkingDirectoryClaim,
WorkingDirectoryRepository, WorkingDirectoryRequest, WorkspaceApiRef,
}; };
use crate::config_bundle::{ use crate::config_bundle::{
ConfigBundle, ConfigBundleMetadata, ConfigBundleProvenance, ConfigDeclaration, ConfigBundle, ConfigBundleMetadata, ConfigBundleProvenance, ConfigDeclaration,
@@ -3134,6 +3147,66 @@ mod tests {
} }
} }
#[test]
fn durable_worker_request_omits_repository_credentials() {
let mut request = task_request("worker-secret-redaction");
request.working_directory_request = Some(WorkingDirectoryRequest {
repository: WorkingDirectoryRepository {
id: "repository-1".to_string(),
provider: "git".to_string(),
source: workspace_api::RepositorySource {
kind: workspace_api::RepositorySourceKind::Ssh,
uri: "ssh://git@example.test/repo.git".to_string(),
},
source_revision: 1,
source_fingerprint: "sha256:source".to_string(),
selector: None,
},
materializer: MaterializerKind::RuntimeGitCache,
backend_workdir_id: Some("working-directory-1".to_string()),
materialization: Some(RepositoryMaterializationContext {
workspace_id: "workspace-1".to_string(),
runtime_id: "runtime-1".to_string(),
operation_id: "operation-1".to_string(),
config_revision: 1,
config_projection_digest: "sha256:projection".to_string(),
cache_generation: 0,
ssh: Some(RepositorySshMaterializationAccess {
credential_id: "credential-1".to_string(),
credential_revision: 1,
host_trust_id: "host-trust-1".to_string(),
host_trust_revision: 1,
access: workspace_api::RepositoryAccessMode::ReadOnly,
expires_at_epoch_seconds: u64::MAX,
private_key: SensitiveString::new("private-key-bytes"),
known_hosts_entry: SensitiveString::new("known-hosts-entry"),
}),
}),
});
let durable = durable_create_worker_request(&request);
assert!(
request
.working_directory_request
.as_ref()
.and_then(|working_directory| working_directory.materialization.as_ref())
.and_then(|materialization| materialization.ssh.as_ref())
.is_some()
);
assert!(
durable
.working_directory_request
.as_ref()
.and_then(|working_directory| working_directory.materialization.as_ref())
.and_then(|materialization| materialization.ssh.as_ref())
.is_none()
);
let serialized = serde_json::to_string(&durable).unwrap();
assert!(!serialized.contains("private-key-bytes"));
assert!(!serialized.contains("known-hosts-entry"));
}
fn scoped_task_request(objective: &str, workspace_id: &str) -> CreateWorkerRequest { fn scoped_task_request(objective: &str, workspace_id: &str) -> CreateWorkerRequest {
let mut request = task_request(objective); let mut request = task_request(objective);
request.workspace_api = Some(WorkspaceApiRef { request.workspace_api = Some(WorkspaceApiRef {
+154 -42
View File
@@ -395,6 +395,44 @@ impl RuntimeGitCacheMaterializer {
}) })
} }
fn cache_repository_access(
&self,
working_directory_id: &str,
ssh: &RepositorySshMaterializationAccess,
) -> Result<(), WorkingDirectoryDiagnostic> {
validate_ssh_materialization_access(ssh)?;
let working_directory_id = working_directory_id.to_string();
let credential_revision = ssh.credential_revision;
let expires_at = ssh.expires_at_epoch_seconds;
self.repository_access
.lock()
.map_err(|_| {
WorkingDirectoryDiagnostic::new(
"working_directory_repository_access_unavailable",
"Runtime Repository access state is unavailable",
)
})?
.insert(working_directory_id.clone(), ssh.clone());
let repository_access = self.repository_access.clone();
std::thread::spawn(move || {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
if expires_at > now {
std::thread::sleep(Duration::from_secs(expires_at - now));
}
if let Ok(mut access) = repository_access.lock()
&& access
.get(&working_directory_id)
.is_some_and(|access| access.credential_revision == credential_revision)
{
access.remove(&working_directory_id);
}
});
Ok(())
}
fn bind_repository_access( fn bind_repository_access(
&self, &self,
working_directory_id: &str, working_directory_id: &str,
@@ -657,13 +695,62 @@ impl RuntimeGitCacheMaterializer {
Ok(cache_path) Ok(cache_path)
} }
fn request_with_authorized_repository_access(
&self,
working_directory_id: &str,
request: &WorkingDirectoryRequest,
) -> Result<WorkingDirectoryRequest, WorkingDirectoryDiagnostic> {
let mut request = request.clone();
if request.repository.provider != "git"
|| request.repository.source.kind != workspace_api::RepositorySourceKind::Ssh
|| request
.materialization
.as_ref()
.and_then(|materialization| materialization.ssh.as_ref())
.is_some()
{
return Ok(request);
}
let access = self
.repository_access
.lock()
.map_err(|_| {
WorkingDirectoryDiagnostic::new(
"working_directory_repository_access_unavailable",
"Runtime Repository access state is unavailable",
)
})?
.get(working_directory_id)
.cloned()
.ok_or_else(|| {
WorkingDirectoryDiagnostic::new(
"working_directory_remote_repository_access_required",
"SSH Repository materialization requires pre-authorized credential and host-trust authority",
)
})?;
validate_ssh_materialization_access(&access)?;
request
.materialization
.as_mut()
.ok_or_else(|| {
WorkingDirectoryDiagnostic::new(
"working_directory_repository_materialization_authority_required",
"remote Repository materialization requires Backend-authored operation authority",
)
})?
.ssh = Some(access);
Ok(request)
}
fn materialize_with_working_directory_id( fn materialize_with_working_directory_id(
&self, &self,
working_directory_id: String, working_directory_id: String,
request: &WorkingDirectoryRequest, request: &WorkingDirectoryRequest,
) -> Result<WorkingDirectoryBinding, WorkingDirectoryDiagnostic> { ) -> Result<WorkingDirectoryBinding, WorkingDirectoryDiagnostic> {
validate_working_directory_id(&working_directory_id)?; validate_working_directory_id(&working_directory_id)?;
let repository_cache = self.ensure_repository_cache(request)?; let request =
self.request_with_authorized_repository_access(&working_directory_id, request)?;
let repository_cache = self.ensure_repository_cache(&request)?;
let selector = request.repository.selector.as_deref().unwrap_or("HEAD"); let selector = request.repository.selector.as_deref().unwrap_or("HEAD");
let resolved_commit = resolve_cached_commit(&repository_cache, selector)?; let resolved_commit = resolve_cached_commit(&repository_cache, selector)?;
let tree_spec = format!("{resolved_commit}^{{tree}}"); let tree_spec = format!("{resolved_commit}^{{tree}}");
@@ -725,7 +812,7 @@ impl RuntimeGitCacheMaterializer {
materializer_kind: MaterializerKind::RuntimeGitCache, materializer_kind: MaterializerKind::RuntimeGitCache,
repository_source_revision: Some(request.repository.source_revision), repository_source_revision: Some(request.repository.source_revision),
repository_source_fingerprint: Some(request.repository.source_fingerprint.clone()), repository_source_fingerprint: Some(request.repository.source_fingerprint.clone()),
repository_cache_key: Some(Self::repository_cache_key(request)), repository_cache_key: Some(Self::repository_cache_key(&request)),
cache_generation: context cache_generation: context
.map(|value| value.cache_generation) .map(|value| value.cache_generation)
.unwrap_or_default(), .unwrap_or_default(),
@@ -741,7 +828,7 @@ impl RuntimeGitCacheMaterializer {
}, },
cleanup_target: WorkingDirectoryCleanupTarget { cleanup_target: WorkingDirectoryCleanupTarget {
kind: "runtime_git_cache_worktree".to_string(), kind: "runtime_git_cache_worktree".to_string(),
working_directory_id, working_directory_id: working_directory_id.clone(),
repository_id: request.repository.id.clone(), repository_id: request.repository.id.clone(),
}, },
status: WorkingDirectoryStatusKind::Active, status: WorkingDirectoryStatusKind::Active,
@@ -760,6 +847,16 @@ impl RuntimeGitCacheMaterializer {
let _ = fs::remove_dir_all(&working_directory_root); let _ = fs::remove_dir_all(&working_directory_root);
return Err(error); return Err(error);
} }
if let Some(ssh) = request
.materialization
.as_ref()
.and_then(|materialization| materialization.ssh.as_ref())
&& let Err(error) = self.cache_repository_access(&working_directory_id, ssh)
{
remove_cached_worktree(&repository_cache, &worktree_root);
let _ = fs::remove_dir_all(&working_directory_root);
return Err(error);
}
Ok(binding) Ok(binding)
} }
} }
@@ -797,7 +894,13 @@ impl WorkingDirectoryMaterializer for RuntimeGitCacheMaterializer {
) )
})?; })?;
validate_ssh_materialization_access(ssh)?; validate_ssh_materialization_access(ssh)?;
let mut binding = self.read_binding(&request.working_directory_id)?; let mut binding = match self.read_binding(&request.working_directory_id) {
Ok(binding) => binding,
Err(error) if error.code == "working_directory_not_found" => {
return self.cache_repository_access(&request.working_directory_id, ssh);
}
Err(error) => return Err(error),
};
if binding if binding
.working_directory .working_directory
.evidence .evidence
@@ -814,36 +917,7 @@ impl WorkingDirectoryMaterializer for RuntimeGitCacheMaterializer {
binding.working_directory.evidence.credential_revision = Some(ssh.credential_revision); binding.working_directory.evidence.credential_revision = Some(ssh.credential_revision);
binding.working_directory.evidence.host_trust_revision = Some(ssh.host_trust_revision); binding.working_directory.evidence.host_trust_revision = Some(ssh.host_trust_revision);
self.write_record(&binding)?; self.write_record(&binding)?;
let working_directory_id = request.working_directory_id.clone(); self.cache_repository_access(&request.working_directory_id, ssh)
let credential_revision = ssh.credential_revision;
let expires_at = ssh.expires_at_epoch_seconds;
self.repository_access
.lock()
.map_err(|_| {
WorkingDirectoryDiagnostic::new(
"working_directory_repository_access_unavailable",
"Runtime Repository access state is unavailable",
)
})?
.insert(working_directory_id.clone(), ssh.clone());
let repository_access = self.repository_access.clone();
std::thread::spawn(move || {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
if expires_at > now {
std::thread::sleep(Duration::from_secs(expires_at - now));
}
if let Ok(mut access) = repository_access.lock()
&& access
.get(&working_directory_id)
.is_some_and(|access| access.credential_revision == credential_revision)
{
access.remove(&working_directory_id);
}
});
Ok(())
} }
fn bind_working_directory( fn bind_working_directory(
@@ -2508,8 +2582,52 @@ mod tests {
drop(command_access); drop(command_access);
assert!(!operation_socket.exists()); assert!(!operation_socket.exists());
let initial_materialization = request.materialization.clone().unwrap();
let working_directory_id = "working-directory-agent".to_string();
materializer
.authorize_repository_access(&WorkingDirectoryRepositoryAccessRequest {
working_directory_id: working_directory_id.clone(),
materialization: initial_materialization.clone(),
})
.unwrap();
request.backend_workdir_id = Some(working_directory_id.clone());
request.materialization.as_mut().unwrap().ssh = None;
let mut authorized_ssh_request = request.clone();
authorized_ssh_request.repository.source = workspace_api::RepositorySource {
kind: workspace_api::RepositorySourceKind::Ssh,
uri: "ssh://git@example.test/repo.git".to_string(),
};
assert!(
materializer
.request_with_authorized_repository_access(
&working_directory_id,
&authorized_ssh_request,
)
.unwrap()
.materialization
.unwrap()
.ssh
.is_some()
);
let created = materializer.create(&request).unwrap(); let created = materializer.create(&request).unwrap();
let id = created.working_directory.id; let id = created.working_directory.id;
assert_eq!(id, working_directory_id);
let mut ssh_backed_binding = materializer.read_binding(&id).unwrap();
ssh_backed_binding
.working_directory
.evidence
.credential_revision = Some(1);
ssh_backed_binding
.working_directory
.evidence
.host_trust_revision = Some(1);
materializer.write_record(&ssh_backed_binding).unwrap();
materializer
.authorize_repository_access(&WorkingDirectoryRepositoryAccessRequest {
working_directory_id: id.clone(),
materialization: initial_materialization.clone(),
})
.unwrap();
assert_eq!(materializer.list_working_directories().unwrap().len(), 1); assert_eq!(materializer.list_working_directories().unwrap().len(), 1);
assert_eq!( assert_eq!(
fs::read_dir(runtime_root.path().join(".repository-agents")) fs::read_dir(runtime_root.path().join(".repository-agents"))
@@ -2517,12 +2635,6 @@ mod tests {
.unwrap_or_default(), .unwrap_or_default(),
0 0
); );
materializer
.authorize_repository_access(&WorkingDirectoryRepositoryAccessRequest {
working_directory_id: id.clone(),
materialization: request.materialization.clone().unwrap(),
})
.unwrap();
let binding = materializer.bind_working_directory(&id, None).unwrap(); let binding = materializer.bind_working_directory(&id, None).unwrap();
let environment = binding.command_environment(); let environment = binding.command_environment();
let socket = PathBuf::from(environment["SSH_AUTH_SOCK"].clone()); let socket = PathBuf::from(environment["SSH_AUTH_SOCK"].clone());
@@ -2543,7 +2655,7 @@ mod tests {
.code, .code,
"working_directory_remote_repository_access_required" "working_directory_remote_repository_access_required"
); );
let mut rotated = request.materialization.clone().unwrap(); let mut rotated = initial_materialization.clone();
rotated.operation_id = "operation-agent-rotated".to_string(); rotated.operation_id = "operation-agent-rotated".to_string();
rotated.ssh.as_mut().unwrap().credential_revision = 2; rotated.ssh.as_mut().unwrap().credential_revision = 2;
rotated.ssh.as_mut().unwrap().access = workspace_api::RepositoryAccessMode::ReadOnly; rotated.ssh.as_mut().unwrap().access = workspace_api::RepositoryAccessMode::ReadOnly;
@@ -2643,7 +2755,7 @@ mod tests {
); );
drop(rebound); drop(rebound);
let mut expired = request.materialization.clone().unwrap(); let mut expired = initial_materialization;
expired.ssh.as_mut().unwrap().expires_at_epoch_seconds = 1; expired.ssh.as_mut().unwrap().expires_at_epoch_seconds = 1;
assert_eq!( assert_eq!(
materializer materializer
+84
View File
@@ -1227,6 +1227,35 @@ pub async fn build_workspace_server_router(
))) )))
} }
fn take_new_workdir_repository_access(
request: &mut WorkerSpawnRequest,
) -> ApiResult<Option<worker_runtime::catalog::WorkingDirectoryRepositoryAccessRequest>> {
let Some(working_directory) = request.resolved_working_directory_request.as_mut() else {
return Ok(None);
};
let Some(materialization) = working_directory.materialization.as_mut() else {
return Ok(None);
};
if materialization.ssh.is_none() {
return Ok(None);
}
let working_directory_id = working_directory
.backend_workdir_id
.clone()
.ok_or_else(|| {
Error::Config(
"repository access authorization requires a Backend WorkingDirectory id"
.to_string(),
)
})?;
let access = worker_runtime::catalog::WorkingDirectoryRepositoryAccessRequest {
working_directory_id,
materialization: materialization.clone(),
};
materialization.ssh = None;
Ok(Some(access))
}
impl WorkspaceApi { impl WorkspaceApi {
pub fn with_config_schema_provider( pub fn with_config_schema_provider(
mut self, mut self,
@@ -1437,6 +1466,11 @@ impl WorkspaceApi {
self.validate_worker_spawn_repository_scope(&request)?; self.validate_worker_spawn_repository_scope(&request)?;
let workspace_api = self.workspace_api_ref(runtime_id); let workspace_api = self.workspace_api_ref(runtime_id);
request.resolved_workspace_api = Some(workspace_api.clone()); request.resolved_workspace_api = Some(workspace_api.clone());
if let Some(access) = take_new_workdir_repository_access(&mut request)? {
self.runtime
.authorize_working_directory_repository_access(runtime_id, access)
.map_err(RuntimeRegistryError::into_error)?;
}
if let Some(working_directory) = request.resolved_working_directory.as_ref() if let Some(working_directory) = request.resolved_working_directory.as_ref()
&& let Some(access) = repository_access_request_for_workdir( && let Some(access) = repository_access_request_for_workdir(
self, self,
@@ -15387,6 +15421,56 @@ mod tests {
api.validate_worker_spawn_repository_scope(&workdir_flow_launch) api.validate_worker_spawn_repository_scope(&workdir_flow_launch)
.is_err() .is_err()
); );
let mut repository_access_launch = workdir_flow_launch;
let working_directory = repository_access_launch
.resolved_working_directory_request
.as_mut()
.unwrap();
working_directory.backend_workdir_id = Some("working-directory-1".to_string());
working_directory.materialization =
Some(worker_runtime::catalog::RepositoryMaterializationContext {
workspace_id: api.config.workspace_id.clone(),
runtime_id: "runtime-1".to_string(),
operation_id: "operation-1".to_string(),
config_revision: 1,
config_projection_digest: "sha256:projection".to_string(),
cache_generation: 0,
ssh: Some(
worker_runtime::catalog::RepositorySshMaterializationAccess {
credential_id: "credential-1".to_string(),
credential_revision: 1,
host_trust_id: "host-trust-1".to_string(),
host_trust_revision: 1,
access: workspace_api::RepositoryAccessMode::ReadOnly,
expires_at_epoch_seconds: u64::MAX,
private_key: worker_runtime::catalog::SensitiveString::new(
"private-key-bytes",
),
known_hosts_entry: worker_runtime::catalog::SensitiveString::new(
"known-hosts-entry",
),
},
),
});
let access = take_new_workdir_repository_access(&mut repository_access_launch)
.unwrap()
.unwrap();
assert_eq!(access.working_directory_id, "working-directory-1");
assert!(access.materialization.ssh.is_some());
assert!(
repository_access_launch
.resolved_working_directory_request
.as_ref()
.and_then(|request| request.materialization.as_ref())
.and_then(|materialization| materialization.ssh.as_ref())
.is_none()
);
let serialized = serde_json::to_string(&repository_access_launch).unwrap();
assert!(!serialized.contains("private-key-bytes"));
assert!(!serialized.contains("known-hosts-entry"));
} }
#[test] #[test]