From mboxrd@z Thu Jan 1 00:00:00 1970 Return-Path: Received: from firstgate.proxmox.com (firstgate.proxmox.com [212.224.123.68]) (using TLSv1.3 with cipher TLS_AES_256_GCM_SHA384 (256/256 bits) key-exchange X25519 server-signature RSA-PSS (2048 bits)) (No client certificate requested) by lists.proxmox.com (Postfix) with ESMTPS id 3FE33DC89 for ; Fri, 22 Sep 2023 09:25:23 +0200 (CEST) Received: from firstgate.proxmox.com (localhost [127.0.0.1]) by firstgate.proxmox.com (Proxmox) with ESMTP id 1E54E71D1 for ; Fri, 22 Sep 2023 09:24:53 +0200 (CEST) Received: from proxmox-new.maurer-it.com (proxmox-new.maurer-it.com [94.136.29.106]) (using TLSv1.3 with cipher TLS_AES_256_GCM_SHA384 (256/256 bits) key-exchange X25519 server-signature RSA-PSS (2048 bits)) (No client certificate requested) by firstgate.proxmox.com (Proxmox) with ESMTPS for ; Fri, 22 Sep 2023 09:24:52 +0200 (CEST) Received: from proxmox-new.maurer-it.com (localhost.localdomain [127.0.0.1]) by proxmox-new.maurer-it.com (Proxmox) with ESMTP id 7C8FE48790 for ; Fri, 22 Sep 2023 09:16:52 +0200 (CEST) From: Christian Ebner To: pbs-devel@lists.proxmox.com Date: Fri, 22 Sep 2023 09:16:17 +0200 Message-Id: <20230922071621.12670-17-c.ebner@proxmox.com> X-Mailer: git-send-email 2.39.2 In-Reply-To: <20230922071621.12670-1-c.ebner@proxmox.com> References: <20230922071621.12670-1-c.ebner@proxmox.com> MIME-Version: 1.0 Content-Transfer-Encoding: 8bit X-SPAM-LEVEL: Spam detection results: 0 AWL 0.103 Adjusted score from AWL reputation of From: address BAYES_00 -1.9 Bayes spam probability is 0 to 1% DMARC_MISSING 0.1 Missing DMARC policy KAM_DMARC_STATUS 0.01 Test Rule for DKIM or SPF Failure with Strict Alignment SPF_HELO_NONE 0.001 SPF: HELO does not publish an SPF Record SPF_PASS -0.001 SPF: sender matches SPF record URIBL_BLOCKED 0.001 ADMINISTRATOR NOTICE: The query to URIBL was blocked. See http://wiki.apache.org/spamassassin/DnsBlocklists#dnsbl-block for more information. [lib.rs] Subject: [pbs-devel] [RFC proxmox-backup 16/20] fix #3174: upload stream: impl reused chunk injector X-BeenThere: pbs-devel@lists.proxmox.com X-Mailman-Version: 2.1.29 Precedence: list List-Id: Proxmox Backup Server development discussion List-Unsubscribe: , List-Archive: List-Post: List-Help: List-Subscribe: , X-List-Received-Date: Fri, 22 Sep 2023 07:25:23 -0000 In order to be included in the backups index file, the reused chunks which store the payload of skipped files during pxar encoding have to be inserted after the encoder has written the pxar appendix entry type. The chunker forces a chunk boundary after this marker and queues the list of chunks to be uploaded thereafter. This implements the logic to inject the chunks into the chunk upload stream after such a boundary is requested, by looping over the queued chunks and inserting them into the stream. Signed-off-by: Christian Ebner --- pbs-client/src/inject_reused_chunks.rs | 123 +++++++++++++++++++++++++ pbs-client/src/lib.rs | 1 + 2 files changed, 124 insertions(+) create mode 100644 pbs-client/src/inject_reused_chunks.rs diff --git a/pbs-client/src/inject_reused_chunks.rs b/pbs-client/src/inject_reused_chunks.rs new file mode 100644 index 00000000..01cb1350 --- /dev/null +++ b/pbs-client/src/inject_reused_chunks.rs @@ -0,0 +1,123 @@ +use std::collections::VecDeque; +use std::pin::Pin; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; +use std::task::{Context, Poll}; + +use anyhow::Error; +use futures::{ready, Stream}; +use pin_project_lite::pin_project; + +use pbs_datastore::dynamic_index::DynamicEntry; + +pin_project! { + pub struct InjectReusedChunksQueue { + #[pin] + input: S, + current: Option, + injection_queue: Arc>>, + stream_len: Arc, + index_csum: Arc>>, + } +} + +#[derive(Debug)] +pub struct InjectChunks { + pub boundary: u64, + pub chunks: Vec, + pub size: usize, +} + +pub enum InjectedChunksInfo { + Known(Vec<(u64, [u8; 32])>), + Raw((u64, bytes::BytesMut)), +} + +pub trait InjectReusedChunks: Sized { + fn inject_reused_chunks( + self, + injection_queue: Arc>>, + stream_len: Arc, + index_csum: Arc>>, + ) -> InjectReusedChunksQueue; +} + +impl InjectReusedChunks for S +where + S: Stream>, +{ + fn inject_reused_chunks( + self, + injection_queue: Arc>>, + stream_len: Arc, + index_csum: Arc>>, + ) -> InjectReusedChunksQueue { + let current = injection_queue.lock().unwrap().pop_front(); + + InjectReusedChunksQueue { + input: self, + current, + injection_queue, + stream_len, + index_csum, + } + } +} + +impl Stream for InjectReusedChunksQueue +where + S: Stream>, +{ + type Item = Result; + + fn poll_next(self: Pin<&mut Self>, cx: &mut Context) -> Poll> { + let mut this = self.project(); + loop { + let current = this.current.take(); + if let Some(current) = current { + let mut chunks = Vec::new(); + let mut guard = this.index_csum.lock().unwrap(); + let csum = guard.as_mut().unwrap(); + + for chunk in current.chunks { + let offset = this + .stream_len + .fetch_add(chunk.end() as usize, Ordering::SeqCst) + as u64; + let digest = chunk.digest(); + chunks.push((offset, digest)); + // Chunk end is assumed to be normalized to chunk size here + let end_offset = offset + chunk.end(); + csum.update(&end_offset.to_le_bytes()); + csum.update(&digest); + } + let chunk_info = InjectedChunksInfo::Known(chunks); + return Poll::Ready(Some(Ok(chunk_info))); + } + + match ready!(this.input.as_mut().poll_next(cx)) { + None => return Poll::Ready(None), + Some(Err(err)) => return Poll::Ready(Some(Err(err))), + Some(Ok(raw)) => { + let chunk_size = raw.len(); + let offset = this.stream_len.fetch_add(chunk_size, Ordering::SeqCst) as u64; + let mut injections = this.injection_queue.lock().unwrap(); + if let Some(inject) = injections.pop_front() { + if inject.boundary == offset { + let _ = this.current.insert(inject); + // Should be injected here, directly jump to next loop iteration + continue; + } else if inject.boundary <= offset + chunk_size as u64 { + let _ = this.current.insert(inject); + } else { + injections.push_front(inject); + } + } + let data = InjectedChunksInfo::Raw((offset, raw)); + + return Poll::Ready(Some(Ok(data))); + } + } + } + } +} diff --git a/pbs-client/src/lib.rs b/pbs-client/src/lib.rs index 21cf8556..8bf26381 100644 --- a/pbs-client/src/lib.rs +++ b/pbs-client/src/lib.rs @@ -8,6 +8,7 @@ pub mod pxar; pub mod tools; mod merge_known_chunks; +mod inject_reused_chunks; pub mod pipe_to_stream; mod http_client; -- 2.39.2