From dbd512ba1fcbdef1458406555030a264e80a2408 Mon Sep 17 00:00:00 2001 From: Yuto Nishida Date: Fri, 26 Jun 2026 02:11:54 -0700 Subject: [PATCH] Use remote builder and fix large file upload in depot --- andref-ipfs-depot/Cargo.lock | 1 + andref-ipfs-depot/Cargo.toml | 2 + andref-ipfs-depot/src/ipfs.rs | 12 ++-- andref-ipfs-depot/src/web.rs | 80 +++++++++++++++------ exports/whale/digests/andref-ipfs-depot.txt | 2 +- milky-way/lib/andref-ipfs-depot.libsonnet | 3 + milky-way/lib/traefik.libsonnet | 13 +++- venus/modules/nixos-darwin/sodium.nix | 54 ++++++++------ 8 files changed, 118 insertions(+), 49 deletions(-) diff --git a/andref-ipfs-depot/Cargo.lock b/andref-ipfs-depot/Cargo.lock index 5bb1dd1..0e5a2f7 100644 --- a/andref-ipfs-depot/Cargo.lock +++ b/andref-ipfs-depot/Cargo.lock @@ -22,6 +22,7 @@ name = "andref-ipfs-depot" version = "0.1.0" dependencies = [ "axum", + "futures-util", "reqwest", "serde", "serde_json", diff --git a/andref-ipfs-depot/Cargo.toml b/andref-ipfs-depot/Cargo.toml index f98f791..018934c 100644 --- a/andref-ipfs-depot/Cargo.toml +++ b/andref-ipfs-depot/Cargo.toml @@ -8,10 +8,12 @@ license = "MIT" [dependencies] axum = { version = "0.8", features = ["multipart"] } +futures-util = "0.3" reqwest = { version = "0.12", default-features = false, features = [ "json", "multipart", "rustls-tls", + "stream", ] } serde = { version = "1", features = ["derive"] } serde_json = "1" diff --git a/andref-ipfs-depot/src/ipfs.rs b/andref-ipfs-depot/src/ipfs.rs index e6a44a5..14e5e72 100644 --- a/andref-ipfs-depot/src/ipfs.rs +++ b/andref-ipfs-depot/src/ipfs.rs @@ -35,11 +35,13 @@ impl From for IpfsError { } } -/// Upload `bytes` to kubo, pinned, as a CIDv1 (base32 -- required so the CID is a DNS-safe -/// subdomain label). Returns the resulting CID. -pub async fn add(state: &AppState, filename: String, bytes: Vec) -> Result { - let part = reqwest::multipart::Part::bytes(bytes) - .file_name(filename) +/// Stream `body` to kubo, pinned, as a CIDv1 (base32 -- required so the CID is a DNS-safe +/// subdomain label). The body is a streaming multipart part, so arbitrarily large files transfer +/// with bounded memory. Returns the resulting CID. The part's filename is irrelevant to the result +/// (the CID is content-addressed and the gateway URL is `.`), so a generic name is used. +pub async fn add(state: &AppState, body: reqwest::Body) -> Result { + let part = reqwest::multipart::Part::stream(body) + .file_name("file") .mime_str("application/octet-stream")?; let form = reqwest::multipart::Form::new().part("file", part); diff --git a/andref-ipfs-depot/src/web.rs b/andref-ipfs-depot/src/web.rs index 6c41c47..6949e32 100644 --- a/andref-ipfs-depot/src/web.rs +++ b/andref-ipfs-depot/src/web.rs @@ -1,6 +1,7 @@ //! The HTTP surface: serves the upload page + assets, accepts the upload, and on success posts //! the resulting link back to the originating Discord channel. +use axum::body::Bytes; use axum::extract::{DefaultBodyLimit, Multipart, Path, State}; use axum::http::{header, StatusCode}; use axum::response::{Html, IntoResponse, Response}; @@ -13,9 +14,10 @@ use crate::assets; use crate::ipfs; use crate::state::AppState; -/// Cap on a single upload. The endpoint is public (token-gated only) and the body is buffered in -/// memory, so this bound is what stops it being a memory-exhaustion DoS. -const MAX_UPLOAD_BYTES: usize = 100 * 1024 * 1024; +/// Sanity ceiling on a single upload. The body is streamed (not buffered), so this isn't a memory +/// bound -- it's a guard against an absurd upload. It's also effectively bounded by kubo's repo PVC +/// (10Gi); grow both together if needed. +const MAX_UPLOAD_BYTES: usize = 8 * 1024 * 1024 * 1024; pub fn router(state: AppState) -> Router { Router::new() @@ -66,9 +68,9 @@ struct UploadResult { async fn upload( State(state): State, Path(token): Path, - mut multipart: Multipart, + multipart: Multipart, ) -> Result, (StatusCode, String)> { - // Consume the token BEFORE reading the body, so a replayed/expired link is rejected without + // Consume the token BEFORE touching the body, so a replayed/expired link is rejected without // streaming the file. This burns the token even on a later failure -- fine, the member just // re-runs `/upload`. let pending = state @@ -76,23 +78,57 @@ async fn upload( .consume(&token) .ok_or((StatusCode::FORBIDDEN, "invalid or expired link".to_string()))?; - let field = multipart - .next_field() - .await - .map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))? - .ok_or((StatusCode::BAD_REQUEST, "no file field".to_string()))?; - let filename = field.file_name().unwrap_or("upload").to_string(); - let bytes = field - .bytes() - .await - .map_err(|e| (StatusCode::BAD_REQUEST, e.to_string()))?; - - let cid = ipfs::add(&state, filename, bytes.to_vec()) - .await - .map_err(|e| { - tracing::error!("kubo add failed: {e}"); - (StatusCode::BAD_GATEWAY, "upload to IPFS failed".to_string()) - })?; + // Stream the upload straight through to kubo so memory stays bounded regardless of file size (a + // GB-scale file must never be buffered in the pod). A task owns the multipart reader and forwards + // each chunk over a bounded channel -- which also applies backpressure to the browser when kubo + // is slow -- and reqwest pulls from it as the request body. + tracing::info!("upload: streaming to kubo (token consumed)"); + let (tx, rx) = tokio::sync::mpsc::channel::>(8); + tokio::spawn(async move { + let mut multipart = multipart; + let mut total: u64 = 0; + match multipart.next_field().await { + Ok(Some(mut field)) => loop { + match field.chunk().await { + Ok(Some(chunk)) => { + total += chunk.len() as u64; + if tx.send(Ok(chunk)).await.is_err() { + // reqwest stopped reading the body -> the kubo side closed/reset. + tracing::warn!("upload: kubo stopped reading after {total} bytes"); + break; + } + } + Ok(None) => { + tracing::info!("upload: read {total} bytes from client (body complete)"); + break; + } + Err(e) => { + // The client->backend body broke (e.g. Traefik aborted the request). + tracing::error!( + "upload: error reading body from client after {total} bytes: {e}" + ); + let _ = tx.send(Err(std::io::Error::other(e.to_string()))).await; + break; + } + } + }, + Ok(None) => tracing::warn!("upload: multipart had no file field"), + Err(e) => { + tracing::error!("upload: multipart error: {e}"); + let _ = tx.send(Err(std::io::Error::other(e.to_string()))).await; + } + } + }); + let body = reqwest::Body::wrap_stream(futures_util::stream::unfold(rx, |mut rx| async move { + rx.recv().await.map(|item| (item, rx)) + })); + + let cid = ipfs::add(&state, body).await.map_err(|e| { + // {e:?} prints the full source chain (reset / timed out / incomplete body / ...). + tracing::error!("kubo add failed: {e:?}"); + (StatusCode::BAD_GATEWAY, "upload to IPFS failed".to_string()) + })?; + tracing::info!("upload: pinned {cid}"); let url = ipfs::gateway_url(&state, &cid); // Announce the result in the channel the command was run in. A bare URL lets Discord unfurl / diff --git a/exports/whale/digests/andref-ipfs-depot.txt b/exports/whale/digests/andref-ipfs-depot.txt index 32f75c6..a9b1f49 100644 --- a/exports/whale/digests/andref-ipfs-depot.txt +++ b/exports/whale/digests/andref-ipfs-depot.txt @@ -1 +1 @@ -sha256:22bfa85fa3c70443427edb38a89cc4448dfd0e6e0064970cfe0356b09161e9b6 \ No newline at end of file +sha256:c6f11b327c94f2234e565a8a5829d65e88f09febe903cb76f523b59c32ba39c0 \ No newline at end of file diff --git a/milky-way/lib/andref-ipfs-depot.libsonnet b/milky-way/lib/andref-ipfs-depot.libsonnet index b463060..60e6c6f 100644 --- a/milky-way/lib/andref-ipfs-depot.libsonnet +++ b/milky-way/lib/andref-ipfs-depot.libsonnet @@ -88,6 +88,9 @@ local images = import 'milky-way/lib/images.libsonnet'; periodSeconds: 30, }, resources: { + // Uploads stream straight through to kubo (never buffered), so memory stays flat + // (tens of MB) regardless of file size -- 256Mi is ample headroom even for a + // multi-GB upload. requests: { memory: '32Mi', cpu: '10m' }, limits: { memory: '256Mi', cpu: '1' }, }, diff --git a/milky-way/lib/traefik.libsonnet b/milky-way/lib/traefik.libsonnet index e278aef..8f30c7d 100644 --- a/milky-way/lib/traefik.libsonnet +++ b/milky-way/lib/traefik.libsonnet @@ -26,7 +26,18 @@ // Listen on standard ports directly (hostNetwork exposes these). ports: { web: { port: 80 }, - websecure: { port: 443 }, + websecure: { + port: 443, + // Traefik v3 defaults respondingTimeouts.readTimeout to 60s, and that covers reading the + // ENTIRE request body. So any HTTPS request whose body takes >60s to upload is aborted + // mid-stream (observed: a slow POST is cut at exactly 60.1s) -- which kills every + // multi-GB upload to andref-ipfs-depot, surfacing as a 502/Bad Gateway. Disable it (0s) + // so long uploads can complete; this is entrypoint-wide, but the other websecure routes + // (kubo gateway, etc.) serve GETs with tiny bodies, so it doesn't weaken them, and + // idleTimeout (180s) still bounds genuinely idle connections. Upload SIZE is still + // bounded by andref-ipfs-depot's own 8 GiB body limit. + transport: { respondingTimeouts: { readTimeout: '0s' } }, + }, }, // Required for binding privileged ports (80, 443) with hostNetwork. securityContext: { diff --git a/venus/modules/nixos-darwin/sodium.nix b/venus/modules/nixos-darwin/sodium.nix index 79fb5b6..bf92bc9 100644 --- a/venus/modules/nixos-darwin/sodium.nix +++ b/venus/modules/nixos-darwin/sodium.nix @@ -83,31 +83,45 @@ in # protocol = "ssh-ng"; # } #]; - # This automatically adds itself (aarch64-linux) to available buildMachines; no need for an entry above + distributedBuilds = true; + # methanol (the cluster/storage node) is a NATIVE x86_64-linux box (16 cores / 27 GB), so x86_64 + # builds -- e.g. whale's container images -- are offloaded to it instead of emulating x86 on the M1 + # via the linux-builder (QEMU TCG, ~5-8x slower). The image BUILDS on methanol; nix then copies the + # result back here, and the whale push-script still runs LOCALLY with the local skopeo credentials + # (see whale/outputs.nix / whale/readme.md). Reached over the LAN as root (same host as deploy-rs); + # the host key is pinned so the daemon's non-interactive SSH doesn't fail verification. (Supersedes + # the commented x86_64.nix.yuto.sh example above.) + buildMachines = [ + { + hostName = "10.0.0.211"; + sshUser = "root"; + sshKey = "/Users/yuto/.ssh/id_rsa"; + systems = [ "x86_64-linux" ]; + maxJobs = 8; + protocol = "ssh-ng"; + supportedFeatures = [ + "big-parallel" + "kvm" + "nixos-test" + "benchmark" + ]; + publicHostKey = "c3NoLWVkMjU1MTkgQUFBQUMzTnphQzFsWkRJMU5URTVBQUFBSU5UNFhYSWt3OWNSNEwzcXVUQzF2a2Z2SVVWUWZleHFYTmRPbllzNW9IOXc="; + } + ]; + # methanol substitutes build inputs from the public caches itself rather than this Mac uploading + # them (this aarch64-darwin host has no x86_64-linux store paths to send anyway). + settings.builders-use-substitutes = true; + + # The auto-managed darwin linux-builder, now for aarch64-linux ONLY (x86_64 goes to methanol above). linux-builder = { enable = true; maxJobs = 4; - # Advertise x86_64-linux in addition to the default aarch64-linux, and let the - # aarch64 guest VM execute x86_64 build steps via QEMU binfmt. Container images - # of prebuilt packages only emulate the tar/gzip assembly, so this stays cheap. - # This is what lets whale's x86_64-linux images build on the M1 (see whale/readme.md). + # aarch64-linux ONLY (the module default). x86_64-linux is built natively by methanol (see + # buildMachines above), so there's no QEMU x86 emulation here and the VM keeps its small default + # size (1 vCPU / 3 GB). # - # Note: after changing `config` here, `darwin-rebuild switch` only reloads the - # service; the running VM keeps the old image until restarted. If x86_64 builds - # fail with `path '…' is not valid` (guest store corruption), reset the builder: + # If an aarch64-linux build fails with `path '…' is not valid` (guest store corruption), reset: # nix run ./flake-profiles/system-sodium#reset-linux-builder - systems = [ "aarch64-linux" "x86_64-linux" ]; - config = { - boot.binfmt.emulatedSystems = [ "x86_64-linux" ]; - # The VM defaults to 1 vCPU / 3 GB, so a single x86_64 cargo build (e.g. whale's - # andref-ipfs-depot) grinds on one emulated core. Give it more so cargo parallelizes. - # Sized to leave the M1 (8 cores / 16 GB) headroom for macOS: ~4 cores + 6 GB are used - # ONLY while a build runs; the VM sits at ~0 when idle. Bump cautiously -- x86_64 steps run - # under TCG (pure emulation), so each vCPU can peg a host core during a build. - # mkForce: the nix-builder-vm profile already pins these (memorySize = 3072), so override. - virtualisation.cores = pkgs.lib.mkForce 4; - virtualisation.memorySize = pkgs.lib.mkForce (6 * 1024); # MiB (was 3072) - }; }; }; -- 2.51.2