feat: add preferences and send empty state

This commit is contained in:
2026-07-08 00:06:00 +02:00
parent 9a7abaca00
commit 72340c6c84
27 changed files with 1436 additions and 78 deletions

View File

@@ -20,6 +20,17 @@ pub trait CoreEventSink: Send + Sync {
fn on_event(&self, event: CoreEvent);
}
#[uniffi::export(with_foreign)]
pub trait ReceiveOutputSink: Send + Sync {
fn start_file(&self, relative_path: String) -> Result<(), crate::error::VnidropError>;
fn write_chunk(
&self,
relative_path: String,
bytes: Vec<u8>,
) -> Result<(), crate::error::VnidropError>;
fn finish_file(&self, relative_path: String) -> Result<(), crate::error::VnidropError>;
}
#[derive(Debug, Clone, Serialize, Deserialize, uniffi::Record)]
pub struct RuntimeStatus {
pub endpoint_id: String,

View File

@@ -13,9 +13,9 @@ mod ticket;
mod util;
pub use api::{
CoreEvent, CoreEventSink, ReceiverRequest, RuntimeStatus, ShareMetadataInput, ShareResult,
ShareSource, SourceKind, StoredTransfer, TicketInspection, TransferAccessMode,
TransferMetadata,
CoreEvent, CoreEventSink, ReceiveOutputSink, ReceiverRequest, RuntimeStatus,
ShareMetadataInput, ShareResult, ShareSource, SourceKind, StoredTransfer, TicketInspection,
TransferAccessMode, TransferMetadata,
};
pub use error::VnidropError;
pub use runtime::VnidropCore;

View File

@@ -32,16 +32,17 @@ use tokio::{
use crate::{
access_policy::{AccessDecision, AccessPolicy},
api::{
CoreEvent, CoreEventSink, ReceiverRequest, RuntimeStatus, ShareMetadataInput, ShareResult,
ShareSource, StoredTransfer, TicketInspection, TransferAccessMode, TransferMetadata,
CoreEvent, CoreEventSink, ReceiveOutputSink, ReceiverRequest, RuntimeStatus,
ShareMetadataInput, ShareResult, ShareSource, StoredTransfer, TicketInspection,
TransferAccessMode, TransferMetadata,
},
approval::ApprovalService,
error::VnidropError,
event_hub::EventHub,
filesystem::{
collect_import_files, default_collection_name, platform_path,
read_stream_from_blocking_reader, safe_output_path, wait_for_writer,
write_stream_to_blocking_writer, TransferImport,
read_stream_from_blocking_reader, safe_output_path, validated_relative_string,
wait_for_writer, write_stream_to_blocking_writer, TransferImport,
},
handshake::{HandshakeResponse, HandshakeService},
logging::init_logging,
@@ -82,6 +83,11 @@ struct CoreInner {
shutdown_started: AtomicBool,
}
enum ReceiveTarget {
Directory(PathBuf),
OutputSink(Arc<dyn ReceiveOutputSink>),
}
#[uniffi::export]
impl VnidropCore {
#[uniffi::constructor]
@@ -139,6 +145,33 @@ impl VnidropCore {
.map_err(VnidropError::transfer)
}
pub fn receive_with_output_sink(
&self,
ticket: String,
output_sink: Arc<dyn ReceiveOutputSink>,
receiver_name: Option<String>,
) -> Result<(), VnidropError> {
if let Err(error) =
parse_transfer_ticket(&ticket).context("failed to parse transfer ticket")
{
self.runtime.block_on(async {
self.inner.emit_endpoint(
"error",
"invalid-ticket",
json!({ "reason": error.to_string() }),
);
self.inner.event_hub.flush().await;
});
return Err(VnidropError::ticket(error));
}
self.runtime
.block_on(
self.inner
.receive_with_output_sink(ticket, output_sink, receiver_name),
)
.map_err(VnidropError::transfer)
}
pub fn cancel_transfer(&self, transfer_id: u64) -> Result<(), VnidropError> {
self.runtime
.block_on(self.inner.cancel_transfer(transfer_id))
@@ -397,6 +430,30 @@ impl CoreInner {
ticket: String,
output_dir: PathBuf,
receiver_name: Option<String>,
) -> Result<()> {
self.receive_to_target(ticket, ReceiveTarget::Directory(output_dir), receiver_name)
.await
}
async fn receive_with_output_sink(
self: &Arc<Self>,
ticket: String,
output_sink: Arc<dyn ReceiveOutputSink>,
receiver_name: Option<String>,
) -> Result<()> {
self.receive_to_target(
ticket,
ReceiveTarget::OutputSink(output_sink),
receiver_name,
)
.await
}
async fn receive_to_target(
self: &Arc<Self>,
ticket: String,
target: ReceiveTarget,
receiver_name: Option<String>,
) -> Result<()> {
let parsed = match parse_transfer_ticket(&ticket).context("failed to parse transfer ticket")
{
@@ -424,7 +481,7 @@ impl CoreInner {
.insert(transfer_id, shutdown_tx);
let result = tokio::select! {
result = self.receive_inner(transfer_id, parsed, output_dir, receiver_name) => result,
result = self.receive_inner(transfer_id, parsed, target, receiver_name) => result,
_ = &mut shutdown_rx => Err(anyhow::anyhow!("transfer cancelled")),
};
@@ -449,7 +506,7 @@ impl CoreInner {
self: &Arc<Self>,
transfer_id: u64,
parsed: ParsedTransferTicket,
output_dir: PathBuf,
target: ReceiveTarget,
receiver_name: Option<String>,
) -> Result<()> {
let metadata_json =
@@ -490,7 +547,9 @@ impl CoreInner {
.unwrap_or_default(),
})
.await?;
tokio::fs::create_dir_all(&output_dir).await?;
if let ReceiveTarget::Directory(output_dir) = &target {
tokio::fs::create_dir_all(output_dir).await?;
}
self.emit_transfer(transfer_id, "receive", "network", "connecting", json!({}));
if let Some(metadata) = &parsed.metadata {
@@ -542,7 +601,7 @@ impl CoreInner {
}
let collection = Collection::load(hash_and_format.hash, self.store.as_ref()).await?;
self.export_collection(transfer_id, total_files, output_dir, collection)
self.export_collection(transfer_id, total_files, target, collection)
.await?;
self.repository
.update_transfer_status(transfer_id, STATUS_DONE)
@@ -795,24 +854,39 @@ impl CoreInner {
&self,
transfer_id: u64,
total_files: u64,
output_dir: PathBuf,
target: ReceiveTarget,
collection: Collection,
) -> Result<()> {
for (i, (name, hash)) in collection.iter().enumerate() {
self.export_blob(
transfer_id,
total_files,
i as u64,
&output_dir,
name.as_ref(),
*hash,
)
.await?;
match &target {
ReceiveTarget::Directory(output_dir) => {
self.export_blob_to_directory(
transfer_id,
total_files,
i as u64,
output_dir,
name.as_ref(),
*hash,
)
.await?;
}
ReceiveTarget::OutputSink(output_sink) => {
self.export_blob_to_sink(
transfer_id,
total_files,
i as u64,
output_sink.as_ref(),
name.as_ref(),
*hash,
)
.await?;
}
}
}
Ok(())
}
async fn export_blob(
async fn export_blob_to_directory(
&self,
transfer_id: u64,
total_files: u64,
@@ -873,6 +947,62 @@ impl CoreInner {
Ok(())
}
async fn export_blob_to_sink(
&self,
transfer_id: u64,
total_files: u64,
current_file_index: u64,
output_sink: &dyn ReceiveOutputSink,
relative_path: &str,
hash: Hash,
) -> Result<()> {
let relative_path = validated_relative_string(relative_path)?;
output_sink
.start_file(relative_path.clone())
.map_err(|error| anyhow::anyhow!(error.to_string()))?;
let mut stream = self.store.export_ranges(hash, 0..u64::MAX).stream();
let mut file_size = 0;
let mut exported = 0;
while let Some(item) = stream.next().await {
match item {
ExportRangesItem::Size(size) => file_size = size,
ExportRangesItem::Data(leaf) => {
if leaf.offset != exported {
anyhow::bail!(
"export stream for {relative_path} yielded out-of-order data"
);
}
exported += leaf.data.len() as u64;
output_sink
.write_chunk(relative_path.clone(), leaf.data.to_vec())
.map_err(|error| anyhow::anyhow!(error.to_string()))?;
self.emit_transfer(
transfer_id,
"receive",
"export",
"progress",
json!({
"total_files": total_files,
"current_file_index": current_file_index,
"file_name": relative_path,
"file_size": file_size,
"exported": exported,
}),
);
}
ExportRangesItem::Error(error) => {
anyhow::bail!("export failed for {relative_path}: {error}");
}
}
}
output_sink
.finish_file(relative_path)
.map_err(|error| anyhow::anyhow!(error.to_string()))?;
Ok(())
}
async fn spawn_provider_event_task(self: &Arc<Self>, mut rx: mpsc::Receiver<ProviderMessage>) {
let core = self.clone();
let task = tokio::spawn(async move {

View File

@@ -1,11 +1,12 @@
use std::{
collections::HashMap,
sync::{Arc, Mutex},
time::{Duration, Instant},
};
use vnidrop::{
CoreEvent, CoreEventSink, ReceiverRequest, ShareMetadataInput, ShareSource, SourceKind,
VnidropCore,
CoreEvent, CoreEventSink, ReceiveOutputSink, ReceiverRequest, ShareMetadataInput, ShareSource,
SourceKind, VnidropCore, VnidropError,
};
#[derive(Default)]
@@ -25,6 +26,44 @@ impl RecordingSink {
}
}
#[derive(Default)]
struct MemoryOutputSink {
files: Mutex<HashMap<String, Vec<u8>>>,
fail_writes: bool,
}
impl ReceiveOutputSink for MemoryOutputSink {
fn start_file(&self, relative_path: String) -> Result<(), VnidropError> {
self.files.lock().unwrap().insert(relative_path, Vec::new());
Ok(())
}
fn write_chunk(&self, relative_path: String, bytes: Vec<u8>) -> Result<(), VnidropError> {
if self.fail_writes {
return Err(VnidropError::Filesystem {
reason: "sink write failed".to_string(),
});
}
self.files
.lock()
.unwrap()
.get_mut(&relative_path)
.expect("file was not started")
.extend(bytes);
Ok(())
}
fn finish_file(&self, _relative_path: String) -> Result<(), VnidropError> {
Ok(())
}
}
impl MemoryOutputSink {
fn file(&self, relative_path: &str) -> Vec<u8> {
self.files.lock().unwrap()[relative_path].clone()
}
}
fn wait_for_receiver_request(sender: &VnidropCore, transfer_id: u64) -> ReceiverRequest {
let started = Instant::now();
loop {
@@ -68,6 +107,31 @@ fn receive_with_response(
handle.join().unwrap()
}
fn receive_with_sink_response(
sender: &VnidropCore,
transfer_id: u64,
receiver: Arc<VnidropCore>,
ticket: String,
output_sink: Arc<dyn ReceiveOutputSink>,
receiver_name: Option<String>,
accepted: bool,
) -> Result<(), String> {
let handle = std::thread::spawn(move || {
receiver
.receive_with_output_sink(ticket, output_sink, receiver_name)
.map_err(|error| error.to_string())
});
let request = wait_for_receiver_request(sender, transfer_id);
sender
.respond_receiver_request(
request.id,
accepted,
(!accepted).then(|| "sender-refused".to_string()),
)
.unwrap();
handle.join().unwrap()
}
#[test]
fn two_local_cores_transfer_file() {
let sender_dir = tempfile::tempdir().unwrap();
@@ -201,6 +265,124 @@ fn two_local_cores_transfer_directory() {
receiver.shutdown();
}
#[test]
fn output_sink_receive_exports_nested_files() {
let sender_dir = tempfile::tempdir().unwrap();
let receiver_dir = tempfile::tempdir().unwrap();
let source_root = sender_dir.path().join("photos");
std::fs::create_dir_all(source_root.join("nested")).unwrap();
std::fs::write(source_root.join("cover.txt"), b"cover").unwrap();
std::fs::write(source_root.join("nested").join("inside.txt"), b"inside").unwrap();
let sender = VnidropCore::initialize(
sender_dir.path().join("core").to_string_lossy().to_string(),
Arc::new(RecordingSink::default()),
)
.unwrap();
let receiver = VnidropCore::initialize(
receiver_dir
.path()
.join("core")
.to_string_lossy()
.to_string(),
Arc::new(RecordingSink::default()),
)
.unwrap();
let share = sender
.share_files(
vec![ShareSource {
kind: SourceKind::Path,
value: source_root.to_string_lossy().to_string(),
display_name: Some("photos".to_string()),
is_directory: true,
}],
ShareMetadataInput {
transfer_id: 18,
transfer_name: Some("photos".to_string()),
sender_name: Some("sender".to_string()),
},
)
.unwrap();
let output_sink = Arc::new(MemoryOutputSink::default());
receive_with_sink_response(
&sender,
share.transfer_id,
receiver.clone(),
share.ticket,
output_sink.clone(),
Some("receiver".to_string()),
true,
)
.unwrap();
assert_eq!(output_sink.file("photos/cover.txt"), b"cover");
assert_eq!(output_sink.file("photos/nested/inside.txt"), b"inside");
sender.shutdown();
receiver.shutdown();
}
#[test]
fn output_sink_receive_fails_when_sink_write_fails() {
let sender_dir = tempfile::tempdir().unwrap();
let receiver_dir = tempfile::tempdir().unwrap();
let source_path = sender_dir.path().join("hello.txt");
std::fs::write(&source_path, b"hello").unwrap();
let sender = VnidropCore::initialize(
sender_dir.path().join("core").to_string_lossy().to_string(),
Arc::new(RecordingSink::default()),
)
.unwrap();
let receiver = VnidropCore::initialize(
receiver_dir
.path()
.join("core")
.to_string_lossy()
.to_string(),
Arc::new(RecordingSink::default()),
)
.unwrap();
let share = sender
.share_files(
vec![ShareSource {
kind: SourceKind::Path,
value: source_path.to_string_lossy().to_string(),
display_name: Some("hello.txt".to_string()),
is_directory: false,
}],
ShareMetadataInput {
transfer_id: 19,
transfer_name: Some("hello".to_string()),
sender_name: Some("sender".to_string()),
},
)
.unwrap();
let output_sink = Arc::new(MemoryOutputSink {
files: Mutex::new(HashMap::new()),
fail_writes: true,
});
let error = receive_with_sink_response(
&sender,
share.transfer_id,
receiver.clone(),
share.ticket,
output_sink,
Some("receiver".to_string()),
true,
)
.unwrap_err();
assert!(error.contains("sink write failed"));
sender.shutdown();
receiver.shutdown();
}
#[test]
fn approval_required_denies_then_allows_receiver() {
let sender_dir = tempfile::tempdir().unwrap();