Source commit: f70e47af061d8f7ba27cd769c8cd2e2f8707dc72 Public tree identity: sha256:03a43db60d295811688bedaa5ad6460df0dbde27fbf37033997f4d8100f83ab2
799 lines
32 KiB
Rust
799 lines
32 KiB
Rust
use std::io::{Read, Seek, SeekFrom, Write};
|
|
|
|
use base64::{engine::general_purpose::STANDARD as BASE64_STANDARD, Engine as _};
|
|
use clusterflux_core::{
|
|
generate_opaque_token, Actor, ArtifactId, AuthContext, DataPlaneObject, DataPlaneScope, Digest,
|
|
DownloadPolicy, NodeEndpoint, NodeId, ProjectId, RendezvousRequest, ResourceLimits,
|
|
ResourceMeter, StorageLocation, TenantId, UserId,
|
|
};
|
|
use sha2::{Digest as _, Sha256};
|
|
|
|
use crate::CoordinatorError;
|
|
|
|
use super::relay::RelayFinishReason;
|
|
use super::{
|
|
bounded_ttl, ArtifactTransferAssignment, CoordinatorResponse, CoordinatorService,
|
|
CoordinatorServiceError,
|
|
};
|
|
|
|
pub(super) const MAX_ARTIFACT_REVERSE_CHUNK_BYTES: u64 = 256 * 1024;
|
|
|
|
#[derive(Debug)]
|
|
pub(super) struct ArtifactReverseTransfer {
|
|
transfer_id: String,
|
|
token_digest: Digest,
|
|
tenant: TenantId,
|
|
project: ProjectId,
|
|
source_node: NodeId,
|
|
artifact: ArtifactId,
|
|
expected_digest: Digest,
|
|
expected_size_bytes: u64,
|
|
expires_at_epoch_seconds: u64,
|
|
spool: tempfile::NamedTempFile,
|
|
received_bytes: u64,
|
|
content_hasher: Sha256,
|
|
delivered_offset: u64,
|
|
error: Option<String>,
|
|
}
|
|
|
|
impl CoordinatorService {
|
|
pub(super) fn handle_create_artifact_download_link(
|
|
&mut self,
|
|
tenant: String,
|
|
project: String,
|
|
actor_user: String,
|
|
artifact: String,
|
|
max_bytes: u64,
|
|
ttl_seconds: u64,
|
|
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
|
|
let context = user_context(tenant, project, actor_user);
|
|
let artifact = ArtifactId::new(artifact);
|
|
let policy = DownloadPolicy { max_bytes };
|
|
let action = self
|
|
.artifact_registry
|
|
.download_action(&context, &artifact, &policy)?;
|
|
self.ensure_download_source_connectivity(&action.source)?;
|
|
let downloadable_size = self
|
|
.artifact_registry
|
|
.downloadable_size(&context, &artifact, &policy)?;
|
|
let now_epoch_seconds = self.current_epoch_seconds()?;
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.expire(now_epoch_seconds);
|
|
Ok(())
|
|
})?;
|
|
self.quota.can_charge_download(
|
|
&context.tenant,
|
|
&context.project,
|
|
downloadable_size,
|
|
now_epoch_seconds,
|
|
)?;
|
|
let token_nonce = generate_opaque_token("artifact_download")
|
|
.map_err(CoordinatorServiceError::Protocol)?;
|
|
let ttl_seconds = bounded_ttl(
|
|
ttl_seconds,
|
|
self.admission.max_artifact_download_ttl_seconds,
|
|
);
|
|
let expires_at_epoch_seconds = now_epoch_seconds.saturating_add(ttl_seconds);
|
|
let account = match &context.actor {
|
|
Actor::User(user) => user.clone(),
|
|
Actor::Agent(_) | Actor::Node(_) | Actor::Task(_) => {
|
|
return Err(CoordinatorServiceError::Protocol(
|
|
"artifact relay download requires a user account".to_owned(),
|
|
))
|
|
}
|
|
};
|
|
let mut relay_candidate = self.artifact_relay.clone();
|
|
relay_candidate
|
|
.reserve(
|
|
token_nonce.clone(),
|
|
context.tenant.clone(),
|
|
context.project.clone(),
|
|
account,
|
|
downloadable_size,
|
|
MAX_ARTIFACT_REVERSE_CHUNK_BYTES,
|
|
expires_at_epoch_seconds,
|
|
now_epoch_seconds,
|
|
)
|
|
.map_err(|error| CoordinatorServiceError::Protocol(error.to_string()))?;
|
|
let link = match self.artifact_registry.create_download_link(
|
|
&context,
|
|
&artifact,
|
|
&policy,
|
|
&token_nonce,
|
|
now_epoch_seconds,
|
|
ttl_seconds,
|
|
) {
|
|
Ok(link) => link,
|
|
Err(error) => return Err(error.into()),
|
|
};
|
|
if let Err(error) =
|
|
relay_candidate.rekey(&token_nonce, link.scoped_token_digest.as_str().to_owned())
|
|
{
|
|
let _ = self.artifact_registry.revoke_download_link(
|
|
&context,
|
|
&artifact,
|
|
&link.scoped_token_digest,
|
|
);
|
|
return Err(CoordinatorServiceError::Protocol(error.to_string()));
|
|
}
|
|
if let Err(error) = self.commit_artifact_relay(relay_candidate) {
|
|
let _ = self.artifact_registry.revoke_download_link(
|
|
&context,
|
|
&artifact,
|
|
&link.scoped_token_digest,
|
|
);
|
|
return Err(error);
|
|
}
|
|
Ok(CoordinatorResponse::ArtifactDownloadLink { link })
|
|
}
|
|
|
|
pub(super) fn handle_open_artifact_download_stream(
|
|
&mut self,
|
|
tenant: String,
|
|
project: String,
|
|
actor_user: String,
|
|
artifact: String,
|
|
max_bytes: u64,
|
|
token_digest: Digest,
|
|
chunk_bytes: u64,
|
|
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
|
|
let context = user_context(tenant, project, actor_user);
|
|
let artifact = ArtifactId::new(artifact);
|
|
let policy = DownloadPolicy { max_bytes };
|
|
let now_epoch_seconds = self.current_epoch_seconds()?;
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.expire(now_epoch_seconds);
|
|
Ok(())
|
|
})?;
|
|
self.artifact_registry
|
|
.expire_download_links(now_epoch_seconds);
|
|
let downloadable_size = self
|
|
.artifact_registry
|
|
.downloadable_size(&context, &artifact, &policy)?;
|
|
let validation_limits = ResourceLimits::unlimited();
|
|
let mut validation_meter = ResourceMeter::default();
|
|
let mut stream = self.artifact_registry.open_download_stream(
|
|
clusterflux_core::DownloadStreamRequest {
|
|
context: &context,
|
|
artifact: &artifact,
|
|
policy: &policy,
|
|
presented_token_digest: &token_digest,
|
|
now_epoch_seconds,
|
|
limits: &validation_limits,
|
|
},
|
|
&mut validation_meter,
|
|
)?;
|
|
self.ensure_download_source_connectivity(&stream.link.source)?;
|
|
self.expire_artifact_reverse_transfers(now_epoch_seconds)?;
|
|
|
|
if let Some(transfer_id) = self.artifact_transfer_by_token.get(&token_digest).cloned() {
|
|
if let Some(message) = self
|
|
.artifact_reverse_transfers
|
|
.get(&transfer_id)
|
|
.and_then(|transfer| transfer.error.clone())
|
|
{
|
|
self.artifact_reverse_transfers.remove(&transfer_id);
|
|
self.artifact_transfer_by_token.remove(&token_digest);
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.finish(token_digest.as_str(), RelayFinishReason::Failed);
|
|
Ok(())
|
|
})?;
|
|
return Err(CoordinatorServiceError::Protocol(format!(
|
|
"retaining node could not stream artifact: {message}"
|
|
)));
|
|
}
|
|
let (content_offset, content, end, complete) = {
|
|
let transfer = self
|
|
.artifact_reverse_transfers
|
|
.get_mut(&transfer_id)
|
|
.ok_or_else(|| {
|
|
CoordinatorServiceError::Protocol(
|
|
"artifact reverse transfer index is inconsistent".to_owned(),
|
|
)
|
|
})?;
|
|
if transfer.tenant != context.tenant
|
|
|| transfer.project != context.project
|
|
|| transfer.artifact != artifact
|
|
|| transfer.token_digest != token_digest
|
|
{
|
|
return Err(clusterflux_core::DownloadError::InvalidToken.into());
|
|
}
|
|
if transfer.received_bytes != transfer.expected_size_bytes {
|
|
let overhead = self.artifact_relay.framing_overhead_bytes();
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.charge_egress(token_digest.as_str(), overhead, now_epoch_seconds)
|
|
})?;
|
|
return Ok(CoordinatorResponse::ArtifactDownloadStream {
|
|
link: stream.link,
|
|
streamed_bytes: 0,
|
|
charged_download_bytes: self.quota.used_download_bytes(
|
|
&context.tenant,
|
|
&context.project,
|
|
now_epoch_seconds,
|
|
),
|
|
content_bytes_available: false,
|
|
content_offset: None,
|
|
content_eof: false,
|
|
content_base64: None,
|
|
content_source: Some("retaining_node_reverse_stream_pending".to_owned()),
|
|
});
|
|
}
|
|
if chunk_bytes == 0 && downloadable_size != 0 {
|
|
return Err(CoordinatorServiceError::Protocol(
|
|
"artifact download chunk_bytes must be greater than zero".to_owned(),
|
|
));
|
|
}
|
|
let requested = chunk_bytes
|
|
.min(downloadable_size)
|
|
.min(MAX_ARTIFACT_REVERSE_CHUNK_BYTES);
|
|
let start = transfer.delivered_offset;
|
|
let end = start.saturating_add(requested).min(transfer.received_bytes);
|
|
let length = usize::try_from(end.saturating_sub(start)).map_err(|_| {
|
|
CoordinatorServiceError::Protocol(
|
|
"artifact download chunk length does not fit memory bounds".to_owned(),
|
|
)
|
|
})?;
|
|
let mut content = vec![0_u8; length];
|
|
let mut spool = transfer.spool.reopen().map_err(|error| {
|
|
CoordinatorServiceError::Protocol(format!(
|
|
"open bounded artifact transfer spool: {error}"
|
|
))
|
|
})?;
|
|
spool.seek(SeekFrom::Start(start)).map_err(|error| {
|
|
CoordinatorServiceError::Protocol(format!(
|
|
"seek bounded artifact transfer spool: {error}"
|
|
))
|
|
})?;
|
|
spool.read_exact(&mut content).map_err(|error| {
|
|
CoordinatorServiceError::Protocol(format!(
|
|
"read bounded artifact transfer spool: {error}"
|
|
))
|
|
})?;
|
|
(start, content, end, end == transfer.received_bytes)
|
|
};
|
|
let streamed_bytes = content.len() as u64;
|
|
self.artifact_registry.stream_download_chunk(
|
|
&mut stream,
|
|
&validation_limits,
|
|
&mut validation_meter,
|
|
streamed_bytes,
|
|
)?;
|
|
let charged_download_bytes = self.quota.charge_download(
|
|
&context.tenant,
|
|
&context.project,
|
|
streamed_bytes,
|
|
now_epoch_seconds,
|
|
)?;
|
|
let content_base64 = BASE64_STANDARD.encode(content);
|
|
let egress_wire_bytes = (content_base64.len() as u64)
|
|
.saturating_add(self.artifact_relay.framing_overhead_bytes());
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.charge_egress(
|
|
token_digest.as_str(),
|
|
egress_wire_bytes,
|
|
now_epoch_seconds,
|
|
)?;
|
|
if complete {
|
|
ledger.finish(token_digest.as_str(), RelayFinishReason::Completed);
|
|
}
|
|
Ok(())
|
|
})?;
|
|
if let Some(transfer) = self.artifact_reverse_transfers.get_mut(&transfer_id) {
|
|
transfer.delivered_offset = end;
|
|
}
|
|
if complete {
|
|
self.artifact_reverse_transfers.remove(&transfer_id);
|
|
self.artifact_transfer_by_token.remove(&token_digest);
|
|
}
|
|
return Ok(CoordinatorResponse::ArtifactDownloadStream {
|
|
link: stream.link,
|
|
streamed_bytes,
|
|
charged_download_bytes,
|
|
content_bytes_available: true,
|
|
content_offset: Some(content_offset),
|
|
content_eof: complete,
|
|
content_base64: Some(content_base64),
|
|
content_source: Some("retaining_node_reverse_stream".to_owned()),
|
|
});
|
|
}
|
|
|
|
let StorageLocation::RetainedNode(source_node) = &stream.link.source else {
|
|
return Err(clusterflux_core::DownloadError::Unavailable.into());
|
|
};
|
|
let metadata = self
|
|
.artifact_registry
|
|
.metadata(&artifact)
|
|
.ok_or(clusterflux_core::DownloadError::NotFound)?;
|
|
let transfer_id = generate_opaque_token("artifact_transfer")
|
|
.map_err(CoordinatorServiceError::Protocol)?;
|
|
let spool = tempfile::Builder::new()
|
|
.prefix("clusterflux-artifact-transfer-")
|
|
.tempfile()
|
|
.map_err(|error| {
|
|
CoordinatorServiceError::Protocol(format!(
|
|
"create bounded artifact transfer spool: {error}"
|
|
))
|
|
})?;
|
|
let transfer = ArtifactReverseTransfer {
|
|
transfer_id: transfer_id.clone(),
|
|
token_digest: token_digest.clone(),
|
|
tenant: context.tenant.clone(),
|
|
project: context.project.clone(),
|
|
source_node: source_node.clone(),
|
|
artifact,
|
|
expected_digest: metadata.digest.clone(),
|
|
expected_size_bytes: metadata.size,
|
|
expires_at_epoch_seconds: stream.link.expires_at_epoch_seconds,
|
|
spool,
|
|
received_bytes: 0,
|
|
content_hasher: Sha256::new(),
|
|
delivered_offset: 0,
|
|
error: None,
|
|
};
|
|
let overhead = self.artifact_relay.framing_overhead_bytes();
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.charge_egress(
|
|
stream.link.scoped_token_digest.as_str(),
|
|
overhead,
|
|
now_epoch_seconds,
|
|
)
|
|
})?;
|
|
self.artifact_reverse_transfers
|
|
.insert(transfer_id.clone(), transfer);
|
|
self.artifact_transfer_by_token
|
|
.insert(token_digest, transfer_id);
|
|
Ok(CoordinatorResponse::ArtifactDownloadStream {
|
|
link: stream.link,
|
|
streamed_bytes: 0,
|
|
charged_download_bytes: self.quota.used_download_bytes(
|
|
&context.tenant,
|
|
&context.project,
|
|
now_epoch_seconds,
|
|
),
|
|
content_bytes_available: false,
|
|
content_offset: None,
|
|
content_eof: false,
|
|
content_base64: None,
|
|
content_source: Some("retaining_node_reverse_stream_pending".to_owned()),
|
|
})
|
|
}
|
|
|
|
pub(super) fn handle_revoke_artifact_download_link(
|
|
&mut self,
|
|
tenant: String,
|
|
project: String,
|
|
actor_user: String,
|
|
artifact: String,
|
|
token_digest: Digest,
|
|
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
|
|
let context = user_context(tenant, project, actor_user);
|
|
let now_epoch_seconds = self.current_epoch_seconds()?;
|
|
self.artifact_registry
|
|
.expire_download_links(now_epoch_seconds);
|
|
let link = self.artifact_registry.revoke_download_link(
|
|
&context,
|
|
&ArtifactId::new(artifact),
|
|
&token_digest,
|
|
)?;
|
|
if let Some(transfer_id) = self.artifact_transfer_by_token.remove(&token_digest) {
|
|
self.artifact_reverse_transfers.remove(&transfer_id);
|
|
}
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.finish(token_digest.as_str(), RelayFinishReason::Cancelled);
|
|
Ok(())
|
|
})?;
|
|
Ok(CoordinatorResponse::ArtifactDownloadLinkRevoked { link })
|
|
}
|
|
|
|
pub(super) fn handle_export_artifact_to_node(
|
|
&mut self,
|
|
tenant: String,
|
|
project: String,
|
|
actor_user: String,
|
|
artifact: String,
|
|
receiver_node: String,
|
|
direct_connectivity: bool,
|
|
failure_reason: String,
|
|
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
|
|
let context = user_context(tenant, project, actor_user);
|
|
let artifact = ArtifactId::new(artifact);
|
|
let receiver_node = NodeId::new(receiver_node);
|
|
let action = self.artifact_registry.download_action(
|
|
&context,
|
|
&artifact,
|
|
&DownloadPolicy {
|
|
max_bytes: u64::MAX,
|
|
},
|
|
)?;
|
|
let StorageLocation::RetainedNode(source_node) = action.source else {
|
|
return Err(clusterflux_core::DownloadError::Unavailable.into());
|
|
};
|
|
let metadata = self
|
|
.artifact_registry
|
|
.metadata(&artifact)
|
|
.ok_or(clusterflux_core::DownloadError::NotFound)?;
|
|
let source = self.export_endpoint(&source_node, &context.tenant, &context.project)?;
|
|
let destination =
|
|
self.export_endpoint(&receiver_node, &context.tenant, &context.project)?;
|
|
let plan = self.transport.plan_authenticated_direct_bulk_transfer(
|
|
RendezvousRequest {
|
|
scope: DataPlaneScope {
|
|
tenant: context.tenant.clone(),
|
|
project: context.project.clone(),
|
|
process: metadata.process.clone(),
|
|
object: DataPlaneObject::Artifact(artifact.clone()),
|
|
authorization_subject: format!(
|
|
"artifact-export:{}-to-{}",
|
|
source_node, receiver_node
|
|
),
|
|
},
|
|
source,
|
|
destination,
|
|
},
|
|
direct_connectivity,
|
|
failure_reason,
|
|
)?;
|
|
Ok(CoordinatorResponse::ArtifactExportPlan {
|
|
plan,
|
|
source_node,
|
|
receiver_node,
|
|
artifact_size_bytes: metadata.size,
|
|
})
|
|
}
|
|
|
|
pub(super) fn handle_poll_artifact_transfer(
|
|
&mut self,
|
|
tenant: String,
|
|
project: String,
|
|
node: String,
|
|
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
|
|
let tenant = TenantId::new(tenant);
|
|
let project = ProjectId::new(project);
|
|
let node = NodeId::new(node);
|
|
self.authorize_artifact_transfer_node(&tenant, &project, &node)?;
|
|
let transfer = self.artifact_reverse_transfers.values().find(|transfer| {
|
|
transfer.tenant == tenant
|
|
&& transfer.project == project
|
|
&& transfer.source_node == node
|
|
&& transfer.error.is_none()
|
|
&& transfer.received_bytes < transfer.expected_size_bytes
|
|
});
|
|
Ok(CoordinatorResponse::ArtifactTransferAssignment {
|
|
transfer: transfer.map(|transfer| ArtifactTransferAssignment {
|
|
transfer_id: transfer.transfer_id.clone(),
|
|
artifact: transfer.artifact.clone(),
|
|
expected_digest: transfer.expected_digest.clone(),
|
|
expected_size_bytes: transfer.expected_size_bytes,
|
|
offset: transfer.received_bytes,
|
|
max_chunk_bytes: MAX_ARTIFACT_REVERSE_CHUNK_BYTES,
|
|
}),
|
|
})
|
|
}
|
|
|
|
#[allow(clippy::too_many_arguments)]
|
|
pub(super) fn handle_upload_artifact_transfer_chunk(
|
|
&mut self,
|
|
tenant: String,
|
|
project: String,
|
|
node: String,
|
|
transfer_id: String,
|
|
artifact: String,
|
|
offset: u64,
|
|
content_base64: String,
|
|
chunk_digest: Digest,
|
|
eof: bool,
|
|
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
|
|
let tenant = TenantId::new(tenant);
|
|
let project = ProjectId::new(project);
|
|
let node = NodeId::new(node);
|
|
self.authorize_artifact_transfer_node(&tenant, &project, &node)?;
|
|
let token_digest = {
|
|
let transfer = self
|
|
.artifact_reverse_transfers
|
|
.get(&transfer_id)
|
|
.ok_or_else(|| {
|
|
CoordinatorServiceError::Protocol(
|
|
"unknown artifact reverse transfer".to_owned(),
|
|
)
|
|
})?;
|
|
if transfer.tenant != tenant
|
|
|| transfer.project != project
|
|
|| transfer.source_node != node
|
|
|| transfer.artifact.as_str() != artifact
|
|
{
|
|
return Err(CoordinatorError::Unauthorized(
|
|
"artifact reverse transfer is outside the signed node scope".to_owned(),
|
|
)
|
|
.into());
|
|
}
|
|
transfer.token_digest.clone()
|
|
};
|
|
let now_epoch_seconds = self.current_epoch_seconds()?;
|
|
let ingress_wire_bytes = (content_base64.len() as u64)
|
|
.saturating_add(self.artifact_relay.framing_overhead_bytes());
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.charge_ingress(token_digest.as_str(), ingress_wire_bytes, now_epoch_seconds)
|
|
})?;
|
|
let transfer = self
|
|
.artifact_reverse_transfers
|
|
.get_mut(&transfer_id)
|
|
.ok_or_else(|| {
|
|
CoordinatorServiceError::Protocol("unknown artifact reverse transfer".to_owned())
|
|
})?;
|
|
if transfer.tenant != tenant
|
|
|| transfer.project != project
|
|
|| transfer.source_node != node
|
|
|| transfer.artifact.as_str() != artifact
|
|
{
|
|
return Err(CoordinatorError::Unauthorized(
|
|
"artifact reverse transfer is outside the signed node scope".to_owned(),
|
|
)
|
|
.into());
|
|
}
|
|
if transfer.received_bytes != offset {
|
|
return Err(CoordinatorServiceError::Protocol(format!(
|
|
"artifact reverse transfer expected offset {}, received {offset}",
|
|
transfer.received_bytes
|
|
)));
|
|
}
|
|
let content = BASE64_STANDARD.decode(content_base64).map_err(|error| {
|
|
CoordinatorServiceError::Protocol(format!(
|
|
"artifact reverse transfer chunk is not valid base64: {error}"
|
|
))
|
|
})?;
|
|
if content.len() as u64 > MAX_ARTIFACT_REVERSE_CHUNK_BYTES {
|
|
return Err(CoordinatorServiceError::Protocol(
|
|
"artifact reverse transfer chunk exceeds the control-plane limit".to_owned(),
|
|
));
|
|
}
|
|
if Digest::sha256(&content) != chunk_digest {
|
|
return Err(CoordinatorServiceError::Protocol(
|
|
"artifact reverse transfer chunk digest mismatch".to_owned(),
|
|
));
|
|
}
|
|
let next_offset = offset.saturating_add(content.len() as u64);
|
|
if next_offset > transfer.expected_size_bytes {
|
|
return Err(CoordinatorServiceError::Protocol(
|
|
"artifact reverse transfer exceeds retained metadata size".to_owned(),
|
|
));
|
|
}
|
|
let complete = next_offset == transfer.expected_size_bytes;
|
|
if eof != complete {
|
|
return Err(CoordinatorServiceError::Protocol(
|
|
"artifact reverse transfer EOF does not match retained metadata size".to_owned(),
|
|
));
|
|
}
|
|
transfer
|
|
.spool
|
|
.as_file_mut()
|
|
.seek(SeekFrom::Start(offset))
|
|
.and_then(|_| transfer.spool.as_file_mut().write_all(&content))
|
|
.map_err(|error| {
|
|
CoordinatorServiceError::Protocol(format!(
|
|
"write bounded artifact transfer spool: {error}"
|
|
))
|
|
})?;
|
|
transfer.content_hasher.update(&content);
|
|
transfer.received_bytes = next_offset;
|
|
if complete {
|
|
transfer.spool.as_file().sync_all().map_err(|error| {
|
|
CoordinatorServiceError::Protocol(format!(
|
|
"sync bounded artifact transfer spool: {error}"
|
|
))
|
|
})?;
|
|
let digest_hex = format!("{:x}", transfer.content_hasher.clone().finalize());
|
|
let digest =
|
|
Digest::from_sha256_hex(&digest_hex).map_err(CoordinatorServiceError::Protocol)?;
|
|
if digest != transfer.expected_digest {
|
|
transfer.spool.as_file_mut().set_len(0).map_err(|error| {
|
|
CoordinatorServiceError::Protocol(format!(
|
|
"reset invalid artifact transfer spool: {error}"
|
|
))
|
|
})?;
|
|
transfer.received_bytes = 0;
|
|
transfer.content_hasher = Sha256::new();
|
|
return Err(CoordinatorServiceError::Protocol(
|
|
"artifact reverse transfer content digest mismatch".to_owned(),
|
|
));
|
|
}
|
|
}
|
|
Ok(CoordinatorResponse::ArtifactTransferChunkAccepted {
|
|
transfer_id,
|
|
next_offset,
|
|
complete,
|
|
})
|
|
}
|
|
|
|
#[allow(clippy::too_many_arguments)]
|
|
pub(super) fn handle_fail_artifact_transfer(
|
|
&mut self,
|
|
tenant: String,
|
|
project: String,
|
|
node: String,
|
|
transfer_id: String,
|
|
artifact: String,
|
|
message: String,
|
|
) -> Result<CoordinatorResponse, CoordinatorServiceError> {
|
|
let tenant = TenantId::new(tenant);
|
|
let project = ProjectId::new(project);
|
|
let node = NodeId::new(node);
|
|
self.authorize_artifact_transfer_node(&tenant, &project, &node)?;
|
|
let token_digest = {
|
|
let transfer = self
|
|
.artifact_reverse_transfers
|
|
.get(&transfer_id)
|
|
.ok_or_else(|| {
|
|
CoordinatorServiceError::Protocol(
|
|
"unknown artifact reverse transfer".to_owned(),
|
|
)
|
|
})?;
|
|
if transfer.tenant != tenant
|
|
|| transfer.project != project
|
|
|| transfer.source_node != node
|
|
|| transfer.artifact.as_str() != artifact
|
|
{
|
|
return Err(CoordinatorError::Unauthorized(
|
|
"artifact reverse transfer failure is outside the signed node scope".to_owned(),
|
|
)
|
|
.into());
|
|
}
|
|
transfer.token_digest.clone()
|
|
};
|
|
let now_epoch_seconds = self.current_epoch_seconds()?;
|
|
let failure_wire_bytes =
|
|
(message.len() as u64).saturating_add(self.artifact_relay.framing_overhead_bytes());
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.charge_ingress(token_digest.as_str(), failure_wire_bytes, now_epoch_seconds)
|
|
})?;
|
|
let transfer = self
|
|
.artifact_reverse_transfers
|
|
.get_mut(&transfer_id)
|
|
.ok_or_else(|| {
|
|
CoordinatorServiceError::Protocol("unknown artifact reverse transfer".to_owned())
|
|
})?;
|
|
if transfer.tenant != tenant
|
|
|| transfer.project != project
|
|
|| transfer.source_node != node
|
|
|| transfer.artifact.as_str() != artifact
|
|
{
|
|
return Err(CoordinatorError::Unauthorized(
|
|
"artifact reverse transfer failure is outside the signed node scope".to_owned(),
|
|
)
|
|
.into());
|
|
}
|
|
let message = message.trim();
|
|
transfer.error = Some(if message.is_empty() {
|
|
"retained artifact bytes are unavailable".to_owned()
|
|
} else {
|
|
message.chars().take(1024).collect()
|
|
});
|
|
self.mutate_artifact_relay(|ledger| {
|
|
ledger.finish(token_digest.as_str(), RelayFinishReason::Failed);
|
|
Ok(())
|
|
})?;
|
|
Ok(CoordinatorResponse::ArtifactTransferFailed { transfer_id })
|
|
}
|
|
|
|
fn authorize_artifact_transfer_node(
|
|
&self,
|
|
tenant: &TenantId,
|
|
project: &ProjectId,
|
|
node: &NodeId,
|
|
) -> Result<(), CoordinatorServiceError> {
|
|
let identity = self
|
|
.coordinator
|
|
.node_identity(node)
|
|
.ok_or(CoordinatorError::UnknownNode)?;
|
|
if &identity.tenant != tenant || &identity.project != project {
|
|
return Err(CoordinatorError::Unauthorized(
|
|
"artifact reverse transfer node is outside its enrolled tenant/project scope"
|
|
.to_owned(),
|
|
)
|
|
.into());
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn expire_artifact_reverse_transfers(
|
|
&mut self,
|
|
now_epoch_seconds: u64,
|
|
) -> Result<(), CoordinatorServiceError> {
|
|
let expired = self
|
|
.artifact_reverse_transfers
|
|
.iter()
|
|
.filter(|(_, transfer)| transfer.expires_at_epoch_seconds < now_epoch_seconds)
|
|
.map(|(id, transfer)| (id.clone(), transfer.token_digest.clone()))
|
|
.collect::<Vec<_>>();
|
|
for (id, token) in &expired {
|
|
self.artifact_reverse_transfers.remove(id);
|
|
self.artifact_transfer_by_token.remove(token);
|
|
}
|
|
self.mutate_artifact_relay(|ledger| {
|
|
for (_, token) in &expired {
|
|
ledger.finish(token.as_str(), RelayFinishReason::Expired);
|
|
}
|
|
ledger.expire(now_epoch_seconds);
|
|
Ok(())
|
|
})
|
|
}
|
|
|
|
fn export_endpoint(
|
|
&self,
|
|
node: &NodeId,
|
|
tenant: &TenantId,
|
|
project: &ProjectId,
|
|
) -> Result<NodeEndpoint, CoordinatorServiceError> {
|
|
let identity = self
|
|
.coordinator
|
|
.node_identity(node)
|
|
.ok_or(CoordinatorError::UnknownNode)?;
|
|
if &identity.tenant != tenant || &identity.project != project {
|
|
return Err(CoordinatorError::Unauthorized(
|
|
"artifact export node is outside the tenant/project scope".to_owned(),
|
|
)
|
|
.into());
|
|
}
|
|
let descriptor = self.node_descriptors.get(node).ok_or_else(|| {
|
|
clusterflux_core::DownloadError::DirectConnectivityUnavailable(format!(
|
|
"node {node} has not reported export connectivity"
|
|
))
|
|
})?;
|
|
if descriptor.tenant != *tenant || descriptor.project != *project {
|
|
return Err(CoordinatorError::Unauthorized(
|
|
"artifact export node descriptor is outside the tenant/project scope".to_owned(),
|
|
)
|
|
.into());
|
|
}
|
|
if !self.node_is_live(node) {
|
|
return Err(
|
|
clusterflux_core::DownloadError::DirectConnectivityUnavailable(format!(
|
|
"node {node} is offline for artifact export"
|
|
))
|
|
.into(),
|
|
);
|
|
}
|
|
if !descriptor.direct_connectivity {
|
|
return Err(
|
|
clusterflux_core::DownloadError::DirectConnectivityUnavailable(format!(
|
|
"direct connectivity unavailable to node {node} for artifact export"
|
|
))
|
|
.into(),
|
|
);
|
|
}
|
|
Ok(NodeEndpoint {
|
|
node: node.clone(),
|
|
advertised_addr: format!("{node}.mesh.invalid:4433"),
|
|
public_key_fingerprint: Digest::sha256(&identity.public_key),
|
|
})
|
|
}
|
|
|
|
fn ensure_download_source_connectivity(
|
|
&self,
|
|
source: &StorageLocation,
|
|
) -> Result<(), clusterflux_core::DownloadError> {
|
|
let StorageLocation::RetainedNode(node) = source else {
|
|
return Ok(());
|
|
};
|
|
let _descriptor = self.node_descriptors.get(node).ok_or_else(|| {
|
|
clusterflux_core::DownloadError::DirectConnectivityUnavailable(format!(
|
|
"retaining node {node} has not reported online status for artifact download"
|
|
))
|
|
})?;
|
|
if !self.node_is_live(node) {
|
|
return Err(
|
|
clusterflux_core::DownloadError::DirectConnectivityUnavailable(format!(
|
|
"retaining node {node} is offline for artifact download"
|
|
)),
|
|
);
|
|
}
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
fn user_context(tenant: String, project: String, actor_user: String) -> AuthContext {
|
|
AuthContext {
|
|
tenant: TenantId::new(tenant),
|
|
project: ProjectId::new(project),
|
|
actor: Actor::User(UserId::new(actor_user)),
|
|
}
|
|
}
|