From 906e26b73ab1f93df8cdfc438dc837422dad1f40 Mon Sep 17 00:00:00 2001 From: kiyo-e Date: Wed, 27 May 2026 23:33:29 +0900 Subject: [PATCH] Prepare Pairlane 0.1.2 transfer checksum fix --- README.ja.md | 2 +- README.md | 2 +- cli/Cargo.lock | 2 +- cli/Cargo.toml | 2 +- cli/npm/package.json | 4 +- cli/src/main.rs | 173 +++++++++++++++++++++++++++++-------- docs/signaling-protocol.md | 29 ++++++- src/client/room.tsx | 43 ++++++--- 8 files changed, 200 insertions(+), 57 deletions(-) diff --git a/README.ja.md b/README.ja.md index 7d113e2..44bc0c8 100644 --- a/README.ja.md +++ b/README.ja.md @@ -10,7 +10,7 @@ WebRTCを使ったP2Pファイル共有ツール。サーバーを経由せず - **P2P転送**: ファイルはサーバーを経由せず、ブラウザ間で直接送受信 - **E2E暗号化(オプション)**: リンクの`#k=...`部分に鍵を含めることで、サーバーに鍵を送らずにAES-GCM暗号化 -- **大容量ファイル対応**: 大きめの転送チャンクを使い、SHA-256検証後にダウンロード完了として扱う +- **大容量ファイル対応**: 安全な転送チャンクを調整し、SHA-256検証後にダウンロード完了として扱う - **サーバーレス**: Cloudflare Workers + Durable Objectsで動作、ファイルはサーバーに保存されない - **複数受信者対応**: 1人の送信者から複数人が同時にファイルを受信可能(同時接続数は設定可能) - **ドラッグ&ドロップ**: ファイル選択UIはドラッグ&ドロップに対応 diff --git a/README.md b/README.md index 40ed966..8bad24a 100644 --- a/README.md +++ b/README.md @@ -10,7 +10,7 @@ A P2P file sharing tool using WebRTC. Transfer files directly between browsers w - **P2P Transfer**: Files are sent directly between browsers, not through a server - **E2E Encryption (Optional)**: AES-GCM encryption with key in URL fragment (`#k=...`), never sent to server -- **Large File Friendly**: Uses larger transfer chunks and verifies SHA-256 before marking a download complete +- **Large File Friendly**: Negotiates safe transfer chunks and verifies SHA-256 before marking a download complete - **Serverless**: Runs on Cloudflare Workers + Durable Objects, no file storage on server - **Multiple Receivers**: One sender can transfer to multiple receivers simultaneously (configurable concurrency) - **Drag & Drop**: File selection UI supports drag and drop diff --git a/cli/Cargo.lock b/cli/Cargo.lock index 71eaee2..c9ea43d 100644 --- a/cli/Cargo.lock +++ b/cli/Cargo.lock @@ -1285,7 +1285,7 @@ dependencies = [ [[package]] name = "pairlane-cli" -version = "0.1.1" +version = "0.1.2" dependencies = [ "aes-gcm", "anyhow", diff --git a/cli/Cargo.toml b/cli/Cargo.toml index be65e39..f7a3875 100644 --- a/cli/Cargo.toml +++ b/cli/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "pairlane-cli" -version = "0.1.1" +version = "0.1.2" edition = "2021" [[bin]] diff --git a/cli/npm/package.json b/cli/npm/package.json index b3372d4..91b4782 100644 --- a/cli/npm/package.json +++ b/cli/npm/package.json @@ -1,11 +1,11 @@ { "name": "pairlane", - "version": "0.1.1", + "version": "0.1.2", "description": "P2P file transfer CLI for Pairlane", "license": "MIT", "repository": { "type": "git", - "url": "https://github.com/kiyo-e/p2p-share-files.git" + "url": "git+https://github.com/kiyo-e/p2p-share-files.git" }, "engines": { "node": ">=18" diff --git a/cli/src/main.rs b/cli/src/main.rs index e928acd..fc8d294 100644 --- a/cli/src/main.rs +++ b/cli/src/main.rs @@ -17,7 +17,7 @@ use std::sync::Arc; use tokio::fs::File; use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::sync::{mpsc, Mutex}; -use tokio::time::{sleep, Duration}; +use tokio::time::{sleep, timeout, Duration}; use tokio_tungstenite::connect_async; use tokio_tungstenite::tungstenite::Message; use url::form_urlencoded; @@ -41,7 +41,8 @@ use webrtc::peer_connection::RTCPeerConnection; const AES_KEY_LEN: usize = 32; const AES_NONCE_LEN: usize = 12; const AES_TAG_LEN: usize = 16; -const MAX_FRAME_SIZE: usize = 256 * 1024; +const CLI_SAFE_CHUNK_SIZE: usize = 16 * 1024; +const PREFERRED_CHUNK_SIZE: usize = 256 * 1024; // Design: see README.md and docs/signaling-protocol.md; related to Command and transfer helpers below. #[derive(Parser, Debug)] @@ -83,6 +84,8 @@ enum Command { key: Option, #[arg(long, help = "Keep running after a successful receive")] stay_open: bool, + #[arg(long, help = "Print detailed receive diagnostics")] + debug: bool, }, } @@ -129,10 +132,14 @@ enum DataMessage { size: u64, mime: String, encrypted: bool, - sha256: String, }, #[serde(rename = "done")] - Done, + Done { sha256: String }, + #[serde(rename = "capabilities")] + Capabilities { + #[serde(rename = "maxChunkSize")] + max_chunk_size: usize, + }, } struct RoomInput { @@ -147,7 +154,6 @@ struct FileInfo { name: String, size: u64, mime: String, - sha256: String, } struct OffererPeerState { @@ -186,6 +192,9 @@ struct ReceiveProgress { expected_sha256: Option, hasher: Option, received: u64, + next_progress_percent: u64, + saw_first_chunk: bool, + debug: bool, encrypted: bool, crypto: Option>, success_tx: Option>, @@ -218,11 +227,12 @@ async fn main() -> Result<()> { endpoint, key, stay_open, + debug, } => { let room_input = room_id .or(room_input) .ok_or_else(|| anyhow!("Room ID or URL is required (usage: receive )"))?; - run_receive(&room_input, &output_dir, endpoint.as_deref(), key.as_deref(), stay_open).await + run_receive(&room_input, &output_dir, endpoint.as_deref(), key.as_deref(), stay_open, debug).await } } } @@ -378,6 +388,7 @@ async fn run_receive( endpoint: Option<&str>, key: Option<&str>, stay_open: bool, + debug: bool, ) -> Result<()> { let parsed = parse_room_input(room_input)?; let mut key_override = parsed.key; @@ -422,6 +433,9 @@ async fn run_receive( expected_sha256: None, hasher: None, received: 0, + next_progress_percent: 10, + saw_first_chunk: false, + debug, encrypted: false, crypto, success_tx, @@ -479,6 +493,7 @@ async fn run_receive( pc.on_data_channel(Box::new(move |dc| { let rx_progress = rx_progress.clone(); Box::pin(async move { + debug_log(&rx_progress, "[rtc] datachannel", "open").await; wire_receiver_channel(dc, rx_progress).await; }) })); @@ -606,6 +621,23 @@ async fn create_offerer_peer( let dc_for_open = dc.clone(); let crypto = crypto.clone(); let success_tx = success_tx.clone(); + let (capability_tx, capability_rx) = mpsc::unbounded_channel::(); + let capability_rx = Arc::new(Mutex::new(capability_rx)); + dc.on_message(Box::new(move |msg: DataChannelMessage| { + let capability_tx = capability_tx.clone(); + Box::pin(async move { + if !msg.is_string { + return; + } + let text = match String::from_utf8(msg.data.to_vec()) { + Ok(text) => text, + Err(_) => return, + }; + if let Ok(DataMessage::Capabilities { max_chunk_size }) = serde_json::from_str::(&text) { + let _ = capability_tx.send(max_chunk_size); + } + }) + })); dc.on_open(Box::new(move || { let send_tx = send_tx.clone(); let send_peer_id = send_peer_id.clone(); @@ -614,6 +646,7 @@ async fn create_offerer_peer( let send_state = send_state.clone(); let crypto = crypto.clone(); let success_tx = success_tx.clone(); + let capability_rx = capability_rx.clone(); Box::pin(async move { let mut guard = send_state.lock().await; if guard.sending { @@ -622,7 +655,16 @@ async fn create_offerer_peer( guard.sending = true; drop(guard); - if let Err(err) = send_file(&dc, &file_info, crypto).await { + let receiver_max = timeout(Duration::from_secs(2), async { + capability_rx.lock().await.recv().await + }) + .await + .ok() + .flatten() + .unwrap_or(CLI_SAFE_CHUNK_SIZE); + let chunk_size = choose_chunk_size(receiver_max); + + if let Err(err) = send_file(&dc, &file_info, crypto, chunk_size).await { log_line("[send] error", &format!("{err:#}")); return; } @@ -734,6 +776,17 @@ async fn flush_receiver_candidates(state: &mut ReceiverState) -> Result<()> { } async fn wire_receiver_channel(dc: Arc, progress: Arc>) { + let dc_for_open = dc.clone(); + dc.on_open(Box::new(move || { + let dc = dc_for_open.clone(); + Box::pin(async move { + let _ = send_capabilities(&dc).await; + }) + })); + if dc.ready_state() == RTCDataChannelState::Open { + let _ = send_capabilities(&dc).await; + } + dc.on_message(Box::new(move |msg: DataChannelMessage| { let progress = progress.clone(); Box::pin(async move { @@ -741,7 +794,7 @@ async fn wire_receiver_channel(dc: Arc, progress: Arc(&text) { match parsed { - DataMessage::Meta { name, size, mime, encrypted, sha256 } => { + DataMessage::Meta { name, size, mime, encrypted } => { let mut guard = progress.lock().await; if encrypted && guard.crypto.is_none() { log_line("[recv] error", "encrypted files need a decryption key"); @@ -758,9 +811,11 @@ async fn wire_receiver_channel(dc: Arc, progress: Arc { @@ -769,10 +824,13 @@ async fn wire_receiver_channel(dc: Arc, progress: Arc { + DataMessage::Done { sha256 } => { let mut guard = progress.lock().await; + debug_log_locked(&guard, "[recv] done", "checksum received"); + guard.expected_sha256 = Some(sha256); finalize_receive(&mut guard).await; } + DataMessage::Capabilities { .. } => {} } } } @@ -811,13 +869,15 @@ async fn wire_receiver_channel(dc: Arc, progress: Arc 0 && guard.received >= guard.expected_size { - finalize_receive(&mut guard).await; - } + log_receive_progress(&mut guard); } else if let Err(err) = write_result { log_line("[recv] error", &format!("{err:#}")); notify_receive_done(&mut guard, false); @@ -827,7 +887,21 @@ async fn wire_receiver_channel(dc: Arc, progress: Arc>) -> Result<()> { +async fn send_capabilities(dc: &RTCDataChannel) -> Result<()> { + let capabilities = serde_json::json!({ + "type": "capabilities", + "maxChunkSize": CLI_SAFE_CHUNK_SIZE, + }); + dc.send_text(capabilities.to_string()).await?; + Ok(()) +} + +async fn send_file( + dc: &RTCDataChannel, + info: &FileInfo, + crypto: Option>, + max_frame_size: usize, +) -> Result<()> { let encrypted = crypto.is_some(); let meta = serde_json::json!({ "type": "meta", @@ -835,23 +909,24 @@ async fn send_file(dc: &RTCDataChannel, info: &FileInfo, crypto: Option Result { .first_or_octet_stream() .essence_str() .to_string(); - let sha256 = file_sha256(path).await?; Ok(FileInfo { path: path.to_path_buf(), name, size, mime, - sha256, }) } -async fn file_sha256(path: &Path) -> Result { - let mut file = File::open(path).await?; - let mut hasher = Sha256::new(); - let mut buffer = vec![0u8; MAX_FRAME_SIZE]; - loop { - let read = file.read(&mut buffer).await?; - if read == 0 { - break; - } - hasher.update(&buffer[..read]); - } - Ok(bytes_to_hex(&hasher.finalize())) -} - async fn finalize_receive(progress: &mut ReceiveProgress) { if progress.file.is_none() { return; @@ -925,7 +986,14 @@ async fn finalize_receive(progress: &mut ReceiveProgress) { return; } }; - let expected_sha256 = progress.expected_sha256.take().unwrap_or_default(); + let expected_sha256 = match progress.expected_sha256.take() { + Some(sha256) => sha256, + None => { + log_line("[recv] error", "checksum missing"); + notify_receive_done(progress, false); + return; + } + }; if actual_sha256 != expected_sha256 { log_line("[recv] error", "sha256 mismatch; keeping partial file"); notify_receive_done(progress, false); @@ -956,6 +1024,41 @@ fn notify_receive_done(progress: &mut ReceiveProgress, ok: bool) { } } +fn choose_chunk_size(receiver_max: usize) -> usize { + receiver_max.clamp(CLI_SAFE_CHUNK_SIZE, PREFERRED_CHUNK_SIZE) +} + +async fn debug_log(progress: &Arc>, label: &str, value: &str) { + let guard = progress.lock().await; + debug_log_locked(&guard, label, value); +} + +fn debug_log_locked(progress: &ReceiveProgress, label: &str, value: &str) { + if progress.debug { + log_line(label, value); + } +} + +fn log_receive_progress(progress: &mut ReceiveProgress) { + if !progress.debug { + return; + } + if progress.expected_size == 0 { + return; + } + let percent = progress.received.saturating_mul(100) / progress.expected_size; + while progress.next_progress_percent <= 100 && percent >= progress.next_progress_percent { + log_line( + "[recv] progress", + &format!( + "{}% ({} / {} bytes)", + progress.next_progress_percent, progress.received, progress.expected_size + ), + ); + progress.next_progress_percent += 10; + } +} + async fn create_peer_connection() -> Result> { let mut media_engine = MediaEngine::default(); media_engine.register_default_codecs()?; diff --git a/docs/signaling-protocol.md b/docs/signaling-protocol.md index 5312913..0dc9c81 100644 --- a/docs/signaling-protocol.md +++ b/docs/signaling-protocol.md @@ -213,13 +213,15 @@ Once WebRTC connection is established: ``` Sender Receiver │ │ - │──── { type: "meta", name, size, sha256, ... } ──►│ + │◄─── { type: "capabilities", maxChunkSize } ─│ + │ │ + │──── { type: "meta", name, size, ... } ──►│ │ │ │──── [binary chunk 1] ───────────────────►│ │──── [binary chunk 2] ───────────────────►│ │──── ... │ │ │ - │──── { type: "done" } ───────────────────►│ + │──── { type: "done", sha256 } ───────────►│ ``` #### Metadata Message @@ -230,12 +232,31 @@ Sender Receiver name: string, // File name size: number, // File size in bytes mime: string, // MIME type - encrypted: boolean, // Whether chunks are encrypted + encrypted: boolean // Whether chunks are encrypted +} +``` + +#### Capabilities Message + +```typescript +{ + type: "capabilities", + maxChunkSize: number // Maximum binary message size this receiver accepts +} +``` + +Receivers send capabilities when the DataChannel opens. Senders choose the transfer chunk size from the receiver limit and their own preferred maximum. CLI receivers using the normal webrtc-rs `on_message` API advertise 16 KiB; browser receivers advertise a larger browser-friendly limit. + +#### Completion Message + +```typescript +{ + type: "done", sha256: string // SHA-256 of the original plaintext file bytes } ``` -Receivers verify the reconstructed plaintext bytes against `sha256` before presenting the transfer as complete. CLI receivers write to a `.partial` file and rename it only after verification succeeds. +Receivers verify the reconstructed plaintext bytes against the completion `sha256` before presenting the transfer as complete. CLI receivers write to a `.partial` file and rename it only after verification succeeds. ### End-to-End Encryption (Optional) diff --git a/src/client/room.tsx b/src/client/room.tsx index 8ab7d19..82bbe43 100644 --- a/src/client/room.tsx +++ b/src/client/room.tsx @@ -33,12 +33,13 @@ type IncomingMeta = { size: number; mime: string; encrypted: boolean; - sha256: string; }; -type DoneMessage = { type: "done" }; +type DoneMessage = { type: "done"; sha256: string }; -type DataMessage = IncomingMeta | DoneMessage; +type CapabilitiesMessage = { type: "capabilities"; maxChunkSize: number }; + +type DataMessage = IncomingMeta | DoneMessage | CapabilitiesMessage; type OutgoingMeta = IncomingMeta; @@ -69,11 +70,13 @@ type OffererPeer = { offerInFlight: boolean; sending: boolean; sent: boolean; + maxChunkSize: number | null; }; const clientId = getClientId(); const t = getT(); -const TRANSFER_CHUNK_SIZE = 256 * 1024; +const CLI_SAFE_CHUNK_SIZE = 16 * 1024; +const WEB_PREFERRED_CHUNK_SIZE = 256 * 1024; const root = document.querySelector("main.container"); const roomId = document.body.dataset.roomId; @@ -197,7 +200,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { } }, [resetPeerSends]); - const finalizeDownload = useCallback(async () => { + const finalizeDownload = useCallback(async (expectedSha256: string) => { const meta = incomingMetaRef.current; if (!meta) return; setStatus(t.status.generatingFile); @@ -208,7 +211,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { const blob = new Blob(recvChunksRef.current, { type: meta.mime }); const actualSha256 = await sha256ArrayBuffer(await blob.arrayBuffer()); - if (actualSha256 !== meta.sha256) { + if (actualSha256 !== expectedSha256) { setDownload(null); setStatus(t.status.checksumMismatch); return; @@ -225,8 +228,6 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { const dc = peer.dc; if (!dc) return; - setStatus(t.status.generatingFile); - const sha256 = await fileSha256(file); const encrypted = !!cryptoKeyRef.current; const meta: OutgoingMeta = { type: "meta", @@ -234,7 +235,6 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { size: file.size, mime: file.type || "application/octet-stream", encrypted, - sha256, }; log("[send] starting:", meta.name, "size:", meta.size, "peer:", peer.peerId); dc.send(JSON.stringify(meta)); @@ -243,6 +243,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { setSendProgress({ sent: 0, total: file.size }); let sent = 0; + const negotiatedChunkSize = peer.maxChunkSize ?? CLI_SAFE_CHUNK_SIZE; dc.bufferedAmountLowThreshold = 4 * 1024 * 1024; const waitDrain = () => @@ -264,7 +265,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { if (dc.bufferedAmount > 8 * 1024 * 1024) await waitDrain(); }; - const chunkSize = encrypted ? TRANSFER_CHUNK_SIZE - 12 - 16 : TRANSFER_CHUNK_SIZE; + const chunkSize = encrypted ? negotiatedChunkSize - 12 - 16 : negotiatedChunkSize; let offset = 0; while (offset < file.size) { const slice = file.slice(offset, offset + chunkSize); @@ -275,7 +276,9 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { offset += value.byteLength; } - dc.send(JSON.stringify({ type: "done" } satisfies DoneMessage)); + setStatus(t.status.generatingFile); + const sha256 = await fileSha256(file); + dc.send(JSON.stringify({ type: "done", sha256 } satisfies DoneMessage)); log("[send] completed, peer:", peer.peerId); setStatus(t.status.sendComplete); peer.sent = true; @@ -292,6 +295,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { if (!file) return; const dc = peer.dc; if (!dc || dc.readyState !== "open") return; + if (peer.maxChunkSize == null) return; log("[send] triggered:", reason, "peer:", peer.peerId); peer.sending = true; @@ -334,6 +338,13 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { ch.onerror = () => { console.warn("[rtc] datachannel error (peer:", peer.peerId + ")"); }; + ch.onmessage = (ev) => { + if (typeof ev.data !== "string") return; + const m = safeJson(ev.data) as DataMessage | null; + if (!m || m.type !== "capabilities") return; + peer.maxChunkSize = chooseChunkSize(m.maxChunkSize); + void trySendPeer(peer, "capabilities"); + }; }; const createOffererPeer = (peerId: string) => { @@ -355,6 +366,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { offerInFlight: false, sending: false, sent: false, + maxChunkSize: null, }; pc.onicecandidate = (ev) => { @@ -424,6 +436,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { ch.binaryType = "arraybuffer"; ch.onopen = () => { log("[rtc] datachannel open (receiver)"); + ch.send(JSON.stringify({ type: "capabilities", maxChunkSize: WEB_PREFERRED_CHUNK_SIZE } satisfies CapabilitiesMessage)); setStatus(t.status.dataChannelReady); }; ch.onclose = () => { @@ -439,6 +452,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { if (typeof ev.data === "string") { const m = safeJson(ev.data) as DataMessage | null; if (!m) return; + if (m.type === "capabilities") return; if (m.type === "meta") { log("[recv] starting:", m.name, "size:", m.size); @@ -457,7 +471,7 @@ function RoomApp({ roomId, maxConcurrent }: RoomAppProps) { if (m.type === "done") { log("[recv] completed"); - await finalizeDownload(); + await finalizeDownload(m.sha256); } return; } @@ -999,6 +1013,11 @@ function progressText(sent: number, total: number) { return `${pct}% (${formatBytes(sent)} / ${formatBytes(total)})`; } +function chooseChunkSize(receiverMax: number) { + if (!Number.isFinite(receiverMax) || receiverMax <= 0) return CLI_SAFE_CHUNK_SIZE; + return Math.max(CLI_SAFE_CHUNK_SIZE, Math.min(WEB_PREFERRED_CHUNK_SIZE, Math.floor(receiverMax))); +} + async function copyText(s: string) { try { await navigator.clipboard.writeText(s);