Public release release-f6c01dca61ee
Source commit: f6c01dca61eefca1a7431c48878a04dbb7a3f710 Public tree identity: sha256:3df7a93eb61c8c5c7841ac9ad6eedec89b167fd77fdd32e9fc31b57c10207387
This commit is contained in:
commit
a2664f4345
210 changed files with 78559 additions and 0 deletions
799
crates/clusterflux-coordinator/src/service/artifacts.rs
Normal file
799
crates/clusterflux-coordinator/src/service/artifacts.rs
Normal file
|
|
@ -0,0 +1,799 @@
|
|||
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)),
|
||||
}
|
||||
}
|
||||
Loading…
Add table
Add a link
Reference in a new issue