chat: attachment cache, download, and transfer hardening (Phase 3)
CI / check (push) Successful in 2m33s
CI / check (push) Successful in 2m33s
Phase 3 of docs/chat-hardening-plan.md — attachments can no longer turn into unbounded memory, bandwidth, decoder, or task pressure (S15 closed; S14's filename half closed). Cache and image cost (3A): AttachmentCache now carries encoded- and decoded-byte budgets (96 MiB / 64 MiB) on top of the count cap, with per-entry weights, replacement accounting, and oldest-first eviction; an individually over-budget fetch services any pending Save/Play from the bytes in hand and is exposed as Evicted instead of retained. validate_image_bytes prechecks header dimensions (per-side AND a new 14 MP total-pixel limit) before any decode; the renderer only ever receives a ≤1600 px downscaled RGBA preview whose w*h*4 cost counts against the decoded budget — originals stay encoded-only for Save. sanitize_filename strips the bidi/zero-width spoofing set (RTL-override extension spoof). Download policy and state (3B): images auto-fetch only when roster- authored AND declared ≤4 MiB, gated by a new deterministic AutoFetchBudget (per-author and session request+byte token buckets, check-then-take, bounded author map) alongside the existing dedup and four-permit bound. Attachment state is now explicit — absence/Loading/ Ready/Failed/Evicted — driven by a new AttachmentFetchStarted event, so skipped or evicted images render a "Load image" button instead of an indefinite "loading…", and repeated clicks can never spawn duplicate fetch tasks. Exact transfers and serve store (3C): fetch_blob requires the received length to equal the declared size (short = local error, overlong = bounded-read reject, empty keeps meaning "sender no longer has it"); the file picker's unbounded read is replaced by a metadata-prechecked cap+1 bounded reader; one Arc<Vec<u8>> now backs the UI cache, command queue, and serve store; served_files is a count- and byte-budgeted FIFO ServeStore (16 entries / 128 MiB). 37 new tests (568 lib total) including a real two-endpoint loopback exercising exact/short/overlong/unknown-id transfers. Plan checkboxes ticked and constant deviations decision-logged. Tests-green-only: the plan's two-machine field-test section remains open. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
+86
-33
@@ -1,11 +1,15 @@
|
|||||||
# Chat hardening — ephemeral implementation plan
|
# Chat hardening — ephemeral implementation plan
|
||||||
|
|
||||||
**Status (2026-07-17):** Phase 1 COMPLETE (shared text policy in `src/sanitize.rs`,
|
**Status (2026-07-18):** Phases 1–3 COMPLETE. Phase 1 = shared text policy in
|
||||||
ceilings enforced at UI input, sign point, and gossip ingress; all gates green).
|
`src/sanitize.rs`, ceilings enforced at UI input, sign point, and gossip ingress.
|
||||||
Phases 2–5 not started. This is a temporary scope
|
Phase 2 = roster-bound authorship (`src/core/chatroster.rs`), replay dedup + rate
|
||||||
contract for hardening the existing room chat. Update the checkboxes and decision
|
limits (`ChatIngressGate` in `src/network/gossip.rs`). Phase 3 = attachment
|
||||||
log as work lands, then delete this file when the work is complete. Do not add
|
cache/serve-store budgets, downscaled previews, auto-fetch byte/request budgets
|
||||||
link previews as part of this effort.
|
(`src/core/fetchbudget.rs`), exact transfers, bounded local reads. All gates
|
||||||
|
green each phase. Phases 4–5 not started. This is a temporary scope contract for
|
||||||
|
hardening the existing room chat. Update the checkboxes and decision log as work
|
||||||
|
lands, then delete this file when the work is complete. Do not add link previews
|
||||||
|
as part of this effort.
|
||||||
|
|
||||||
## Goal
|
## Goal
|
||||||
|
|
||||||
@@ -182,70 +186,70 @@ unbounded memory, bandwidth, decoder, or task pressure.
|
|||||||
|
|
||||||
### 3A. Cache and image cost
|
### 3A. Cache and image cost
|
||||||
|
|
||||||
- [ ] Extend `AttachmentCache` with encoded-byte and decoded-preview-byte counters.
|
- [x] Extend `AttachmentCache` with encoded-byte and decoded-preview-byte counters.
|
||||||
Preserve the count cap, but evict oldest entries until all three budgets fit.
|
Preserve the count cap, but evict oldest entries until all three budgets fit.
|
||||||
- [ ] Give every entry an explicit weight. Replacement must subtract the old
|
- [x] Give every entry an explicit weight. Replacement must subtract the old
|
||||||
weight before checking/inserting the new one.
|
weight before checking/inserting the new one.
|
||||||
- [ ] Decide behavior for a single entry larger than the cache budget: service an
|
- [x] Decide behavior for a single entry larger than the cache budget: service an
|
||||||
immediate pending Save/Play request without retaining it, then expose it as
|
immediate pending Save/Play request without retaining it, then expose it as
|
||||||
evicted/unavailable rather than exceeding the budget.
|
evicted/unavailable rather than exceeding the budget.
|
||||||
- [ ] Add a total-pixel limit to `validate_image_bytes` in addition to the existing
|
- [x] Add a total-pixel limit to `validate_image_bytes` in addition to the existing
|
||||||
width/height limit.
|
width/height limit.
|
||||||
- [ ] Build a downscaled inline preview handle with a maximum 1,600 px side. Keep
|
- [x] Build a downscaled inline preview handle with a maximum 1,600 px side. Keep
|
||||||
original bytes only for Save; do not hand a full-resolution 4K image to the
|
original bytes only for Save; do not hand a full-resolution 4K image to the
|
||||||
renderer merely to display it at chat width.
|
renderer merely to display it at chat width.
|
||||||
- [ ] Count estimated RGBA preview cost (`width * height * 4`) against the decoded
|
- [x] Count estimated RGBA preview cost (`width * height * 4`) against the decoded
|
||||||
budget even if iced internally copies or uploads it.
|
budget even if iced internally copies or uploads it.
|
||||||
- [ ] Strip the same bidi/zero-width spoofing characters used for display labels
|
- [x] Strip the same bidi/zero-width spoofing characters used for display labels
|
||||||
from attachment filenames, while preserving ordinary Unicode filenames.
|
from attachment filenames, while preserving ordinary Unicode filenames.
|
||||||
|
|
||||||
### 3B. Automatic download policy and state
|
### 3B. Automatic download policy and state
|
||||||
|
|
||||||
- [ ] Auto-fetch only roster-authored images whose declared size is at or below
|
- [x] Auto-fetch only roster-authored images whose declared size is at or below
|
||||||
`MAX_AUTO_IMAGE_BYTES`; keep the existing `(author,id)` dedup and four-permit
|
`MAX_AUTO_IMAGE_BYTES`; keep the existing `(author,id)` dedup and four-permit
|
||||||
concurrency bound.
|
concurrency bound.
|
||||||
- [ ] Add per-author and session byte/request budgets for automatic fetches so a
|
- [x] Add per-author and session byte/request budgets for automatic fetches so a
|
||||||
peer cannot drain bandwidth sequentially after each permit is released.
|
peer cannot drain bandwidth sequentially after each permit is released.
|
||||||
- [ ] Represent `NotFetched`, `Loading`, `Ready`, `Failed`, and `Evicted` distinctly
|
- [x] Represent `NotFetched`, `Loading`, `Ready`, `Failed`, and `Evicted` distinctly
|
||||||
enough for the UI to avoid an indefinite “loading…” label when auto-fetch was
|
enough for the UI to avoid an indefinite “loading…” label when auto-fetch was
|
||||||
skipped or the cache evicted an item.
|
skipped or the cache evicted an item.
|
||||||
- [ ] Render a Load image button for large/skipped images. A manual click may use
|
- [x] Render a Load image button for large/skipped images. A manual click may use
|
||||||
the 25 MiB file cap but still observes cache/decoder budgets.
|
the 25 MiB file cap but still observes cache/decoder budgets.
|
||||||
- [ ] Ensure a repeated click cannot create duplicate unguarded fetch tasks.
|
- [x] Ensure a repeated click cannot create duplicate unguarded fetch tasks.
|
||||||
- [ ] Keep non-image attachments manual-only.
|
- [x] Keep non-image attachments manual-only.
|
||||||
|
|
||||||
### 3C. Exact transfers, local reads, and served files
|
### 3C. Exact transfers, local reads, and served files
|
||||||
|
|
||||||
- [ ] In `IrohTransport::fetch_blob`, require `bytes.len() as u64 == declared_size`.
|
- [x] In `IrohTransport::fetch_blob`, require `bytes.len() as u64 == declared_size`.
|
||||||
Reject empty, short, and overlong transfers with a concise local error.
|
Reject empty, short, and overlong transfers with a concise local error.
|
||||||
- [ ] Replace the file picker's unbounded `FileHandle::read()` with a helper that
|
- [x] Replace the file picker's unbounded `FileHandle::read()` with a helper that
|
||||||
reads at most `MAX_ATTACHMENT_BYTES + 1`. Check metadata first where available,
|
reads at most `MAX_ATTACHMENT_BYTES + 1`. Check metadata first where available,
|
||||||
but retain the bounded read because metadata can race or be unavailable through
|
but retain the bounded read because metadata can race or be unavailable through
|
||||||
a portal.
|
a portal.
|
||||||
- [ ] Avoid duplicating a full attachment across UI, command queue, and serve store.
|
- [x] Avoid duplicating a full attachment across UI, command queue, and serve store.
|
||||||
Prefer `Arc<Vec<u8>>`/`Arc<[u8]>` through `AttachmentState`, `CoreCommand`, and
|
Prefer `Arc<Vec<u8>>`/`Arc<[u8]>` through `AttachmentState`, `CoreCommand`, and
|
||||||
`serve_attachment`, subject to iced handle API constraints.
|
`serve_attachment`, subject to iced handle API constraints.
|
||||||
- [ ] Replace the unbounded session `served_files` map with a count- and byte-
|
- [x] Replace the unbounded session `served_files` map with a count- and byte-
|
||||||
budgeted FIFO store. Evicted ids should produce the existing “sender no longer
|
budgeted FIFO store. Evicted ids should produce the existing “sender no longer
|
||||||
has the file” response rather than stale or aliased data.
|
has the file” response rather than stale or aliased data.
|
||||||
- [ ] Keep attachment ids keyed by author on receipt and preserve all existing
|
- [x] Keep attachment ids keyed by author on receipt and preserve all existing
|
||||||
request-length, timeout, filename, and decoder checks.
|
request-length, timeout, filename, and decoder checks.
|
||||||
|
|
||||||
### Phase 3 tests
|
### Phase 3 tests
|
||||||
|
|
||||||
- [ ] Byte-budget eviction, count eviction, replacement accounting, clear/reset,
|
- [x] Byte-budget eviction, count eviction, replacement accounting, clear/reset,
|
||||||
and an individually overweight entry.
|
and an individually overweight entry.
|
||||||
- [ ] Decoded-preview budget and downscale dimensions for wide, tall, square, and
|
- [x] Decoded-preview budget and downscale dimensions for wide, tall, square, and
|
||||||
boundary images.
|
boundary images.
|
||||||
- [ ] Image with valid per-side dimensions but excessive total pixels is rejected.
|
- [x] Image with valid per-side dimensions but excessive total pixels is rejected.
|
||||||
- [ ] A declared 4 MiB image auto-fetches; the first byte over the limit requires a
|
- [x] A declared 4 MiB image auto-fetches; the first byte over the limit requires a
|
||||||
click.
|
click.
|
||||||
- [ ] Per-author/session auto-fetch budgets recover according to their policy and
|
- [x] Per-author/session auto-fetch budgets recover according to their policy and
|
||||||
never exceed task concurrency.
|
never exceed task concurrency.
|
||||||
- [ ] Short, exact, and overlong file responses.
|
- [x] Short, exact, and overlong file responses.
|
||||||
- [ ] Local file reader stops at cap + 1 instead of allocating the full source.
|
- [x] Local file reader stops at cap + 1 instead of allocating the full source.
|
||||||
- [ ] Served-file FIFO/byte eviction and replacement accounting.
|
- [x] Served-file FIFO/byte eviction and replacement accounting.
|
||||||
- [ ] Same attachment id from two authors remains isolated throughout fetch, cache,
|
- [x] Same attachment id from two authors remains isolated throughout fetch, cache,
|
||||||
save, and display.
|
save, and display.
|
||||||
|
|
||||||
## Phase 4 — URL and rendering resilience
|
## Phase 4 — URL and rendering resilience
|
||||||
@@ -431,3 +435,52 @@ The plan is complete when:
|
|||||||
security boundary (every peer independently enforces the buckets), and a
|
security boundary (every peer independently enforces the buckets), and a
|
||||||
local silent drop would be a UX regression better designed alongside Phase
|
local silent drop would be a UX regression better designed alongside Phase
|
||||||
5's honest send status.
|
5's honest send status.
|
||||||
|
- **2026-07-18 (Phase 3):** Constants that deviate from the proposed table, all
|
||||||
|
bound-tested: total decoded pixels **14 MP** (not 16 MP) so the bound clears
|
||||||
|
12 MP phone photos (4032×3024) yet actually binds inside the 4096²≈16.8 MP
|
||||||
|
per-side envelope; cache encoded budget **96 MiB** (not 128) — still several
|
||||||
|
full-size files, tighter worst case; serve store **128 MiB + 16 entries**
|
||||||
|
(not 256 MiB) — a sender's own session should not pin a quarter GiB.
|
||||||
|
- **2026-07-18 (Phase 3):** `validate_image_bytes`/`decode_preview` precheck
|
||||||
|
dimensions from the container HEADER (`into_dimensions`) before any pixel
|
||||||
|
decode, so an over-limit decode bomb is rejected without paying its decode
|
||||||
|
cost; the decode-time `image::Limits` remain as defense in depth, and the
|
||||||
|
decoded dimensions must equal the prechecked header dimensions.
|
||||||
|
- **2026-07-18 (Phase 3):** Budget-pressure evictions leave NO cache entry
|
||||||
|
(absence = NotFetched → the same Load/Download affordance), while the
|
||||||
|
explicit `Evicted` state marks only an *individually over-budget* fetch whose
|
||||||
|
bytes were used once (pending Save/Play serviced from hand) and dropped. Both
|
||||||
|
render load-on-demand; only the bookkeeping differs.
|
||||||
|
- **2026-07-18 (Phase 3):** The core still runs `validate_image_bytes` before
|
||||||
|
emitting `AttachmentReady`, and the UI decodes once more to build the ≤1600px
|
||||||
|
preview. Two bounded decodes per image were accepted over shipping decoded
|
||||||
|
RGBA across the channel (which would defeat the encoded-only Arc sharing).
|
||||||
|
- **2026-07-18 (Phase 3):** The image lightbox now enlarges the ≤1600px preview
|
||||||
|
handle, not the original bitmap — originals are retained encoded-only for
|
||||||
|
Save. At the lightbox's window-sized draw area the visual difference is nil
|
||||||
|
for the chat use case; full fidelity remains one Save away.
|
||||||
|
- **2026-07-18 (Phase 3):** `AutoFetchBudget` checks all four buckets
|
||||||
|
(author/session × requests/bytes) and only then consumes atomically, so a
|
||||||
|
rejection burns nothing (no refund path like Phase 2's room bucket needed).
|
||||||
|
Tokens ARE consumed if the four-permit semaphore then rejects the spawn —
|
||||||
|
that only happens mid-flood, when charging the author is the intent.
|
||||||
|
- **2026-07-18 (Phase 3):** The auto-fetch budget's author map prunes
|
||||||
|
least-recently-active past 64 entries instead of wiring roster eviction into
|
||||||
|
the event task: authors are roster-gated upstream (≤32 live members), so
|
||||||
|
strangers cannot churn the map, and a pruned author returning with full
|
||||||
|
buckets is within policy.
|
||||||
|
- **2026-07-18 (Phase 3):** Music-track serving shares the bounded serve store
|
||||||
|
with chat attachments. A user who sends enough large attachments during a
|
||||||
|
broadcast can evict their own current track; listeners then get the standard
|
||||||
|
"sender no longer has the file" failure. Accepted: budget honesty over a
|
||||||
|
second store, and the store comfortably fits current+next track plus a
|
||||||
|
normal chat working set.
|
||||||
|
- **2026-07-18 (Phase 3):** The clip player's command channel still takes one
|
||||||
|
owned byte copy at the moment of a Play click (small, human-initiated). The
|
||||||
|
Arc de-duplication targeted the send path (UI cache / command queue / serve
|
||||||
|
store), which now shares a single allocation.
|
||||||
|
- **2026-07-18 (Phase 3):** Overlong transfers are rejected by the transport
|
||||||
|
read itself (`read_to_end(size)` errors past the bound) rather than an
|
||||||
|
explicit length compare; short transfers get the explicit
|
||||||
|
`len == declared_size` check. Music fetches ride `fetch_blob`, so they
|
||||||
|
inherit exactness for free.
|
||||||
|
|||||||
+442
-82
@@ -321,13 +321,23 @@ struct ChatEntry {
|
|||||||
attachment: Option<crate::files::ChatAttachment>,
|
attachment: Option<crate::files::ChatAttachment>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Fetch state of a chat attachment's bytes (session-only).
|
/// Fetch state of a chat attachment's bytes (session-only). Absence from the
|
||||||
|
/// cache means NOT FETCHED (auto-fetch skipped or never attempted) — the UI
|
||||||
|
/// renders a Load/Download button for it, never a loading label (Phase 3B).
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
enum AttachmentState {
|
enum AttachmentState {
|
||||||
/// Bytes in hand (image decoded-valid, or a file ready to save).
|
/// A fetch task is running for this key (core announced it, or the user
|
||||||
Ready(Vec<u8>),
|
/// clicked). Suppresses duplicate fetches on repeated clicks.
|
||||||
|
Loading,
|
||||||
|
/// Bytes in hand (image decoded-valid, or a file ready to save). Shared so
|
||||||
|
/// the cache, save tasks, and the serve path hold one allocation.
|
||||||
|
Ready(std::sync::Arc<Vec<u8>>),
|
||||||
/// Fetch or decode failed; carries a short reason for the UI.
|
/// Fetch or decode failed; carries a short reason for the UI.
|
||||||
Failed(String),
|
Failed(String),
|
||||||
|
/// The bytes were fetched but could not be retained within the cache
|
||||||
|
/// budgets (or were dropped under budget pressure after any pending
|
||||||
|
/// save/play was serviced). Rendered like NotFetched: load-on-demand.
|
||||||
|
Evicted,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Identifies one fetched attachment by BOTH the authoring peer and the
|
/// Identifies one fetched attachment by BOTH the authoring peer and the
|
||||||
@@ -341,13 +351,31 @@ type AttachmentKey = (EndpointId, crate::files::AttachmentId);
|
|||||||
/// above any realistic on-screen image working set.
|
/// above any realistic on-screen image working set.
|
||||||
const ATTACHMENT_CACHE_CAP: usize = 64;
|
const ATTACHMENT_CACHE_CAP: usize = 64;
|
||||||
|
|
||||||
|
/// Budget for retained ENCODED attachment bytes (the original file bytes kept
|
||||||
|
/// for Save). The count cap alone would admit 64 × 25 MiB = 1.6 GiB; this keeps
|
||||||
|
/// the worst case at a few large attachments plus the normal small working set
|
||||||
|
/// (Phase 3A).
|
||||||
|
const ATTACHMENT_CACHE_ENCODED_BUDGET: usize = 96 * 1024 * 1024;
|
||||||
|
|
||||||
|
/// Budget for DECODED preview bytes (estimated RGBA cost `w*h*4` of each inline
|
||||||
|
/// preview, counted even though iced may copy/upload internally). Previews are
|
||||||
|
/// capped at 1600 px per side (≈10 MiB RGBA worst case), so this comfortably
|
||||||
|
/// covers the on-screen set while bounding an image-spam session (Phase 3A).
|
||||||
|
const ATTACHMENT_CACHE_DECODED_BUDGET: usize = 64 * 1024 * 1024;
|
||||||
|
|
||||||
/// One cached attachment result: its fetch state plus, for ready images, the
|
/// One cached attachment result: its fetch state plus, for ready images, the
|
||||||
/// pre-built iced image handle (built once on arrival, not per redraw — the
|
/// pre-built iced image handle (built once on arrival, not per redraw — the
|
||||||
/// e917c53 flicker fix). `handle` is `None` for files and failures.
|
/// e917c53 flicker fix). `handle` is `None` for files and failures.
|
||||||
|
///
|
||||||
|
/// `encoded` / `decoded` are this entry's weights against the cache byte
|
||||||
|
/// budgets: the retained encoded length and the preview's estimated RGBA cost.
|
||||||
|
/// Both are zero for non-`Ready` states.
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
struct AttachmentEntry {
|
struct AttachmentEntry {
|
||||||
state: AttachmentState,
|
state: AttachmentState,
|
||||||
handle: Option<iced::widget::image::Handle>,
|
handle: Option<PreviewHandle>,
|
||||||
|
encoded: usize,
|
||||||
|
decoded: usize,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Bounded store of fetched chat-attachment results, keyed by [`AttachmentKey`].
|
/// Bounded store of fetched chat-attachment results, keyed by [`AttachmentKey`].
|
||||||
@@ -367,38 +395,100 @@ struct AttachmentCache {
|
|||||||
/// exactly the present keys (one entry each), so it is bounded by `cap`.
|
/// exactly the present keys (one entry each), so it is bounded by `cap`.
|
||||||
order: VecDeque<AttachmentKey>,
|
order: VecDeque<AttachmentKey>,
|
||||||
cap: usize,
|
cap: usize,
|
||||||
|
/// Byte budgets and their running totals (Phase 3A). Totals are exactly the
|
||||||
|
/// sum of the present entries' weights — every remove/replace subtracts the
|
||||||
|
/// old weights before new ones are counted.
|
||||||
|
encoded_budget: usize,
|
||||||
|
decoded_budget: usize,
|
||||||
|
encoded_total: usize,
|
||||||
|
decoded_total: usize,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl AttachmentCache {
|
impl AttachmentCache {
|
||||||
fn new(cap: usize) -> Self {
|
fn new(cap: usize, encoded_budget: usize, decoded_budget: usize) -> Self {
|
||||||
Self {
|
Self {
|
||||||
entries: HashMap::new(),
|
entries: HashMap::new(),
|
||||||
order: VecDeque::new(),
|
order: VecDeque::new(),
|
||||||
cap: cap.max(1),
|
cap: cap.max(1),
|
||||||
|
encoded_budget,
|
||||||
|
decoded_budget,
|
||||||
|
encoded_total: 0,
|
||||||
|
decoded_total: 0,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Insert or replace an entry. A brand-new key evicts the oldest entries
|
/// Insert or replace an entry, keeping the count cap AND both byte budgets
|
||||||
/// until there is room; replacing an existing key keeps its position (and so
|
/// satisfied by evicting oldest entries first. Replacing an existing key
|
||||||
/// its age), only updating the value.
|
/// keeps its position (and so its age); its old weights are subtracted
|
||||||
|
/// before the new ones are checked.
|
||||||
|
///
|
||||||
|
/// Returns `false` if the entry ALONE exceeds a byte budget: it is not
|
||||||
|
/// retained (and an existing entry under the key is dropped), leaving no
|
||||||
|
/// entry for `key`. The caller services any immediate need from the bytes
|
||||||
|
/// in hand and stores an [`AttachmentState::Evicted`] marker instead.
|
||||||
fn insert(
|
fn insert(
|
||||||
&mut self,
|
&mut self,
|
||||||
key: AttachmentKey,
|
key: AttachmentKey,
|
||||||
state: AttachmentState,
|
state: AttachmentState,
|
||||||
handle: Option<iced::widget::image::Handle>,
|
handle: Option<PreviewHandle>,
|
||||||
) {
|
) -> bool {
|
||||||
if !self.entries.contains_key(&key) {
|
let encoded = match &state {
|
||||||
while self.entries.len() >= self.cap {
|
AttachmentState::Ready(data) => data.len(),
|
||||||
match self.order.pop_front() {
|
_ => 0,
|
||||||
Some(old) => {
|
};
|
||||||
self.entries.remove(&old);
|
let decoded = handle_decoded_cost(&state, &handle);
|
||||||
}
|
// Replacement: subtract the old weights first so they don't count
|
||||||
None => break,
|
// against the new entry's fit.
|
||||||
}
|
if let Some(old) = self.entries.get(&key) {
|
||||||
|
self.encoded_total -= old.encoded;
|
||||||
|
self.decoded_total -= old.decoded;
|
||||||
|
}
|
||||||
|
if encoded > self.encoded_budget || decoded > self.decoded_budget {
|
||||||
|
// Individually overweight: never retained, budgets never exceeded.
|
||||||
|
if self.entries.remove(&key).is_some() {
|
||||||
|
self.order.retain(|k| k != &key);
|
||||||
}
|
}
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
let replacing = self.entries.contains_key(&key);
|
||||||
|
// Evict oldest (skipping the key being replaced, which keeps its slot)
|
||||||
|
// until the count cap and both budgets accommodate the new entry.
|
||||||
|
while self.would_overflow(replacing, encoded, decoded) {
|
||||||
|
let Some(victim) = self.order.iter().find(|k| **k != key).copied() else {
|
||||||
|
break;
|
||||||
|
};
|
||||||
|
self.remove_entry(&victim);
|
||||||
|
}
|
||||||
|
if !replacing {
|
||||||
self.order.push_back(key);
|
self.order.push_back(key);
|
||||||
}
|
}
|
||||||
self.entries.insert(key, AttachmentEntry { state, handle });
|
self.encoded_total += encoded;
|
||||||
|
self.decoded_total += decoded;
|
||||||
|
self.entries.insert(
|
||||||
|
key,
|
||||||
|
AttachmentEntry {
|
||||||
|
state,
|
||||||
|
handle,
|
||||||
|
encoded,
|
||||||
|
decoded,
|
||||||
|
},
|
||||||
|
);
|
||||||
|
true
|
||||||
|
}
|
||||||
|
|
||||||
|
fn would_overflow(&self, replacing: bool, encoded: usize, decoded: usize) -> bool {
|
||||||
|
let count_full = !replacing && self.entries.len() >= self.cap;
|
||||||
|
count_full
|
||||||
|
|| self.encoded_total + encoded > self.encoded_budget
|
||||||
|
|| self.decoded_total + decoded > self.decoded_budget
|
||||||
|
}
|
||||||
|
|
||||||
|
fn remove_entry(&mut self, key: &AttachmentKey) {
|
||||||
|
if let Some(old) = self.entries.remove(key) {
|
||||||
|
self.encoded_total -= old.encoded;
|
||||||
|
self.decoded_total -= old.decoded;
|
||||||
|
self.order.retain(|k| k != key);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
fn get(&self, key: &AttachmentKey) -> Option<&AttachmentState> {
|
fn get(&self, key: &AttachmentKey) -> Option<&AttachmentState> {
|
||||||
@@ -406,7 +496,9 @@ impl AttachmentCache {
|
|||||||
}
|
}
|
||||||
|
|
||||||
fn handle(&self, key: &AttachmentKey) -> Option<&iced::widget::image::Handle> {
|
fn handle(&self, key: &AttachmentKey) -> Option<&iced::widget::image::Handle> {
|
||||||
self.entries.get(key).and_then(|e| e.handle.as_ref())
|
self.entries
|
||||||
|
.get(key)
|
||||||
|
.and_then(|e| e.handle.as_ref().map(|h| &h.handle))
|
||||||
}
|
}
|
||||||
|
|
||||||
fn is_ready(&self, key: &AttachmentKey) -> bool {
|
fn is_ready(&self, key: &AttachmentKey) -> bool {
|
||||||
@@ -419,9 +511,16 @@ impl AttachmentCache {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Whether a fetch for this key is already running (dedup for clicks).
|
||||||
|
fn is_loading(&self, key: &AttachmentKey) -> bool {
|
||||||
|
matches!(self.get(key), Some(AttachmentState::Loading))
|
||||||
|
}
|
||||||
|
|
||||||
fn clear(&mut self) {
|
fn clear(&mut self) {
|
||||||
self.entries.clear();
|
self.entries.clear();
|
||||||
self.order.clear();
|
self.order.clear();
|
||||||
|
self.encoded_total = 0;
|
||||||
|
self.decoded_total = 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
@@ -430,6 +529,38 @@ impl AttachmentCache {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The decoded-budget weight of an entry: the preview's RGBA estimate for a
|
||||||
|
/// ready image (`w*h*4` via the handle's known preview dimensions), else zero.
|
||||||
|
/// The handle is built from the downscaled preview, so its dimensions are the
|
||||||
|
/// preview's — but iced's `Handle` doesn't expose them, so the caller threads
|
||||||
|
/// the cost through [`PreviewHandle`].
|
||||||
|
fn handle_decoded_cost(state: &AttachmentState, handle: &Option<PreviewHandle>) -> usize {
|
||||||
|
match (state, handle) {
|
||||||
|
(AttachmentState::Ready(_), Some(h)) => h.decoded_cost,
|
||||||
|
_ => 0,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// An iced image handle built from a downscaled preview, paired with the
|
||||||
|
/// preview's decoded RGBA cost (iced doesn't expose handle dimensions).
|
||||||
|
#[derive(Debug, Clone)]
|
||||||
|
struct PreviewHandle {
|
||||||
|
handle: iced::widget::image::Handle,
|
||||||
|
decoded_cost: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Build the inline preview handle (and its decoded weight) for image bytes, or
|
||||||
|
/// `None` when the bytes don't validate as a displayable image. The renderer
|
||||||
|
/// only ever receives the ≤1600 px downscale; originals stay encoded-only.
|
||||||
|
fn build_preview_handle(bytes: &[u8]) -> Option<PreviewHandle> {
|
||||||
|
let p = crate::files::decode_preview(bytes)?;
|
||||||
|
let decoded_cost = crate::files::preview_rgba_cost(p.width, p.height);
|
||||||
|
Some(PreviewHandle {
|
||||||
|
handle: iced::widget::image::Handle::from_rgba(p.width, p.height, p.rgba),
|
||||||
|
decoded_cost,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
/// Cap on retained chat history so a long call can't grow it without bound.
|
/// Cap on retained chat history so a long call can't grow it without bound.
|
||||||
const CHAT_HISTORY_MAX: usize = 300;
|
const CHAT_HISTORY_MAX: usize = 300;
|
||||||
|
|
||||||
@@ -612,11 +743,17 @@ pub enum AppMessage {
|
|||||||
/// Open the native picker to attach a file to the chat.
|
/// Open the native picker to attach a file to the chat.
|
||||||
PickAttachmentFile,
|
PickAttachmentFile,
|
||||||
/// Result of the attach picker: (filename, bytes), or None if cancelled.
|
/// Result of the attach picker: (filename, bytes), or None if cancelled.
|
||||||
AttachmentFilePicked(Option<(String, Vec<u8>)>),
|
/// Picker result: `None` = cancelled; `Some((name, None))` = file over the
|
||||||
|
/// attachment cap or unreadable; `Some((name, Some(bytes)))` = bytes in hand.
|
||||||
|
AttachmentFilePicked(Option<(String, Option<Vec<u8>>)>),
|
||||||
/// Save (downloading first if needed) a received attachment to disk. Carries
|
/// Save (downloading first if needed) a received attachment to disk. Carries
|
||||||
/// the full `(author, id)` key so the correct sender's bytes are fetched and
|
/// the full `(author, id)` key so the correct sender's bytes are fetched and
|
||||||
/// saved even if another peer reused the same attachment id (Tier C F-12).
|
/// saved even if another peer reused the same attachment id (Tier C F-12).
|
||||||
SaveAttachment(AttachmentKey),
|
SaveAttachment(AttachmentKey),
|
||||||
|
/// Fetch an image that wasn't auto-fetched (over the auto-size cap, budget-
|
||||||
|
/// skipped, or evicted) for inline display (Phase 3B). Manual, so it may use
|
||||||
|
/// the full attachment cap; duplicates are suppressed via the Loading state.
|
||||||
|
LoadImage(AttachmentKey),
|
||||||
/// Open the click-to-enlarge image lightbox for this attachment (author + id).
|
/// Open the click-to-enlarge image lightbox for this attachment (author + id).
|
||||||
OpenImageLightbox(AttachmentKey),
|
OpenImageLightbox(AttachmentKey),
|
||||||
/// Close the image lightbox overlay.
|
/// Close the image lightbox overlay.
|
||||||
@@ -1220,7 +1357,11 @@ impl Default for AppState {
|
|||||||
recording: false,
|
recording: false,
|
||||||
recording_started: None,
|
recording_started: None,
|
||||||
chat_messages: Vec::new(),
|
chat_messages: Vec::new(),
|
||||||
attachments: AttachmentCache::new(ATTACHMENT_CACHE_CAP),
|
attachments: AttachmentCache::new(
|
||||||
|
ATTACHMENT_CACHE_CAP,
|
||||||
|
ATTACHMENT_CACHE_ENCODED_BUDGET,
|
||||||
|
ATTACHMENT_CACHE_DECODED_BUDGET,
|
||||||
|
),
|
||||||
image_lightbox: None,
|
image_lightbox: None,
|
||||||
pending_saves: HashSet::new(),
|
pending_saves: HashSet::new(),
|
||||||
pending_plays: HashSet::new(),
|
pending_plays: HashSet::new(),
|
||||||
@@ -2021,25 +2162,40 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
UiEvent::AttachmentReady { from, id, data } => {
|
UiEvent::AttachmentReady { from, id, data } => {
|
||||||
// Bytes arrived for this specific (author, id). For images we
|
// Bytes arrived for this specific (author, id). For images
|
||||||
// can cache the iced handle now (built once, not per redraw).
|
// the ≤1600px preview handle is built once here (not per
|
||||||
// If the user was waiting to save this file, the save dialog
|
// redraw); only its downscale ever reaches the renderer.
|
||||||
// is opened from update() below by checking pending_saves —
|
|
||||||
// done lazily so this arm stays simple.
|
|
||||||
let key = (from, id);
|
let key = (from, id);
|
||||||
let handle = crate::files::validate_image_bytes(&data)
|
let handle = build_preview_handle(&data);
|
||||||
.is_some()
|
|
||||||
.then(|| iced::widget::image::Handle::from_bytes(data.clone()));
|
|
||||||
let needs_save = state.pending_saves.remove(&key);
|
let needs_save = state.pending_saves.remove(&key);
|
||||||
let needs_play = state.pending_plays.remove(&id);
|
let needs_play = state.pending_plays.remove(&id);
|
||||||
state
|
let retained =
|
||||||
.attachments
|
state
|
||||||
.insert(key, AttachmentState::Ready(data), handle);
|
.attachments
|
||||||
|
.insert(key, AttachmentState::Ready(data.clone()), handle);
|
||||||
|
if !retained {
|
||||||
|
// Individually over-budget: any immediate save/play is
|
||||||
|
// serviced from the bytes in hand below, and the entry
|
||||||
|
// is exposed as evicted, never exceeding the budgets.
|
||||||
|
state
|
||||||
|
.attachments
|
||||||
|
.insert(key, AttachmentState::Evicted, None);
|
||||||
|
}
|
||||||
if needs_play {
|
if needs_play {
|
||||||
play_ready_audio(state, key);
|
play_audio_data(state, id, &data);
|
||||||
}
|
}
|
||||||
if needs_save {
|
if needs_save {
|
||||||
return save_attachment_task(state, key);
|
return save_attachment_data_task(state, key, data);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
UiEvent::AttachmentFetchStarted { from, id } => {
|
||||||
|
// A fetch task exists for this key; show a real loading
|
||||||
|
// state (and dedup further clicks). Never demote Ready.
|
||||||
|
let key = (from, id);
|
||||||
|
if !state.attachments.is_ready(&key) {
|
||||||
|
state
|
||||||
|
.attachments
|
||||||
|
.insert(key, AttachmentState::Loading, None);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
UiEvent::AttachmentFailed { from, id, error } => {
|
UiEvent::AttachmentFailed { from, id, error } => {
|
||||||
@@ -2933,7 +3089,9 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
AppMessage::PickAttachmentFile => {
|
AppMessage::PickAttachmentFile => {
|
||||||
// Native picker off the UI thread; returns (filename, bytes).
|
// Native picker off the UI thread; returns (filename, bytes). The
|
||||||
|
// read is BOUNDED (metadata precheck + cap+1 read, Phase 3C) — an
|
||||||
|
// oversized pick reports `(name, None)` without buffering the file.
|
||||||
return Task::perform(
|
return Task::perform(
|
||||||
async {
|
async {
|
||||||
let handle = rfd::AsyncFileDialog::new()
|
let handle = rfd::AsyncFileDialog::new()
|
||||||
@@ -2941,7 +3099,18 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
.pick_file()
|
.pick_file()
|
||||||
.await;
|
.await;
|
||||||
match handle {
|
match handle {
|
||||||
Some(h) => Some((h.file_name(), h.read().await)),
|
Some(h) => {
|
||||||
|
let name = h.file_name();
|
||||||
|
let path = h.path().to_path_buf();
|
||||||
|
let bytes = tokio::task::spawn_blocking(move || {
|
||||||
|
crate::files::read_file_capped(&path)
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.ok()
|
||||||
|
.and_then(|r| r.ok())
|
||||||
|
.flatten();
|
||||||
|
Some((name, bytes))
|
||||||
|
}
|
||||||
None => None,
|
None => None,
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
@@ -2950,6 +3119,13 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
}
|
}
|
||||||
AppMessage::AttachmentFilePicked(picked) => {
|
AppMessage::AttachmentFilePicked(picked) => {
|
||||||
if let Some((name, bytes)) = picked {
|
if let Some((name, bytes)) = picked {
|
||||||
|
let Some(bytes) = bytes else {
|
||||||
|
state.status_message = format!(
|
||||||
|
"Can't attach {name} — unreadable or larger than {}.",
|
||||||
|
crate::files::human_size(crate::files::MAX_ATTACHMENT_BYTES)
|
||||||
|
);
|
||||||
|
return Task::none();
|
||||||
|
};
|
||||||
let size = bytes.len() as u64;
|
let size = bytes.len() as u64;
|
||||||
if !crate::files::size_within_cap(size) {
|
if !crate::files::size_within_cap(size) {
|
||||||
state.status_message = format!(
|
state.status_message = format!(
|
||||||
@@ -2958,6 +3134,7 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
);
|
);
|
||||||
return Task::none();
|
return Task::none();
|
||||||
}
|
}
|
||||||
|
let bytes = Arc::new(bytes);
|
||||||
let kind = crate::files::classify(&bytes);
|
let kind = crate::files::classify(&bytes);
|
||||||
// Random 32-byte handle for this attachment.
|
// Random 32-byte handle for this attachment.
|
||||||
let id: crate::files::AttachmentId = rand::random();
|
let id: crate::files::AttachmentId = rand::random();
|
||||||
@@ -2970,12 +3147,13 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
// Keep our own bytes locally so we see our own attachment inline
|
// Keep our own bytes locally so we see our own attachment inline
|
||||||
// immediately (others fetch it off the file plane). Cache under
|
// immediately (others fetch it off the file plane). Cache under
|
||||||
// our own (author, id) key so the render path — which keys the
|
// our own (author, id) key so the render path — which keys the
|
||||||
// line by `from` = self_id — finds them.
|
// line by `from` = self_id — finds them. The same Arc backs the
|
||||||
|
// cache, the command queue, and the serve store (Phase 3C).
|
||||||
if let Ok(self_eid) = state.self_id.parse::<EndpointId>() {
|
if let Ok(self_eid) = state.self_id.parse::<EndpointId>() {
|
||||||
let key = (self_eid, id);
|
let key = (self_eid, id);
|
||||||
let handle = (kind == crate::files::AttachmentKind::Image
|
let handle = (kind == crate::files::AttachmentKind::Image)
|
||||||
&& crate::files::validate_image_bytes(&bytes).is_some())
|
.then(|| build_preview_handle(&bytes))
|
||||||
.then(|| iced::widget::image::Handle::from_bytes(bytes.clone()));
|
.flatten();
|
||||||
state
|
state
|
||||||
.attachments
|
.attachments
|
||||||
.insert(key, AttachmentState::Ready(bytes.clone()), handle);
|
.insert(key, AttachmentState::Ready(bytes.clone()), handle);
|
||||||
@@ -2998,13 +3176,20 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
AppMessage::SaveAttachment(key) => {
|
AppMessage::SaveAttachment(key) => {
|
||||||
// If we already have the bytes, save now; otherwise fetch from the
|
// If we already have the bytes, save now; if a fetch is already
|
||||||
// sender and save when AttachmentReady arrives (pending_saves). The
|
// running, just record the save intent (no duplicate fetch task);
|
||||||
// key's author half is the exact sender of the clicked line.
|
// otherwise fetch from the sender and save when AttachmentReady
|
||||||
|
// arrives (pending_saves). The key's author half is the exact
|
||||||
|
// sender of the clicked line.
|
||||||
if state.attachments.is_ready(&key) {
|
if state.attachments.is_ready(&key) {
|
||||||
return save_attachment_task(state, key);
|
return save_attachment_task(state, key);
|
||||||
|
} else if state.attachments.is_loading(&key) {
|
||||||
|
state.pending_saves.insert(key);
|
||||||
} else if let Some(att) = find_attachment_source(state, key) {
|
} else if let Some(att) = find_attachment_source(state, key) {
|
||||||
state.pending_saves.insert(key);
|
state.pending_saves.insert(key);
|
||||||
|
state
|
||||||
|
.attachments
|
||||||
|
.insert(key, AttachmentState::Loading, None);
|
||||||
state.status_message = format!("Downloading {}…", att.name);
|
state.status_message = format!("Downloading {}…", att.name);
|
||||||
let _ = state.controller.send(CoreCommand::FetchAttachment {
|
let _ = state.controller.send(CoreCommand::FetchAttachment {
|
||||||
from: key.0,
|
from: key.0,
|
||||||
@@ -3012,6 +3197,23 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
AppMessage::LoadImage(key) => {
|
||||||
|
// Manual inline-image fetch for a skipped/evicted image. The
|
||||||
|
// Loading state both drives the label and suppresses duplicates.
|
||||||
|
if !state.attachments.is_ready(&key)
|
||||||
|
&& !state.attachments.is_loading(&key)
|
||||||
|
&& let Some(att) = find_attachment_source(state, key)
|
||||||
|
{
|
||||||
|
state
|
||||||
|
.attachments
|
||||||
|
.insert(key, AttachmentState::Loading, None);
|
||||||
|
state.status_message = format!("Loading {}…", att.name);
|
||||||
|
let _ = state.controller.send(CoreCommand::FetchAttachment {
|
||||||
|
from: key.0,
|
||||||
|
attachment: att,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
AppMessage::AttachmentSaved(msg) => {
|
AppMessage::AttachmentSaved(msg) => {
|
||||||
if let Some(m) = msg {
|
if let Some(m) = msg {
|
||||||
state.status_message = m;
|
state.status_message = m;
|
||||||
@@ -3020,11 +3222,18 @@ fn update(state: &mut AppState, message: AppMessage) -> Task<AppMessage> {
|
|||||||
AppMessage::PlayAudio(key) => {
|
AppMessage::PlayAudio(key) => {
|
||||||
if state.attachments.is_ready(&key) {
|
if state.attachments.is_ready(&key) {
|
||||||
play_ready_audio(state, key);
|
play_ready_audio(state, key);
|
||||||
|
} else if state.attachments.is_loading(&key) {
|
||||||
|
// A fetch is already running (e.g. a Save click); just record
|
||||||
|
// the play intent for when the bytes arrive.
|
||||||
|
state.pending_plays.insert(key.1);
|
||||||
} else if let Some(att) = find_attachment_source(state, key) {
|
} else if let Some(att) = find_attachment_source(state, key) {
|
||||||
// Repeated clicks while the transfer is pending must not launch
|
// Repeated clicks while the transfer is pending must not launch
|
||||||
// duplicate fetches. (pending_plays is id-keyed — it's coupled to
|
// duplicate fetches. (pending_plays is id-keyed — it's coupled to
|
||||||
// the id-keyed clip player; same-id collisions are cosmetic.)
|
// the id-keyed clip player; same-id collisions are cosmetic.)
|
||||||
if state.pending_plays.insert(key.1) {
|
if state.pending_plays.insert(key.1) {
|
||||||
|
state
|
||||||
|
.attachments
|
||||||
|
.insert(key, AttachmentState::Loading, None);
|
||||||
state.status_message = format!("Loading {}…", att.name);
|
state.status_message = format!("Loading {}…", att.name);
|
||||||
let _ = state.controller.send(CoreCommand::FetchAttachment {
|
let _ = state.controller.send(CoreCommand::FetchAttachment {
|
||||||
from: key.0,
|
from: key.0,
|
||||||
@@ -3650,9 +3859,17 @@ fn play_ready_audio(state: &mut AppState, key: AttachmentKey) {
|
|||||||
let Some(AttachmentState::Ready(data)) = state.attachments.get(&key) else {
|
let Some(AttachmentState::Ready(data)) = state.attachments.get(&key) else {
|
||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
let id = key.1;
|
let data = data.clone();
|
||||||
|
play_audio_data(state, key.1, &data);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Play audio from bytes in hand — used both for cached entries and for a
|
||||||
|
/// just-fetched attachment too large to retain (Phase 3A overweight path).
|
||||||
|
fn play_audio_data(state: &mut AppState, id: crate::files::AttachmentId, data: &Arc<Vec<u8>>) {
|
||||||
if crate::files::is_probably_audio(data) {
|
if crate::files::is_probably_audio(data) {
|
||||||
let bytes = data.clone();
|
// The clip player's command channel owns its bytes; one copy here at
|
||||||
|
// the click, the cache keeps the shared original.
|
||||||
|
let bytes = data.as_ref().clone();
|
||||||
state.invalid_audio.remove(&id);
|
state.invalid_audio.remove(&id);
|
||||||
state.clip_player.play(id, bytes);
|
state.clip_player.play(id, bytes);
|
||||||
// Apply this clip's effective gain; the command lands after Play so it
|
// Apply this clip's effective gain; the command lands after Play so it
|
||||||
@@ -4065,6 +4282,16 @@ fn save_attachment_task(state: &AppState, key: AttachmentKey) -> Task<AppMessage
|
|||||||
return Task::none();
|
return Task::none();
|
||||||
};
|
};
|
||||||
let data = data.clone();
|
let data = data.clone();
|
||||||
|
save_attachment_data_task(state, key, data)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The dialog+write half of a save, from bytes in hand — also used when a
|
||||||
|
/// just-fetched attachment was too large to retain in the cache (Phase 3A).
|
||||||
|
fn save_attachment_data_task(
|
||||||
|
state: &AppState,
|
||||||
|
key: AttachmentKey,
|
||||||
|
data: Arc<Vec<u8>>,
|
||||||
|
) -> Task<AppMessage> {
|
||||||
let default_name = attachment_default_name(&state.chat_messages, key);
|
let default_name = attachment_default_name(&state.chat_messages, key);
|
||||||
Task::perform(
|
Task::perform(
|
||||||
async move {
|
async move {
|
||||||
@@ -4074,7 +4301,7 @@ fn save_attachment_task(state: &AppState, key: AttachmentKey) -> Task<AppMessage
|
|||||||
.save_file()
|
.save_file()
|
||||||
.await;
|
.await;
|
||||||
match handle {
|
match handle {
|
||||||
Some(h) => match std::fs::write(h.path(), &data) {
|
Some(h) => match std::fs::write(h.path(), data.as_slice()) {
|
||||||
Ok(()) => Some(format!("Saved {}", h.path().display())),
|
Ok(()) => Some(format!("Saved {}", h.path().display())),
|
||||||
Err(e) => Some(format!("Save failed: {e}")),
|
Err(e) => Some(format!("Save failed: {e}")),
|
||||||
},
|
},
|
||||||
@@ -7124,10 +7351,33 @@ fn view(state: &AppState) -> Element<'_, AppMessage> {
|
|||||||
None => img.into(),
|
None => img.into(),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
None => text(format!("🖼 {} — loading…", att.name))
|
// "loading…" ONLY while a fetch task actually runs
|
||||||
|
// (Phase 3B). Absent/evicted images — over the
|
||||||
|
// auto-fetch cap, budget-skipped, or dropped by the
|
||||||
|
// cache — get an explicit Load button instead of an
|
||||||
|
// indefinite label.
|
||||||
|
None if matches!(data, Some(AttachmentState::Loading)) => {
|
||||||
|
text(format!("🖼 {} — loading…", att.name))
|
||||||
|
.size(12)
|
||||||
|
.color(color_subtext)
|
||||||
|
.into()
|
||||||
|
}
|
||||||
|
None => row![
|
||||||
|
text(format!(
|
||||||
|
"🖼 {} ({})",
|
||||||
|
att.name,
|
||||||
|
crate::files::human_size(att.size)
|
||||||
|
))
|
||||||
.size(12)
|
.size(12)
|
||||||
.color(color_subtext)
|
.color(color_text),
|
||||||
.into(),
|
button(text("Load image").size(12))
|
||||||
|
.on_press_maybe(key.map(AppMessage::LoadImage))
|
||||||
|
.style(b_style(color_blue, color_lavender, color_crust, 6.0,))
|
||||||
|
.padding(6),
|
||||||
|
]
|
||||||
|
.spacing(8)
|
||||||
|
.align_y(iced::alignment::Vertical::Center)
|
||||||
|
.into(),
|
||||||
}
|
}
|
||||||
} else if crate::files::looks_like_audio_name(&att.name)
|
} else if crate::files::looks_like_audio_name(&att.name)
|
||||||
&& !state.invalid_audio.contains(&att.id)
|
&& !state.invalid_audio.contains(&att.id)
|
||||||
@@ -7167,10 +7417,10 @@ fn view(state: &AppState) -> Element<'_, AppMessage> {
|
|||||||
.size(12)
|
.size(12)
|
||||||
.color(color_text),
|
.color(color_text),
|
||||||
button(
|
button(
|
||||||
text(if matches!(data, Some(AttachmentState::Ready(_))) {
|
text(match data {
|
||||||
"Save"
|
Some(AttachmentState::Ready(_)) => "Save",
|
||||||
} else {
|
Some(AttachmentState::Loading) => "Downloading…",
|
||||||
"Download"
|
_ => "Download",
|
||||||
})
|
})
|
||||||
.size(12)
|
.size(12)
|
||||||
)
|
)
|
||||||
@@ -7212,8 +7462,14 @@ fn view(state: &AppState) -> Element<'_, AppMessage> {
|
|||||||
.spacing(4)
|
.spacing(4)
|
||||||
.into()
|
.into()
|
||||||
} else {
|
} else {
|
||||||
let ready = matches!(data, Some(AttachmentState::Ready(_)));
|
// A repeated click during Loading only re-records the
|
||||||
let btn_label = if ready { "Save" } else { "Download" };
|
// save intent (the update arm dedups the fetch), so the
|
||||||
|
// button stays enabled with an honest label.
|
||||||
|
let btn_label = match data {
|
||||||
|
Some(AttachmentState::Ready(_)) => "Save",
|
||||||
|
Some(AttachmentState::Loading) => "Downloading…",
|
||||||
|
_ => "Download",
|
||||||
|
};
|
||||||
row![
|
row![
|
||||||
text(format!(
|
text(format!(
|
||||||
"📎 {} ({})",
|
"📎 {} ({})",
|
||||||
@@ -9018,17 +9274,35 @@ mod tests {
|
|||||||
(SecretKey::generate().public(), id)
|
(SecretKey::generate().public(), id)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A `Ready` state around bytes (tests only).
|
||||||
|
fn ready(bytes: Vec<u8>) -> AttachmentState {
|
||||||
|
AttachmentState::Ready(std::sync::Arc::new(bytes))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A cache with byte budgets too large to bind, isolating count-cap tests.
|
||||||
|
fn count_only_cache(cap: usize) -> AttachmentCache {
|
||||||
|
AttachmentCache::new(cap, usize::MAX, usize::MAX)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A dummy preview handle carrying an explicit decoded weight (tests only).
|
||||||
|
fn preview(decoded_cost: usize) -> super::PreviewHandle {
|
||||||
|
super::PreviewHandle {
|
||||||
|
handle: iced::widget::image::Handle::from_rgba(1, 1, vec![0u8; 4]),
|
||||||
|
decoded_cost,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn attachment_cache_evicts_oldest_when_full() {
|
fn attachment_cache_evicts_oldest_when_full() {
|
||||||
let mut cache = AttachmentCache::new(2);
|
let mut cache = count_only_cache(2);
|
||||||
let k1 = peer_key([1u8; 32]);
|
let k1 = peer_key([1u8; 32]);
|
||||||
let k2 = peer_key([2u8; 32]);
|
let k2 = peer_key([2u8; 32]);
|
||||||
let k3 = peer_key([3u8; 32]);
|
let k3 = peer_key([3u8; 32]);
|
||||||
cache.insert(k1, AttachmentState::Ready(vec![1]), None);
|
cache.insert(k1, ready(vec![1]), None);
|
||||||
cache.insert(k2, AttachmentState::Ready(vec![2]), None);
|
cache.insert(k2, ready(vec![2]), None);
|
||||||
assert_eq!(cache.len(), 2);
|
assert_eq!(cache.len(), 2);
|
||||||
// Inserting a third NEW key evicts the oldest (k1), not the newest.
|
// Inserting a third NEW key evicts the oldest (k1), not the newest.
|
||||||
cache.insert(k3, AttachmentState::Ready(vec![3]), None);
|
cache.insert(k3, ready(vec![3]), None);
|
||||||
assert_eq!(cache.len(), 2);
|
assert_eq!(cache.len(), 2);
|
||||||
assert!(cache.get(&k1).is_none(), "oldest should be evicted");
|
assert!(cache.get(&k1).is_none(), "oldest should be evicted");
|
||||||
assert!(cache.get(&k2).is_some());
|
assert!(cache.get(&k2).is_some());
|
||||||
@@ -9037,17 +9311,17 @@ mod tests {
|
|||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn attachment_cache_replace_keeps_position_and_count() {
|
fn attachment_cache_replace_keeps_position_and_count() {
|
||||||
let mut cache = AttachmentCache::new(2);
|
let mut cache = count_only_cache(2);
|
||||||
let k1 = peer_key([1u8; 32]);
|
let k1 = peer_key([1u8; 32]);
|
||||||
let k2 = peer_key([2u8; 32]);
|
let k2 = peer_key([2u8; 32]);
|
||||||
cache.insert(k1, AttachmentState::Failed("pending".into()), None);
|
cache.insert(k1, AttachmentState::Failed("pending".into()), None);
|
||||||
cache.insert(k2, AttachmentState::Ready(vec![2]), None);
|
cache.insert(k2, ready(vec![2]), None);
|
||||||
// Replacing k1 (Failed -> Ready) must NOT bump it to newest; it stays the
|
// Replacing k1 (Failed -> Ready) must NOT bump it to newest; it stays the
|
||||||
// eviction candidate, and the count is unchanged (no order leak).
|
// eviction candidate, and the count is unchanged (no order leak).
|
||||||
cache.insert(k1, AttachmentState::Ready(vec![1]), None);
|
cache.insert(k1, ready(vec![1]), None);
|
||||||
assert_eq!(cache.len(), 2);
|
assert_eq!(cache.len(), 2);
|
||||||
let k3 = peer_key([3u8; 32]);
|
let k3 = peer_key([3u8; 32]);
|
||||||
cache.insert(k3, AttachmentState::Ready(vec![3]), None);
|
cache.insert(k3, ready(vec![3]), None);
|
||||||
assert!(cache.get(&k1).is_none(), "replaced entry kept its old age");
|
assert!(cache.get(&k1).is_none(), "replaced entry kept its old age");
|
||||||
assert!(cache.get(&k2).is_some());
|
assert!(cache.get(&k2).is_some());
|
||||||
assert!(cache.get(&k3).is_some());
|
assert!(cache.get(&k3).is_some());
|
||||||
@@ -9057,26 +9331,26 @@ mod tests {
|
|||||||
fn attachment_cache_same_id_distinct_authors_do_not_alias() {
|
fn attachment_cache_same_id_distinct_authors_do_not_alias() {
|
||||||
// The F-12 core property: two peers sending the SAME attachment id keep
|
// The F-12 core property: two peers sending the SAME attachment id keep
|
||||||
// separate bytes — one never overwrites or aliases the other.
|
// separate bytes — one never overwrites or aliases the other.
|
||||||
let mut cache = AttachmentCache::new(8);
|
let mut cache = count_only_cache(8);
|
||||||
let shared_id = [7u8; 32];
|
let shared_id = [7u8; 32];
|
||||||
let victim = peer_key(shared_id);
|
let victim = peer_key(shared_id);
|
||||||
let attacker = peer_key(shared_id);
|
let attacker = peer_key(shared_id);
|
||||||
cache.insert(victim, AttachmentState::Ready(vec![1, 1, 1]), None);
|
cache.insert(victim, ready(vec![1, 1, 1]), None);
|
||||||
cache.insert(attacker, AttachmentState::Ready(vec![9, 9, 9]), None);
|
cache.insert(attacker, ready(vec![9, 9, 9]), None);
|
||||||
assert_eq!(cache.len(), 2);
|
assert_eq!(cache.len(), 2);
|
||||||
assert!(matches!(cache.get(&victim), Some(AttachmentState::Ready(b)) if b == &[1, 1, 1]));
|
assert!(
|
||||||
assert!(matches!(cache.get(&attacker), Some(AttachmentState::Ready(b)) if b == &[9, 9, 9]));
|
matches!(cache.get(&victim), Some(AttachmentState::Ready(b)) if b.as_slice() == [1, 1, 1])
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
matches!(cache.get(&attacker), Some(AttachmentState::Ready(b)) if b.as_slice() == [9, 9, 9])
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn attachment_cache_is_ready_and_handle_and_clear() {
|
fn attachment_cache_is_ready_and_handle_and_clear() {
|
||||||
let mut cache = AttachmentCache::new(4);
|
let mut cache = count_only_cache(4);
|
||||||
let k = peer_key([5u8; 32]);
|
let k = peer_key([5u8; 32]);
|
||||||
cache.insert(
|
cache.insert(k, ready(vec![1]), Some(preview(4)));
|
||||||
k,
|
|
||||||
AttachmentState::Ready(vec![1]),
|
|
||||||
Some(iced::widget::image::Handle::from_bytes(vec![1])),
|
|
||||||
);
|
|
||||||
assert!(cache.is_ready(&k));
|
assert!(cache.is_ready(&k));
|
||||||
assert!(cache.handle(&k).is_some());
|
assert!(cache.handle(&k).is_some());
|
||||||
let failed = peer_key([6u8; 32]);
|
let failed = peer_key([6u8; 32]);
|
||||||
@@ -9086,17 +9360,105 @@ mod tests {
|
|||||||
cache.clear();
|
cache.clear();
|
||||||
assert_eq!(cache.len(), 0);
|
assert_eq!(cache.len(), 0);
|
||||||
assert!(cache.get(&k).is_none());
|
assert!(cache.get(&k).is_none());
|
||||||
|
assert_eq!(cache.encoded_total, 0);
|
||||||
|
assert_eq!(cache.decoded_total, 0);
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn attachment_cache_cap_zero_is_clamped_to_one() {
|
fn attachment_cache_cap_zero_is_clamped_to_one() {
|
||||||
let mut cache = AttachmentCache::new(0);
|
let mut cache = count_only_cache(0);
|
||||||
let k = peer_key([1u8; 32]);
|
let k = peer_key([1u8; 32]);
|
||||||
cache.insert(k, AttachmentState::Ready(vec![1]), None);
|
cache.insert(k, ready(vec![1]), None);
|
||||||
assert_eq!(cache.len(), 1);
|
assert_eq!(cache.len(), 1);
|
||||||
assert!(cache.get(&k).is_some());
|
assert!(cache.get(&k).is_some());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn attachment_cache_encoded_budget_evicts_oldest() {
|
||||||
|
// Budget of 100 encoded bytes; entries of 60 + 30 fit, adding 40 must
|
||||||
|
// evict the oldest (60) — not fail, not exceed the budget.
|
||||||
|
let mut cache = AttachmentCache::new(64, 100, usize::MAX);
|
||||||
|
let k1 = peer_key([1u8; 32]);
|
||||||
|
let k2 = peer_key([2u8; 32]);
|
||||||
|
let k3 = peer_key([3u8; 32]);
|
||||||
|
assert!(cache.insert(k1, ready(vec![0; 60]), None));
|
||||||
|
assert!(cache.insert(k2, ready(vec![0; 30]), None));
|
||||||
|
assert_eq!(cache.encoded_total, 90);
|
||||||
|
assert!(cache.insert(k3, ready(vec![0; 40]), None));
|
||||||
|
assert!(cache.get(&k1).is_none(), "oldest evicted for byte budget");
|
||||||
|
assert!(cache.get(&k2).is_some());
|
||||||
|
assert!(cache.get(&k3).is_some());
|
||||||
|
assert_eq!(cache.encoded_total, 70);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn attachment_cache_decoded_budget_evicts_oldest() {
|
||||||
|
let mut cache = AttachmentCache::new(64, usize::MAX, 100);
|
||||||
|
let k1 = peer_key([1u8; 32]);
|
||||||
|
let k2 = peer_key([2u8; 32]);
|
||||||
|
assert!(cache.insert(k1, ready(vec![1]), Some(preview(80))));
|
||||||
|
assert!(cache.insert(k2, ready(vec![2]), Some(preview(40))));
|
||||||
|
assert!(
|
||||||
|
cache.get(&k1).is_none(),
|
||||||
|
"oldest evicted for decoded budget"
|
||||||
|
);
|
||||||
|
assert!(cache.get(&k2).is_some());
|
||||||
|
assert_eq!(cache.decoded_total, 40);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn attachment_cache_replacement_subtracts_old_weight_first() {
|
||||||
|
// k1 holds 80 of a 100-byte budget. Replacing it with 90 must fit
|
||||||
|
// WITHOUT evicting k2 (10): the old 80 is subtracted before the check.
|
||||||
|
let mut cache = AttachmentCache::new(64, 100, usize::MAX);
|
||||||
|
let k1 = peer_key([1u8; 32]);
|
||||||
|
let k2 = peer_key([2u8; 32]);
|
||||||
|
assert!(cache.insert(k1, ready(vec![0; 80]), None));
|
||||||
|
assert!(cache.insert(k2, ready(vec![0; 10]), None));
|
||||||
|
assert!(cache.insert(k1, ready(vec![0; 90]), None));
|
||||||
|
assert!(cache.get(&k2).is_some(), "replacement must not over-evict");
|
||||||
|
assert_eq!(cache.encoded_total, 100);
|
||||||
|
// A non-Ready replacement releases the weight entirely.
|
||||||
|
cache.insert(k1, AttachmentState::Evicted, None);
|
||||||
|
assert_eq!(cache.encoded_total, 10);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn attachment_cache_overweight_entry_is_not_retained() {
|
||||||
|
let mut cache = AttachmentCache::new(64, 100, usize::MAX);
|
||||||
|
let k1 = peer_key([1u8; 32]);
|
||||||
|
let k2 = peer_key([2u8; 32]);
|
||||||
|
assert!(cache.insert(k2, ready(vec![0; 50]), None));
|
||||||
|
// Individually larger than the whole budget: rejected, nothing evicted.
|
||||||
|
assert!(!cache.insert(k1, ready(vec![0; 101]), None));
|
||||||
|
assert!(cache.get(&k1).is_none());
|
||||||
|
assert!(cache.get(&k2).is_some(), "others untouched by the reject");
|
||||||
|
assert_eq!(cache.encoded_total, 50);
|
||||||
|
// Rejecting a REPLACEMENT drops the stale entry too (its old bytes
|
||||||
|
// must not linger under a key the caller will mark Evicted).
|
||||||
|
assert!(!cache.insert(k2, ready(vec![0; 200]), None));
|
||||||
|
assert!(cache.get(&k2).is_none());
|
||||||
|
assert_eq!(cache.encoded_total, 0);
|
||||||
|
// The caller's follow-up Evicted marker is retained (zero weight).
|
||||||
|
assert!(cache.insert(k1, AttachmentState::Evicted, None));
|
||||||
|
assert!(matches!(cache.get(&k1), Some(AttachmentState::Evicted)));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn attachment_cache_loading_states_are_weightless_and_tracked() {
|
||||||
|
let mut cache = AttachmentCache::new(64, 100, 100);
|
||||||
|
let k = peer_key([1u8; 32]);
|
||||||
|
assert!(cache.insert(k, AttachmentState::Loading, None));
|
||||||
|
assert!(cache.is_loading(&k));
|
||||||
|
assert!(!cache.is_ready(&k));
|
||||||
|
assert_eq!(cache.encoded_total, 0);
|
||||||
|
// Ready replaces Loading, weights counted; not loading any more.
|
||||||
|
assert!(cache.insert(k, ready(vec![0; 10]), None));
|
||||||
|
assert!(!cache.is_loading(&k));
|
||||||
|
assert!(cache.is_ready(&k));
|
||||||
|
assert_eq!(cache.encoded_total, 10);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn attachment_default_name_matches_full_key_not_bare_id() {
|
fn attachment_default_name_matches_full_key_not_bare_id() {
|
||||||
// Two chat lines carry the SAME attachment id but come from different
|
// Two chat lines carry the SAME attachment id but come from different
|
||||||
@@ -9166,11 +9528,9 @@ mod tests {
|
|||||||
});
|
});
|
||||||
state.chat_input = "draft".to_string();
|
state.chat_input = "draft".to_string();
|
||||||
let att_key = (peer, attachment_id);
|
let att_key = (peer, attachment_id);
|
||||||
state.attachments.insert(
|
state
|
||||||
att_key,
|
.attachments
|
||||||
AttachmentState::Ready(vec![1]),
|
.insert(att_key, ready(vec![1]), Some(preview(4)));
|
||||||
Some(iced::widget::image::Handle::from_bytes(vec![1])),
|
|
||||||
);
|
|
||||||
state.pending_saves.insert(att_key);
|
state.pending_saves.insert(att_key);
|
||||||
state.pending_plays.insert(attachment_id);
|
state.pending_plays.insert(attachment_id);
|
||||||
state.invalid_audio.insert(attachment_id);
|
state.invalid_audio.insert(attachment_id);
|
||||||
|
|||||||
@@ -0,0 +1,252 @@
|
|||||||
|
//! Byte/request budgets for AUTOMATIC chat-attachment fetches (Phase 3B).
|
||||||
|
//!
|
||||||
|
//! The four-permit semaphore bounds how many auto-fetch tasks run at once, but
|
||||||
|
//! not how much a peer can make us download over time: with permits released
|
||||||
|
//! after each transfer, an insider could stream distinct ≤4 MiB images
|
||||||
|
//! sequentially forever. This budget adds per-author and session (room-wide)
|
||||||
|
//! token buckets over both request COUNT and declared BYTES. Like the Phase 2
|
||||||
|
//! chat gate, time is passed in — never read from a clock — so every refill
|
||||||
|
//! boundary is unit-testable.
|
||||||
|
//!
|
||||||
|
//! Only the automatic path consults this; a user's explicit click (Save /
|
||||||
|
//! Download / Load image) is human-rate-limited and always allowed through to
|
||||||
|
//! the fetch (still subject to the transfer cap and cache/decoder budgets).
|
||||||
|
|
||||||
|
use iroh::EndpointId;
|
||||||
|
use std::collections::HashMap;
|
||||||
|
|
||||||
|
/// Per-author request burst: how many auto-fetches one author can trigger
|
||||||
|
/// back-to-back before refill pacing binds.
|
||||||
|
pub const AUTHOR_REQ_BURST: f64 = 8.0;
|
||||||
|
/// Per-author request refill: one recovered every 10 s.
|
||||||
|
pub const AUTHOR_REQ_REFILL_PER_MS: f64 = 1.0 / 10_000.0;
|
||||||
|
/// Per-author byte burst (declared sizes): a couple of full-size auto images
|
||||||
|
/// plus a normal working set.
|
||||||
|
pub const AUTHOR_BYTES_BURST: f64 = (16 * 1024 * 1024) as f64;
|
||||||
|
/// Per-author byte refill: 64 KiB/s (~one 4 MiB auto image per minute).
|
||||||
|
pub const AUTHOR_BYTES_REFILL_PER_MS: f64 = (64 * 1024) as f64 / 1000.0;
|
||||||
|
|
||||||
|
/// Session-wide request burst across all authors.
|
||||||
|
pub const SESSION_REQ_BURST: f64 = 16.0;
|
||||||
|
/// Session-wide request refill: one recovered every 5 s.
|
||||||
|
pub const SESSION_REQ_REFILL_PER_MS: f64 = 1.0 / 5_000.0;
|
||||||
|
/// Session-wide byte burst across all authors.
|
||||||
|
pub const SESSION_BYTES_BURST: f64 = (48 * 1024 * 1024) as f64;
|
||||||
|
/// Session-wide byte refill: 128 KiB/s.
|
||||||
|
pub const SESSION_BYTES_REFILL_PER_MS: f64 = (128 * 1024) as f64 / 1000.0;
|
||||||
|
|
||||||
|
/// Bound on the per-author bucket map. Authors are roster members (≤32 live),
|
||||||
|
/// so this tracks the roster plus recently departed; the least-recently-active
|
||||||
|
/// entry is pruned past the cap.
|
||||||
|
pub const AUTHOR_MAP_CAP: usize = 64;
|
||||||
|
|
||||||
|
/// A deterministic token bucket that can take a WEIGHTED cost (bytes), unlike
|
||||||
|
/// the unit-cost bucket in the gossip chat gate.
|
||||||
|
#[derive(Debug, Clone, Copy)]
|
||||||
|
struct WeightedBucket {
|
||||||
|
tokens: f64,
|
||||||
|
last_ms: u64,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl WeightedBucket {
|
||||||
|
fn full(burst: f64, now_ms: u64) -> Self {
|
||||||
|
Self {
|
||||||
|
tokens: burst,
|
||||||
|
last_ms: now_ms,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Refill for elapsed time (capped at `burst`) without consuming.
|
||||||
|
fn refill(&mut self, burst: f64, refill_per_ms: f64, now_ms: u64) {
|
||||||
|
let elapsed = now_ms.saturating_sub(self.last_ms) as f64;
|
||||||
|
self.tokens = (self.tokens + elapsed * refill_per_ms).min(burst);
|
||||||
|
self.last_ms = now_ms;
|
||||||
|
}
|
||||||
|
|
||||||
|
fn has(&self, cost: f64) -> bool {
|
||||||
|
self.tokens >= cost
|
||||||
|
}
|
||||||
|
|
||||||
|
fn take(&mut self, cost: f64) {
|
||||||
|
self.tokens -= cost;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// One author's pair of buckets plus last activity (for idle pruning).
|
||||||
|
#[derive(Debug)]
|
||||||
|
struct AuthorBudget {
|
||||||
|
reqs: WeightedBucket,
|
||||||
|
bytes: WeightedBucket,
|
||||||
|
last_seen_ms: u64,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Admission budget for automatic attachment fetches. All four buckets are
|
||||||
|
/// checked BEFORE any is consumed, so a rejection never burns tokens (no
|
||||||
|
/// refund bookkeeping — the check-then-take is atomic within `admit`).
|
||||||
|
#[derive(Debug)]
|
||||||
|
pub struct AutoFetchBudget {
|
||||||
|
session_reqs: WeightedBucket,
|
||||||
|
session_bytes: WeightedBucket,
|
||||||
|
authors: HashMap<EndpointId, AuthorBudget>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl AutoFetchBudget {
|
||||||
|
pub fn new(now_ms: u64) -> Self {
|
||||||
|
Self {
|
||||||
|
session_reqs: WeightedBucket::full(SESSION_REQ_BURST, now_ms),
|
||||||
|
session_bytes: WeightedBucket::full(SESSION_BYTES_BURST, now_ms),
|
||||||
|
authors: HashMap::new(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Whether an auto-fetch of `size` declared bytes for `author` may start
|
||||||
|
/// now. Consumes one request token and `size` byte tokens from BOTH the
|
||||||
|
/// author's and the session's buckets — or nothing at all on rejection.
|
||||||
|
pub fn admit(&mut self, author: EndpointId, size: u64, now_ms: u64) -> bool {
|
||||||
|
self.prune(author, now_ms);
|
||||||
|
let entry = self.authors.entry(author).or_insert_with(|| AuthorBudget {
|
||||||
|
reqs: WeightedBucket::full(AUTHOR_REQ_BURST, now_ms),
|
||||||
|
bytes: WeightedBucket::full(AUTHOR_BYTES_BURST, now_ms),
|
||||||
|
last_seen_ms: now_ms,
|
||||||
|
});
|
||||||
|
entry.last_seen_ms = now_ms;
|
||||||
|
entry
|
||||||
|
.reqs
|
||||||
|
.refill(AUTHOR_REQ_BURST, AUTHOR_REQ_REFILL_PER_MS, now_ms);
|
||||||
|
entry
|
||||||
|
.bytes
|
||||||
|
.refill(AUTHOR_BYTES_BURST, AUTHOR_BYTES_REFILL_PER_MS, now_ms);
|
||||||
|
self.session_reqs
|
||||||
|
.refill(SESSION_REQ_BURST, SESSION_REQ_REFILL_PER_MS, now_ms);
|
||||||
|
self.session_bytes
|
||||||
|
.refill(SESSION_BYTES_BURST, SESSION_BYTES_REFILL_PER_MS, now_ms);
|
||||||
|
|
||||||
|
let cost = size as f64;
|
||||||
|
let ok = entry.reqs.has(1.0)
|
||||||
|
&& entry.bytes.has(cost)
|
||||||
|
&& self.session_reqs.has(1.0)
|
||||||
|
&& self.session_bytes.has(cost);
|
||||||
|
if ok {
|
||||||
|
let entry = self.authors.get_mut(&author).expect("just inserted");
|
||||||
|
entry.reqs.take(1.0);
|
||||||
|
entry.bytes.take(cost);
|
||||||
|
self.session_reqs.take(1.0);
|
||||||
|
self.session_bytes.take(cost);
|
||||||
|
}
|
||||||
|
ok
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Keep the author map bounded: past the cap, drop the least-recently
|
||||||
|
/// active entry that isn't the author being admitted. A pruned author
|
||||||
|
/// returns with full buckets, but authors are roster-gated upstream, so
|
||||||
|
/// the map can't be churned by strangers.
|
||||||
|
fn prune(&mut self, keep: EndpointId, _now_ms: u64) {
|
||||||
|
while self.authors.len() >= AUTHOR_MAP_CAP {
|
||||||
|
let Some(victim) = self
|
||||||
|
.authors
|
||||||
|
.iter()
|
||||||
|
.filter(|(id, _)| **id != keep)
|
||||||
|
.min_by_key(|(_, b)| b.last_seen_ms)
|
||||||
|
.map(|(id, _)| *id)
|
||||||
|
else {
|
||||||
|
break;
|
||||||
|
};
|
||||||
|
self.authors.remove(&victim);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
fn author_count(&self) -> usize {
|
||||||
|
self.authors.len()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
use iroh::SecretKey;
|
||||||
|
|
||||||
|
const T0: u64 = 1_000_000;
|
||||||
|
const MIB: u64 = 1024 * 1024;
|
||||||
|
|
||||||
|
fn author() -> EndpointId {
|
||||||
|
SecretKey::generate().public()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn author_request_burst_then_refill_recovers() {
|
||||||
|
let mut b = AutoFetchBudget::new(T0);
|
||||||
|
let a = author();
|
||||||
|
// Tiny sizes so only the REQUEST buckets can bind.
|
||||||
|
for _ in 0..AUTHOR_REQ_BURST as usize {
|
||||||
|
assert!(b.admit(a, 1, T0));
|
||||||
|
}
|
||||||
|
assert!(!b.admit(a, 1, T0), "author request burst exhausted");
|
||||||
|
// One request refills after 10 s.
|
||||||
|
assert!(b.admit(a, 1, T0 + 10_000));
|
||||||
|
assert!(!b.admit(a, 1, T0 + 10_000));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn author_byte_budget_binds_and_recovers() {
|
||||||
|
let mut b = AutoFetchBudget::new(T0);
|
||||||
|
let a = author();
|
||||||
|
// 4 × 4 MiB = the full 16 MiB author byte burst (well under the
|
||||||
|
// 8-request burst, so bytes are the binding constraint).
|
||||||
|
for _ in 0..4 {
|
||||||
|
assert!(b.admit(a, 4 * MIB, T0));
|
||||||
|
}
|
||||||
|
assert!(!b.admit(a, 4 * MIB, T0), "author byte burst exhausted");
|
||||||
|
// 64 KiB/s → a 4 MiB image is affordable again after 64 s (which also
|
||||||
|
// refills 6 request tokens, so bytes stay the binding constraint).
|
||||||
|
assert!(!b.admit(a, 4 * MIB, T0 + 32_000));
|
||||||
|
assert!(b.admit(a, 4 * MIB, T0 + 64_000));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn session_budget_binds_across_authors_without_burning_author_tokens() {
|
||||||
|
let mut b = AutoFetchBudget::new(T0);
|
||||||
|
// Three authors × 16 MiB exhausts the 48 MiB session byte burst even
|
||||||
|
// though each author is within their own budget.
|
||||||
|
for _ in 0..3 {
|
||||||
|
let a = author();
|
||||||
|
for _ in 0..4 {
|
||||||
|
assert!(b.admit(a, 4 * MIB, T0));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
let fresh = author();
|
||||||
|
assert!(!b.admit(fresh, 4 * MIB, T0), "session bytes exhausted");
|
||||||
|
// The rejection consumed NOTHING: once the session refills enough for
|
||||||
|
// one image (4 MiB / 128 KiB/s = 32 s), the fresh author's own full
|
||||||
|
// burst is intact and admits immediately.
|
||||||
|
assert!(b.admit(fresh, 4 * MIB, T0 + 32_000));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn session_request_bucket_binds_across_authors() {
|
||||||
|
let mut b = AutoFetchBudget::new(T0);
|
||||||
|
// 16 tiny requests from distinct authors exhaust the session request
|
||||||
|
// burst while every author bucket stays nearly full.
|
||||||
|
for _ in 0..SESSION_REQ_BURST as usize {
|
||||||
|
assert!(b.admit(author(), 1, T0));
|
||||||
|
}
|
||||||
|
assert!(!b.admit(author(), 1, T0), "session requests exhausted");
|
||||||
|
assert!(b.admit(author(), 1, T0 + 5_000), "one recovers after 5 s");
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn author_map_stays_bounded_pruning_least_recent() {
|
||||||
|
let mut b = AutoFetchBudget::new(T0);
|
||||||
|
// Session request refill would bind over a naive loop; space the
|
||||||
|
// admissions out so only the map bound is under test.
|
||||||
|
let mut t = T0;
|
||||||
|
let first = author();
|
||||||
|
assert!(b.admit(first, 1, t));
|
||||||
|
for _ in 0..(AUTHOR_MAP_CAP + 10) {
|
||||||
|
t += 10_000;
|
||||||
|
assert!(b.admit(author(), 1, t));
|
||||||
|
assert!(b.author_count() <= AUTHOR_MAP_CAP);
|
||||||
|
}
|
||||||
|
assert!(b.author_count() <= AUTHOR_MAP_CAP);
|
||||||
|
}
|
||||||
|
}
|
||||||
+13
-2
@@ -74,7 +74,10 @@ pub enum CoreCommand {
|
|||||||
SendChatFile {
|
SendChatFile {
|
||||||
text: String,
|
text: String,
|
||||||
attachment: crate::files::ChatAttachment,
|
attachment: crate::files::ChatAttachment,
|
||||||
data: Vec<u8>,
|
/// Shared, not owned: the same allocation is retained by the UI cache
|
||||||
|
/// and handed to the serve store, so a 25 MiB attachment is held once,
|
||||||
|
/// not copied across UI / command queue / serve store (Phase 3C).
|
||||||
|
data: std::sync::Arc<Vec<u8>>,
|
||||||
},
|
},
|
||||||
/// Fetch a received attachment's bytes from its sender over the file plane
|
/// Fetch a received attachment's bytes from its sender over the file plane
|
||||||
/// (used for on-demand file/chip downloads; images are auto-fetched on
|
/// (used for on-demand file/chip downloads; images are auto-fetched on
|
||||||
@@ -434,7 +437,15 @@ pub enum UiEvent {
|
|||||||
AttachmentReady {
|
AttachmentReady {
|
||||||
from: EndpointId,
|
from: EndpointId,
|
||||||
id: crate::files::AttachmentId,
|
id: crate::files::AttachmentId,
|
||||||
data: Vec<u8>,
|
data: std::sync::Arc<Vec<u8>>,
|
||||||
|
},
|
||||||
|
/// An attachment fetch task was spawned (auto or on demand). Lets the UI
|
||||||
|
/// show a real "loading" state instead of inferring it from cache absence —
|
||||||
|
/// absence now means NOT fetched (e.g. auto-fetch was skipped), which
|
||||||
|
/// renders a Load button rather than an indefinite "loading…" (Phase 3B).
|
||||||
|
AttachmentFetchStarted {
|
||||||
|
from: EndpointId,
|
||||||
|
id: crate::files::AttachmentId,
|
||||||
},
|
},
|
||||||
/// An attachment fetch failed (sender gone, too large, decode error, etc.).
|
/// An attachment fetch failed (sender gone, too large, decode error, etc.).
|
||||||
AttachmentFailed {
|
AttachmentFailed {
|
||||||
|
|||||||
+61
-15
@@ -1,5 +1,6 @@
|
|||||||
pub mod chatroster;
|
pub mod chatroster;
|
||||||
pub mod connstats;
|
pub mod connstats;
|
||||||
|
pub mod fetchbudget;
|
||||||
pub mod jitter;
|
pub mod jitter;
|
||||||
pub mod messages;
|
pub mod messages;
|
||||||
mod recovery;
|
mod recovery;
|
||||||
@@ -1017,6 +1018,15 @@ const MAX_INFLIGHT_ATTACHMENT_FETCHES: usize = 4;
|
|||||||
/// fetches (Tier C F-02). Bounded by [`MAX_INFLIGHT_ATTACHMENT_FETCHES`].
|
/// fetches (Tier C F-02). Bounded by [`MAX_INFLIGHT_ATTACHMENT_FETCHES`].
|
||||||
type InflightAttachments = Arc<std::sync::Mutex<HashSet<(EndpointId, crate::files::AttachmentId)>>>;
|
type InflightAttachments = Arc<std::sync::Mutex<HashSet<(EndpointId, crate::files::AttachmentId)>>>;
|
||||||
|
|
||||||
|
/// Milliseconds since the Unix epoch — the time source handed to the
|
||||||
|
/// deterministic auto-fetch budget (mirrors the gossip plane's timestamps).
|
||||||
|
fn unix_now_ms() -> u64 {
|
||||||
|
std::time::SystemTime::now()
|
||||||
|
.duration_since(std::time::UNIX_EPOCH)
|
||||||
|
.map(|d| d.as_millis() as u64)
|
||||||
|
.unwrap_or(0)
|
||||||
|
}
|
||||||
|
|
||||||
/// RAII bookkeeping for one bounded auto-fetch: holds the concurrency permit for
|
/// RAII bookkeeping for one bounded auto-fetch: holds the concurrency permit for
|
||||||
/// the task's lifetime and clears the in-flight `(author, id)` marker when the
|
/// the task's lifetime and clears the in-flight `(author, id)` marker when the
|
||||||
/// fetch finishes (success OR failure), so the same image can be retried later.
|
/// fetch finishes (success OR failure), so the same image can be retried later.
|
||||||
@@ -1033,11 +1043,22 @@ impl Drop for AutoFetchGuard {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Whether to AUTO-fetch a chat image attachment. Only authenticated roster
|
/// Whether to AUTO-fetch a chat image attachment. Only authenticated roster
|
||||||
/// authors qualify (closing the non-roster injection vector), and a `(author,
|
/// authors qualify (closing the non-roster injection vector), only declared
|
||||||
/// id)` already being fetched is skipped (dedup). The concurrency bound itself is
|
/// sizes at or under [`crate::files::MAX_AUTO_IMAGE_BYTES`] (larger images get
|
||||||
/// enforced separately by the permit. Pure → unit-testable (Tier C F-02).
|
/// a Load button instead, Phase 3B), and a `(author, id)` already being fetched
|
||||||
fn should_auto_fetch(is_image: bool, author_in_roster: bool, already_inflight: bool) -> bool {
|
/// is skipped (dedup). The concurrency bound (permits) and the byte/request
|
||||||
is_image && author_in_roster && !already_inflight
|
/// budgets ([`fetchbudget::AutoFetchBudget`]) are enforced separately. Pure →
|
||||||
|
/// unit-testable (Tier C F-02).
|
||||||
|
fn should_auto_fetch(
|
||||||
|
is_image: bool,
|
||||||
|
author_in_roster: bool,
|
||||||
|
already_inflight: bool,
|
||||||
|
declared_size: u64,
|
||||||
|
) -> bool {
|
||||||
|
is_image
|
||||||
|
&& author_in_roster
|
||||||
|
&& !already_inflight
|
||||||
|
&& declared_size <= crate::files::MAX_AUTO_IMAGE_BYTES
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Fetch a chat attachment's bytes from `from` over the file plane in a detached
|
/// Fetch a chat attachment's bytes from `from` over the file plane in a detached
|
||||||
@@ -1060,6 +1081,12 @@ fn spawn_attachment_fetch(
|
|||||||
tokio::spawn(async move {
|
tokio::spawn(async move {
|
||||||
// Held for the whole fetch; dropped here on completion (Tier C F-02).
|
// Held for the whole fetch; dropped here on completion (Tier C F-02).
|
||||||
let _guard = guard;
|
let _guard = guard;
|
||||||
|
// Tell the UI a real fetch task exists for this key, so it can show a
|
||||||
|
// genuine loading state and dedup further clicks (Phase 3B). Sent from
|
||||||
|
// this task's channel handle, so it always precedes Ready/Failed.
|
||||||
|
let _ = ui_tx
|
||||||
|
.send(UiEvent::AttachmentFetchStarted { from, id: att.id })
|
||||||
|
.await;
|
||||||
match transport.fetch_attachment(from, &att).await {
|
match transport.fetch_attachment(from, &att).await {
|
||||||
Ok(data) => {
|
Ok(data) => {
|
||||||
if is_image && crate::files::validate_image_bytes(&data).is_none() {
|
if is_image && crate::files::validate_image_bytes(&data).is_none() {
|
||||||
@@ -1076,7 +1103,7 @@ fn spawn_attachment_fetch(
|
|||||||
.send(UiEvent::AttachmentReady {
|
.send(UiEvent::AttachmentReady {
|
||||||
from,
|
from,
|
||||||
id: att.id,
|
id: att.id,
|
||||||
data,
|
data: Arc::new(data),
|
||||||
})
|
})
|
||||||
.await;
|
.await;
|
||||||
}
|
}
|
||||||
@@ -2362,6 +2389,10 @@ async fn run_core_loop(
|
|||||||
let chat_roster_events = chat_roster.clone();
|
let chat_roster_events = chat_roster.clone();
|
||||||
let event_task = tokio::spawn(async move {
|
let event_task = tokio::spawn(async move {
|
||||||
let roster = chat_roster_events;
|
let roster = chat_roster_events;
|
||||||
|
// Byte/request budgets for automatic attachment fetches
|
||||||
|
// (Phase 3B). Only this sequential task consults it, so it
|
||||||
|
// needs no lock; time is passed in for testability.
|
||||||
|
let mut auto_fetch_budget = fetchbudget::AutoFetchBudget::new(unix_now_ms());
|
||||||
while let Some(event) = room_events.recv().await {
|
while let Some(event) = room_events.recv().await {
|
||||||
match event {
|
match event {
|
||||||
RoomEvent::PeerJoined(peer_id, state) => {
|
RoomEvent::PeerJoined(peer_id, state) => {
|
||||||
@@ -2532,7 +2563,9 @@ async fn run_core_loop(
|
|||||||
// any attachment handling.
|
// any attachment handling.
|
||||||
true,
|
true,
|
||||||
already_inflight,
|
already_inflight,
|
||||||
) {
|
att.size,
|
||||||
|
) && auto_fetch_budget.admit(from, att.size, unix_now_ms())
|
||||||
|
{
|
||||||
// Reserve the dedup slot, then a permit. If the
|
// Reserve the dedup slot, then a permit. If the
|
||||||
// pool is exhausted, drop the auto-fetch (and the
|
// pool is exhausted, drop the auto-fetch (and the
|
||||||
// dedup marker) — the descriptor still shows and
|
// dedup marker) — the descriptor still shows and
|
||||||
@@ -3242,9 +3275,7 @@ async fn run_core_loop(
|
|||||||
if let Some(session) = &active_session {
|
if let Some(session) = &active_session {
|
||||||
// Make the bytes fetchable by room members, then broadcast the
|
// Make the bytes fetchable by room members, then broadcast the
|
||||||
// descriptor alongside the (possibly empty) caption text.
|
// descriptor alongside the (possibly empty) caption text.
|
||||||
session
|
session.transport.serve_attachment(attachment.id, data);
|
||||||
.transport
|
|
||||||
.serve_attachment(attachment.id, Arc::new(data));
|
|
||||||
if let Err(e) = session.room_state.send_chat(text, Some(attachment)).await {
|
if let Err(e) = session.room_state.send_chat(text, Some(attachment)).await {
|
||||||
crate::log_msg(&format!("Failed to send chat file: {e}"));
|
crate::log_msg(&format!("Failed to send chat file: {e}"));
|
||||||
}
|
}
|
||||||
@@ -3753,15 +3784,30 @@ mod tests {
|
|||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn auto_fetch_only_for_roster_images_not_already_inflight() {
|
fn auto_fetch_only_for_roster_images_not_already_inflight() {
|
||||||
// The happy path: a roster author's brand-new image attachment.
|
const OK_SIZE: u64 = 1024;
|
||||||
assert!(should_auto_fetch(true, true, false));
|
// The happy path: a roster author's brand-new, small-enough image.
|
||||||
|
assert!(should_auto_fetch(true, true, false, OK_SIZE));
|
||||||
// A non-image (generic file) never auto-fetches — it waits for "Save".
|
// A non-image (generic file) never auto-fetches — it waits for "Save".
|
||||||
assert!(!should_auto_fetch(false, true, false));
|
assert!(!should_auto_fetch(false, true, false, OK_SIZE));
|
||||||
// A non-roster author (e.g. a sock puppet that never announced) is rejected,
|
// A non-roster author (e.g. a sock puppet that never announced) is rejected,
|
||||||
// closing the F-02 unbounded-task vector.
|
// closing the F-02 unbounded-task vector.
|
||||||
assert!(!should_auto_fetch(true, false, false));
|
assert!(!should_auto_fetch(true, false, false, OK_SIZE));
|
||||||
// An identical (author,id) already being fetched is deduped.
|
// An identical (author,id) already being fetched is deduped.
|
||||||
assert!(!should_auto_fetch(true, true, true));
|
assert!(!should_auto_fetch(true, true, true, OK_SIZE));
|
||||||
|
// The declared-size gate (Phase 3B): at the cap auto-fetches, the first
|
||||||
|
// byte over requires a click.
|
||||||
|
assert!(should_auto_fetch(
|
||||||
|
true,
|
||||||
|
true,
|
||||||
|
false,
|
||||||
|
crate::files::MAX_AUTO_IMAGE_BYTES
|
||||||
|
));
|
||||||
|
assert!(!should_auto_fetch(
|
||||||
|
true,
|
||||||
|
true,
|
||||||
|
false,
|
||||||
|
crate::files::MAX_AUTO_IMAGE_BYTES + 1
|
||||||
|
));
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
|
|||||||
+341
-10
@@ -22,6 +22,22 @@ pub const MAX_ATTACHMENT_BYTES: u64 = 25 * 1024 * 1024;
|
|||||||
/// the byte cap. Applied via `image::Limits` when validating/decoding.
|
/// the byte cap. Applied via `image::Limits` when validating/decoding.
|
||||||
pub const MAX_IMAGE_PX: u32 = 4096;
|
pub const MAX_IMAGE_PX: u32 = 4096;
|
||||||
|
|
||||||
|
/// Max total decoded pixels, applied on top of the per-side [`MAX_IMAGE_PX`]
|
||||||
|
/// limit. The per-side cap alone still admits a 4096×4096 ≈ 16.8 MP bitmap
|
||||||
|
/// (~64 MiB transient RGBA); this bounds the worst-case decode allocation while
|
||||||
|
/// still clearing common 12 MP phone photos (4032×3024 ≈ 12.2 MP).
|
||||||
|
pub const MAX_IMAGE_TOTAL_PIXELS: u64 = 14_000_000;
|
||||||
|
|
||||||
|
/// Max pixels per side of the downscaled inline preview handed to the renderer.
|
||||||
|
/// Original bytes are kept only for Save; the chat column never needs more than
|
||||||
|
/// this (it displays at ~260 px, and the lightbox at window size).
|
||||||
|
pub const IMAGE_PREVIEW_MAX_SIDE: u32 = 1600;
|
||||||
|
|
||||||
|
/// Largest declared size an image attachment may auto-fetch at. Anything larger
|
||||||
|
/// (or any skipped/evicted image) renders a "Load image" button instead; a
|
||||||
|
/// manual click may use the full [`MAX_ATTACHMENT_BYTES`] cap.
|
||||||
|
pub const MAX_AUTO_IMAGE_BYTES: u64 = 4 * 1024 * 1024;
|
||||||
|
|
||||||
/// Longest filename we keep and display. Keeps the gossip descriptor compact and
|
/// Longest filename we keep and display. Keeps the gossip descriptor compact and
|
||||||
/// the UI tidy; the real bytes are unaffected.
|
/// the UI tidy; the real bytes are unaffected.
|
||||||
pub const MAX_FILENAME_LEN: usize = 96;
|
pub const MAX_FILENAME_LEN: usize = 96;
|
||||||
@@ -75,8 +91,14 @@ pub fn sanitize_filename(raw: &str) -> String {
|
|||||||
.unwrap_or("")
|
.unwrap_or("")
|
||||||
.trim();
|
.trim();
|
||||||
|
|
||||||
// Drop control chars; turn other whitespace into single spaces later.
|
// Drop control chars and the same bidi/zero-width spoofing format chars
|
||||||
let cleaned: String = base.chars().filter(|c| !c.is_control()).collect();
|
// stripped from display names (a U+202E override can visually reverse an
|
||||||
|
// extension, e.g. "photo\u{202E}gnp.exe" renders as "photoexe.png").
|
||||||
|
// Ordinary non-ASCII filenames pass through untouched.
|
||||||
|
let cleaned: String = base
|
||||||
|
.chars()
|
||||||
|
.filter(|c| !c.is_control() && !crate::sanitize::is_spoofing_format_char(*c))
|
||||||
|
.collect();
|
||||||
let collapsed = cleaned.split_whitespace().collect::<Vec<_>>().join(" ");
|
let collapsed = cleaned.split_whitespace().collect::<Vec<_>>().join(" ");
|
||||||
let collapsed = collapsed.trim_matches('.').trim();
|
let collapsed = collapsed.trim_matches('.').trim();
|
||||||
|
|
||||||
@@ -164,20 +186,185 @@ pub fn classify(bytes: &[u8]) -> AttachmentKind {
|
|||||||
/// regardless of [`MAX_ATTACHMENT_BYTES`]. Only PNG/JPEG are buildable in our
|
/// regardless of [`MAX_ATTACHMENT_BYTES`]. Only PNG/JPEG are buildable in our
|
||||||
/// `image` feature set; anything else returns `None` and the caller shows a chip.
|
/// `image` feature set; anything else returns `None` and the caller shows a chip.
|
||||||
pub fn validate_image_bytes(bytes: &[u8]) -> Option<(u32, u32)> {
|
pub fn validate_image_bytes(bytes: &[u8]) -> Option<(u32, u32)> {
|
||||||
let mut limits = image::Limits::default();
|
let img = decode_image_bounded(bytes)?;
|
||||||
limits.max_image_width = Some(MAX_IMAGE_PX);
|
Some((img.width(), img.height()))
|
||||||
limits.max_image_height = Some(MAX_IMAGE_PX);
|
}
|
||||||
|
|
||||||
|
/// Shared bounded decode: header-check the dimensions (per-side AND total-pixel
|
||||||
|
/// limits) BEFORE decoding, then decode under `image::Limits` as defense in
|
||||||
|
/// depth. The precheck reads only the container header, so an over-limit bomb is
|
||||||
|
/// rejected without paying its decode cost.
|
||||||
|
fn decode_image_bounded(bytes: &[u8]) -> Option<image::DynamicImage> {
|
||||||
let reader = image::ImageReader::new(std::io::Cursor::new(bytes))
|
let reader = image::ImageReader::new(std::io::Cursor::new(bytes))
|
||||||
.with_guessed_format()
|
.with_guessed_format()
|
||||||
.ok()?;
|
.ok()?;
|
||||||
let mut reader = reader;
|
let (w, h) = reader.into_dimensions().ok()?;
|
||||||
reader.limits(limits);
|
|
||||||
let img = reader.decode().ok()?;
|
|
||||||
let (w, h) = (img.width(), img.height());
|
|
||||||
if w == 0 || h == 0 || w > MAX_IMAGE_PX || h > MAX_IMAGE_PX {
|
if w == 0 || h == 0 || w > MAX_IMAGE_PX || h > MAX_IMAGE_PX {
|
||||||
return None;
|
return None;
|
||||||
}
|
}
|
||||||
Some((w, h))
|
if u64::from(w) * u64::from(h) > MAX_IMAGE_TOTAL_PIXELS {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
let mut limits = image::Limits::default();
|
||||||
|
limits.max_image_width = Some(MAX_IMAGE_PX);
|
||||||
|
limits.max_image_height = Some(MAX_IMAGE_PX);
|
||||||
|
let mut reader = image::ImageReader::new(std::io::Cursor::new(bytes))
|
||||||
|
.with_guessed_format()
|
||||||
|
.ok()?;
|
||||||
|
reader.limits(limits);
|
||||||
|
let img = reader.decode().ok()?;
|
||||||
|
// Decoded size must match the header the precheck approved.
|
||||||
|
if img.width() != w || img.height() != h {
|
||||||
|
return None;
|
||||||
|
}
|
||||||
|
Some(img)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A decoded, display-ready inline preview: RGBA pixels downscaled so neither
|
||||||
|
/// side exceeds [`IMAGE_PREVIEW_MAX_SIDE`]. `rgba.len() == width * height * 4`,
|
||||||
|
/// which is also the preview's decoded-budget weight in the attachment cache.
|
||||||
|
pub struct ImagePreview {
|
||||||
|
pub width: u32,
|
||||||
|
pub height: u32,
|
||||||
|
pub rgba: Vec<u8>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Decode image bytes under the same limits as [`validate_image_bytes`] and
|
||||||
|
/// build the downscaled inline preview. The full-resolution bitmap exists only
|
||||||
|
/// transiently here; the renderer is never handed more than
|
||||||
|
/// [`IMAGE_PREVIEW_MAX_SIDE`]² pixels. Returns `None` for anything that fails
|
||||||
|
/// validation (caller falls back to a chip / failure row).
|
||||||
|
pub fn decode_preview(bytes: &[u8]) -> Option<ImagePreview> {
|
||||||
|
let img = decode_image_bounded(bytes)?;
|
||||||
|
let img = if img.width() > IMAGE_PREVIEW_MAX_SIDE || img.height() > IMAGE_PREVIEW_MAX_SIDE {
|
||||||
|
// `thumbnail` preserves aspect ratio within the bounding box.
|
||||||
|
img.thumbnail(IMAGE_PREVIEW_MAX_SIDE, IMAGE_PREVIEW_MAX_SIDE)
|
||||||
|
} else {
|
||||||
|
img
|
||||||
|
};
|
||||||
|
let rgba = img.into_rgba8();
|
||||||
|
let (width, height) = (rgba.width(), rgba.height());
|
||||||
|
Some(ImagePreview {
|
||||||
|
width,
|
||||||
|
height,
|
||||||
|
rgba: rgba.into_raw(),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Estimated decoded RGBA cost of a preview, the weight counted against the
|
||||||
|
/// attachment cache's decoded-byte budget (`width * height * 4`).
|
||||||
|
pub fn preview_rgba_cost(width: u32, height: u32) -> usize {
|
||||||
|
(width as usize)
|
||||||
|
.saturating_mul(height as usize)
|
||||||
|
.saturating_mul(4)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Read at most [`MAX_ATTACHMENT_BYTES`] bytes from `r`. Returns `Ok(None)` if
|
||||||
|
/// the source holds even one byte more (detected by reading cap + 1), so a huge
|
||||||
|
/// or unbounded source is never fully buffered. Pure over `Read` for tests; the
|
||||||
|
/// picker wraps it via [`read_file_capped`].
|
||||||
|
pub fn read_capped<R: std::io::Read>(r: R) -> std::io::Result<Option<Vec<u8>>> {
|
||||||
|
use std::io::Read as _;
|
||||||
|
let mut buf = Vec::new();
|
||||||
|
let mut limited = r.take(MAX_ATTACHMENT_BYTES + 1);
|
||||||
|
limited.read_to_end(&mut buf)?;
|
||||||
|
if buf.len() as u64 > MAX_ATTACHMENT_BYTES {
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
Ok(Some(buf))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Read a picked file, bounded by [`MAX_ATTACHMENT_BYTES`]. Checks metadata
|
||||||
|
/// first to reject an obviously-oversized file without opening it, but keeps the
|
||||||
|
/// bounded read regardless — metadata can race (the file can grow after the
|
||||||
|
/// check) or be unavailable through a portal. `Ok(None)` = over the cap.
|
||||||
|
pub fn read_file_capped(path: &std::path::Path) -> std::io::Result<Option<Vec<u8>>> {
|
||||||
|
if let Ok(meta) = std::fs::metadata(path)
|
||||||
|
&& meta.len() > MAX_ATTACHMENT_BYTES
|
||||||
|
{
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
read_capped(std::fs::File::open(path)?)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Cap on how many blobs the session serve store retains at once (sent chat
|
||||||
|
/// attachments plus the current/next broadcast music tracks).
|
||||||
|
pub const SERVED_FILES_MAX_ENTRIES: usize = 16;
|
||||||
|
|
||||||
|
/// Byte budget for the serve store. Without it, a sender's own session could
|
||||||
|
/// grow unbounded at up to [`MAX_ATTACHMENT_BYTES`] per send (Phase 3C).
|
||||||
|
pub const SERVED_FILES_MAX_BYTES: usize = 128 * 1024 * 1024;
|
||||||
|
|
||||||
|
/// Count- and byte-budgeted FIFO store of blobs we serve to room members over
|
||||||
|
/// the file plane. Evicting an id makes a later request for it read as an empty
|
||||||
|
/// body — the existing "sender no longer has the file" response — never stale
|
||||||
|
/// or aliased bytes. Pure (no locks/IO) so budgets are unit-testable; the
|
||||||
|
/// transport wraps it in its own mutex.
|
||||||
|
#[derive(Debug, Default)]
|
||||||
|
pub struct ServeStore {
|
||||||
|
entries: std::collections::HashMap<AttachmentId, std::sync::Arc<Vec<u8>>>,
|
||||||
|
/// Present ids in insertion order; the front is the eviction candidate.
|
||||||
|
order: std::collections::VecDeque<AttachmentId>,
|
||||||
|
total_bytes: usize,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ServeStore {
|
||||||
|
/// Insert or replace a blob, evicting oldest entries until the count and
|
||||||
|
/// byte budgets fit. Replacement keeps the id's age and subtracts the old
|
||||||
|
/// bytes before the new ones are counted. Returns `false` for a blob that
|
||||||
|
/// alone exceeds the byte budget (not stored; an existing entry under the
|
||||||
|
/// id is dropped rather than left stale).
|
||||||
|
pub fn insert(&mut self, id: AttachmentId, bytes: std::sync::Arc<Vec<u8>>) -> bool {
|
||||||
|
if let Some(old) = self.entries.get(&id) {
|
||||||
|
self.total_bytes -= old.len();
|
||||||
|
}
|
||||||
|
if bytes.len() > SERVED_FILES_MAX_BYTES {
|
||||||
|
if self.entries.remove(&id).is_some() {
|
||||||
|
self.order.retain(|k| k != &id);
|
||||||
|
}
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
let replacing = self.entries.contains_key(&id);
|
||||||
|
loop {
|
||||||
|
let count_full = !replacing && self.entries.len() >= SERVED_FILES_MAX_ENTRIES;
|
||||||
|
let bytes_full = self.total_bytes + bytes.len() > SERVED_FILES_MAX_BYTES;
|
||||||
|
if !count_full && !bytes_full {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
let Some(victim) = self.order.iter().find(|k| **k != id).copied() else {
|
||||||
|
break;
|
||||||
|
};
|
||||||
|
self.remove(&victim);
|
||||||
|
}
|
||||||
|
if !replacing {
|
||||||
|
self.order.push_back(id);
|
||||||
|
}
|
||||||
|
self.total_bytes += bytes.len();
|
||||||
|
self.entries.insert(id, bytes);
|
||||||
|
true
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn get(&self, id: &AttachmentId) -> Option<std::sync::Arc<Vec<u8>>> {
|
||||||
|
self.entries.get(id).cloned()
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn remove(&mut self, id: &AttachmentId) {
|
||||||
|
if let Some(old) = self.entries.remove(id) {
|
||||||
|
self.total_bytes -= old.len();
|
||||||
|
self.order.retain(|k| k != id);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub fn clear(&mut self) {
|
||||||
|
self.entries.clear();
|
||||||
|
self.order.clear();
|
||||||
|
self.total_bytes = 0;
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
fn len(&self) -> usize {
|
||||||
|
self.entries.len()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Parse a file-plane request: it must be exactly one [`AttachmentId`] (32
|
/// Parse a file-plane request: it must be exactly one [`AttachmentId`] (32
|
||||||
@@ -352,6 +539,87 @@ mod tests {
|
|||||||
assert_eq!(validate_image_bytes(&buf.into_inner()), Some((4, 3)));
|
assert_eq!(validate_image_bytes(&buf.into_inner()), Some((4, 3)));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn sanitize_strips_bidi_and_zero_width_spoofing_chars() {
|
||||||
|
// U+202E would visually reverse the tail, disguising the extension.
|
||||||
|
assert_eq!(sanitize_filename("photo\u{202E}gnp.exe"), "photognp.exe");
|
||||||
|
assert_eq!(sanitize_filename("a\u{200B}b\u{FEFF}.txt"), "ab.txt");
|
||||||
|
// Ordinary Unicode filenames pass through.
|
||||||
|
assert_eq!(sanitize_filename("família_fotos.png"), "família_fotos.png");
|
||||||
|
assert_eq!(sanitize_filename("日本語.pdf"), "日本語.pdf");
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Encode a solid PNG of the given dimensions for limit tests.
|
||||||
|
fn png_bytes(w: u32, h: u32) -> Vec<u8> {
|
||||||
|
let img = image::RgbImage::from_pixel(w, h, image::Rgb([10, 20, 30]));
|
||||||
|
let mut buf = std::io::Cursor::new(Vec::new());
|
||||||
|
image::DynamicImage::ImageRgb8(img)
|
||||||
|
.write_to(&mut buf, image::ImageFormat::Png)
|
||||||
|
.unwrap();
|
||||||
|
buf.into_inner()
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn validate_image_rejects_excessive_total_pixels() {
|
||||||
|
// Both sides within MAX_IMAGE_PX, but 4096 * 4096 > MAX_IMAGE_TOTAL_PIXELS.
|
||||||
|
assert!(u64::from(MAX_IMAGE_PX) * u64::from(MAX_IMAGE_PX) > MAX_IMAGE_TOTAL_PIXELS);
|
||||||
|
assert_eq!(validate_image_bytes(&png_bytes(4096, 4096)), None);
|
||||||
|
// A 12 MP phone-photo shape passes both limits.
|
||||||
|
assert_eq!(
|
||||||
|
validate_image_bytes(&png_bytes(4032, 3024)),
|
||||||
|
Some((4032, 3024))
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn preview_downscales_to_max_side_preserving_aspect() {
|
||||||
|
// Wide: 3200x400 → 1600x200.
|
||||||
|
let p = decode_preview(&png_bytes(3200, 400)).unwrap();
|
||||||
|
assert_eq!((p.width, p.height), (1600, 200));
|
||||||
|
assert_eq!(p.rgba.len(), preview_rgba_cost(1600, 200));
|
||||||
|
// Tall: 400x3200 → 200x1600.
|
||||||
|
let p = decode_preview(&png_bytes(400, 3200)).unwrap();
|
||||||
|
assert_eq!((p.width, p.height), (200, 1600));
|
||||||
|
// Square over the side cap: 2000x2000 → 1600x1600.
|
||||||
|
let p = decode_preview(&png_bytes(2000, 2000)).unwrap();
|
||||||
|
assert_eq!((p.width, p.height), (1600, 1600));
|
||||||
|
// At/under the cap is untouched.
|
||||||
|
let p = decode_preview(&png_bytes(1600, 900)).unwrap();
|
||||||
|
assert_eq!((p.width, p.height), (1600, 900));
|
||||||
|
let p = decode_preview(&png_bytes(4, 3)).unwrap();
|
||||||
|
assert_eq!((p.width, p.height), (4, 3));
|
||||||
|
assert_eq!(p.rgba.len(), preview_rgba_cost(4, 3));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn preview_rejects_what_validation_rejects() {
|
||||||
|
assert!(decode_preview(b"not an image").is_none());
|
||||||
|
assert!(decode_preview(&png_bytes(4096, 4096)).is_none());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn read_capped_stops_at_cap_plus_one() {
|
||||||
|
// Under the cap: full read.
|
||||||
|
let small = vec![7u8; 1024];
|
||||||
|
assert_eq!(
|
||||||
|
read_capped(std::io::Cursor::new(&small))
|
||||||
|
.unwrap()
|
||||||
|
.as_deref(),
|
||||||
|
Some(&small[..])
|
||||||
|
);
|
||||||
|
// Exactly at the cap: accepted. `repeat` is endless, `take` proves the
|
||||||
|
// reader is bounded rather than draining the source.
|
||||||
|
let at_cap = std::io::Read::take(std::io::repeat(1), MAX_ATTACHMENT_BYTES);
|
||||||
|
let got = read_capped(at_cap).unwrap().unwrap();
|
||||||
|
assert_eq!(got.len() as u64, MAX_ATTACHMENT_BYTES);
|
||||||
|
// One byte over: rejected, and only cap + 1 bytes were ever buffered
|
||||||
|
// (an unbounded source returns instead of allocating forever).
|
||||||
|
let over = std::io::Read::take(std::io::repeat(1), MAX_ATTACHMENT_BYTES + 1);
|
||||||
|
assert_eq!(read_capped(over).unwrap(), None);
|
||||||
|
let endless = std::io::repeat(1);
|
||||||
|
assert_eq!(read_capped(endless).unwrap(), None);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn human_size_units() {
|
fn human_size_units() {
|
||||||
assert_eq!(human_size(40), "40 B");
|
assert_eq!(human_size(40), "40 B");
|
||||||
@@ -359,6 +627,69 @@ mod tests {
|
|||||||
assert_eq!(human_size(3 * 1024 * 1024 + 300 * 1024), "3.3 MB");
|
assert_eq!(human_size(3 * 1024 * 1024 + 300 * 1024), "3.3 MB");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn serve_store_count_and_byte_eviction_fifo() {
|
||||||
|
use std::sync::Arc;
|
||||||
|
let mut s = ServeStore::default();
|
||||||
|
let blob = |n: u8, len: usize| ([n; 32], Arc::new(vec![n; len]));
|
||||||
|
// Count cap: entry 0 is evicted when the 17th arrives.
|
||||||
|
for n in 0..=SERVED_FILES_MAX_ENTRIES as u8 {
|
||||||
|
let (id, b) = blob(n, 8);
|
||||||
|
assert!(s.insert(id, b));
|
||||||
|
}
|
||||||
|
assert_eq!(s.len(), SERVED_FILES_MAX_ENTRIES);
|
||||||
|
assert!(s.get(&[0u8; 32]).is_none(), "oldest evicted by count");
|
||||||
|
assert!(s.get(&[1u8; 32]).is_some());
|
||||||
|
// Byte budget: two ~half-budget blobs evict everything older.
|
||||||
|
let half = SERVED_FILES_MAX_BYTES / 2;
|
||||||
|
let (a, ab) = blob(100, half);
|
||||||
|
let (b, bb) = blob(101, half);
|
||||||
|
assert!(s.insert(a, ab));
|
||||||
|
assert!(s.insert(b, bb));
|
||||||
|
assert!(s.get(&a).is_some());
|
||||||
|
assert!(s.get(&b).is_some());
|
||||||
|
assert!(s.get(&[1u8; 32]).is_none(), "evicted for byte budget");
|
||||||
|
// A third half-budget blob evicts `a` (oldest), keeps `b`.
|
||||||
|
let (c, cb) = blob(102, half);
|
||||||
|
assert!(s.insert(c, cb));
|
||||||
|
assert!(s.get(&a).is_none());
|
||||||
|
assert!(s.get(&b).is_some());
|
||||||
|
assert!(s.get(&c).is_some());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn serve_store_replacement_accounting_and_remove_clear() {
|
||||||
|
use std::sync::Arc;
|
||||||
|
let mut s = ServeStore::default();
|
||||||
|
let id = [9u8; 32];
|
||||||
|
assert!(s.insert(id, Arc::new(vec![1; SERVED_FILES_MAX_BYTES - 10])));
|
||||||
|
// Replacing the near-budget blob must subtract its old bytes first —
|
||||||
|
// otherwise this same-id replacement would evict itself.
|
||||||
|
assert!(s.insert(id, Arc::new(vec![2; SERVED_FILES_MAX_BYTES - 5])));
|
||||||
|
assert_eq!(s.get(&id).unwrap()[0], 2);
|
||||||
|
assert_eq!(s.len(), 1);
|
||||||
|
s.remove(&id);
|
||||||
|
assert!(s.get(&id).is_none());
|
||||||
|
// Removed bytes were released: the budget admits a full-size blob again.
|
||||||
|
assert!(s.insert(id, Arc::new(vec![3; SERVED_FILES_MAX_BYTES])));
|
||||||
|
s.clear();
|
||||||
|
assert_eq!(s.len(), 0);
|
||||||
|
assert!(s.insert(id, Arc::new(vec![4; SERVED_FILES_MAX_BYTES])));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn serve_store_rejects_individually_overweight_blob() {
|
||||||
|
use std::sync::Arc;
|
||||||
|
let mut s = ServeStore::default();
|
||||||
|
let id = [7u8; 32];
|
||||||
|
assert!(s.insert(id, Arc::new(vec![1; 8])));
|
||||||
|
assert!(!s.insert(id, Arc::new(vec![2; SERVED_FILES_MAX_BYTES + 1])));
|
||||||
|
// The stale small blob is gone too — a fetch reads "no longer has it",
|
||||||
|
// never old bytes under a replaced id.
|
||||||
|
assert!(s.get(&id).is_none());
|
||||||
|
assert_eq!(s.len(), 0);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn attachment_descriptor_round_trips_json() {
|
fn attachment_descriptor_round_trips_json() {
|
||||||
let a = ChatAttachment {
|
let a = ChatAttachment {
|
||||||
|
|||||||
@@ -69,7 +69,9 @@ struct Shared {
|
|||||||
/// the random attachment id. Populated when we send a chat file; read by the
|
/// the random attachment id. Populated when we send a chat file; read by the
|
||||||
/// file protocol handler to answer a member's fetch. Cleared on leave. Each
|
/// file protocol handler to answer a member's fetch. Cleared on leave. Each
|
||||||
/// blob is already byte-capped at send time.
|
/// blob is already byte-capped at send time.
|
||||||
served_files: StdMutex<HashMap<crate::files::AttachmentId, Arc<Vec<u8>>>>,
|
/// Blobs we serve to room members, bounded by count and byte budgets
|
||||||
|
/// (Phase 3C) — an evicted id reads as "sender no longer has the file".
|
||||||
|
served_files: StdMutex<crate::files::ServeStore>,
|
||||||
incoming_tx: mpsc::Sender<(EndpointId, Bytes)>,
|
incoming_tx: mpsc::Sender<(EndpointId, Bytes)>,
|
||||||
/// Best-effort link-state notifications for the UI (connecting / connected).
|
/// Best-effort link-state notifications for the UI (connecting / connected).
|
||||||
conn_events_tx: mpsc::Sender<ConnEvent>,
|
conn_events_tx: mpsc::Sender<ConnEvent>,
|
||||||
@@ -518,7 +520,7 @@ impl iroh::protocol::ProtocolHandler for FileRouter {
|
|||||||
let Some(id) = crate::files::parse_request(&req) else {
|
let Some(id) = crate::files::parse_request(&req) else {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
};
|
};
|
||||||
let blob = shared.served_files.lock().unwrap().get(&id).cloned();
|
let blob = shared.served_files.lock().unwrap().get(&id);
|
||||||
if let Some(blob) = blob {
|
if let Some(blob) = blob {
|
||||||
let _ = send.write_all(&blob).await;
|
let _ = send.write_all(&blob).await;
|
||||||
}
|
}
|
||||||
@@ -563,7 +565,7 @@ impl IrohTransport {
|
|||||||
peers: tokio::sync::Mutex::new(HashMap::new()),
|
peers: tokio::sync::Mutex::new(HashMap::new()),
|
||||||
live_conns: StdMutex::new(HashMap::new()),
|
live_conns: StdMutex::new(HashMap::new()),
|
||||||
admitted_audio: StdMutex::new(HashSet::new()),
|
admitted_audio: StdMutex::new(HashSet::new()),
|
||||||
served_files: StdMutex::new(HashMap::new()),
|
served_files: StdMutex::new(crate::files::ServeStore::default()),
|
||||||
incoming_tx,
|
incoming_tx,
|
||||||
conn_events_tx,
|
conn_events_tx,
|
||||||
});
|
});
|
||||||
@@ -634,7 +636,9 @@ impl IrohTransport {
|
|||||||
/// session (served by the [`FileRouter`] handler). Called by core when we
|
/// session (served by the [`FileRouter`] handler). Called by core when we
|
||||||
/// send a chat file. The blob is cleared on leave.
|
/// send a chat file. The blob is cleared on leave.
|
||||||
pub fn serve_attachment(&self, id: AttachmentId, bytes: Arc<Vec<u8>>) {
|
pub fn serve_attachment(&self, id: AttachmentId, bytes: Arc<Vec<u8>>) {
|
||||||
self.shared.served_files.lock().unwrap().insert(id, bytes);
|
if !self.shared.served_files.lock().unwrap().insert(id, bytes) {
|
||||||
|
crate::log_msg("Transport: refused to serve an over-budget blob");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Drop a previously-served blob (e.g. a music track no longer current-or-next).
|
/// Drop a previously-served blob (e.g. a music track no longer current-or-next).
|
||||||
@@ -676,6 +680,8 @@ impl IrohTransport {
|
|||||||
send.finish()
|
send.finish()
|
||||||
.map_err(|e| NetError::Other(format!("file fetch: request finish failed: {e}")))?;
|
.map_err(|e| NetError::Other(format!("file fetch: request finish failed: {e}")))?;
|
||||||
|
|
||||||
|
// `read_to_end(size)` errors if the stream exceeds `size`, rejecting an
|
||||||
|
// overlong transfer; the exact-length check below rejects a short one.
|
||||||
let read = recv.read_to_end(size as usize);
|
let read = recv.read_to_end(size as usize);
|
||||||
let bytes = tokio::time::timeout(FILE_FETCH_TIMEOUT, read)
|
let bytes = tokio::time::timeout(FILE_FETCH_TIMEOUT, read)
|
||||||
.await
|
.await
|
||||||
@@ -686,6 +692,14 @@ impl IrohTransport {
|
|||||||
"file fetch: sender no longer has the file".to_string(),
|
"file fetch: sender no longer has the file".to_string(),
|
||||||
));
|
));
|
||||||
}
|
}
|
||||||
|
// Exact transfer required (Phase 3C): a truncated body must not be
|
||||||
|
// cached/saved/decoded as if it were the declared attachment.
|
||||||
|
if bytes.len() as u64 != size {
|
||||||
|
return Err(NetError::Other(format!(
|
||||||
|
"file fetch: incomplete transfer ({} of {size} bytes)",
|
||||||
|
bytes.len()
|
||||||
|
)));
|
||||||
|
}
|
||||||
Ok(bytes)
|
Ok(bytes)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+1
-1
@@ -18,7 +18,7 @@ pub const NAME_MAX_CHARS: usize = 48;
|
|||||||
/// zero-width / BOM characters (invisible, can hide or fake content). Listed
|
/// zero-width / BOM characters (invisible, can hide or fake content). Listed
|
||||||
/// explicitly so the sanitizer stays dependency-free (std exposes no category
|
/// explicitly so the sanitizer stays dependency-free (std exposes no category
|
||||||
/// query). Stripped outright rather than replaced.
|
/// query). Stripped outright rather than replaced.
|
||||||
fn is_spoofing_format_char(c: char) -> bool {
|
pub(crate) fn is_spoofing_format_char(c: char) -> bool {
|
||||||
matches!(c,
|
matches!(c,
|
||||||
'\u{200B}'..='\u{200F}' // zero-width space, ZWNJ, ZWJ, LRM, RLM
|
'\u{200B}'..='\u{200F}' // zero-width space, ZWNJ, ZWJ, LRM, RLM
|
||||||
| '\u{202A}'..='\u{202E}' // LRE, RLE, PDF, LRO, RLO (bidi overrides)
|
| '\u{202A}'..='\u{202E}' // LRE, RLE, PDF, LRO, RLO (bidi overrides)
|
||||||
|
|||||||
@@ -570,3 +570,94 @@ async fn dialer_connects_and_reconnects_without_an_address_lookup() {
|
|||||||
"audio should resume after reconnecting via the retained address; got {received} frames"
|
"audio should resume after reconnecting via the retained address; got {received} frames"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A node that also serves the file plane (`FILES_ALPN`), mirroring how core
|
||||||
|
/// registers the `FileRouter` for a session.
|
||||||
|
async fn spawn_file_server() -> Node {
|
||||||
|
let lookup = MemoryLookup::new();
|
||||||
|
let endpoint = Endpoint::builder(presets::Minimal)
|
||||||
|
.secret_key(iroh::SecretKey::generate())
|
||||||
|
.relay_mode(RelayMode::Disabled)
|
||||||
|
.address_lookup(lookup.clone())
|
||||||
|
.bind()
|
||||||
|
.await
|
||||||
|
.expect("bind endpoint");
|
||||||
|
let transport = Arc::new(IrohTransport::new(endpoint.clone()));
|
||||||
|
let audio_router = AudioRouter::new();
|
||||||
|
audio_router.bind(&transport);
|
||||||
|
let file_router = peerspeak::network::iroh_impl::FileRouter::new();
|
||||||
|
file_router.bind(&transport);
|
||||||
|
let router = Router::builder(endpoint.clone())
|
||||||
|
.accept(AUDIO_ALPN, audio_router)
|
||||||
|
.accept(peerspeak::protocol::FILES_ALPN, file_router)
|
||||||
|
.spawn();
|
||||||
|
Node {
|
||||||
|
endpoint,
|
||||||
|
transport,
|
||||||
|
_router: router,
|
||||||
|
lookup,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Phase 3C: a file fetch must deliver EXACTLY the declared size — short,
|
||||||
|
/// overlong, and unknown-id transfers are all rejected with local errors, and
|
||||||
|
/// an exact transfer round-trips byte-identically.
|
||||||
|
#[tokio::test]
|
||||||
|
async fn file_plane_requires_exact_declared_size() {
|
||||||
|
let fetcher = spawn_node().await;
|
||||||
|
let server = spawn_file_server().await;
|
||||||
|
fetcher.lookup.add_endpoint_info(server.endpoint.addr());
|
||||||
|
server.lookup.add_endpoint_info(fetcher.endpoint.addr());
|
||||||
|
|
||||||
|
let server_id = server.endpoint.id();
|
||||||
|
// Member gating: the server only serves current room members.
|
||||||
|
server.transport.admit_audio_sender(fetcher.endpoint.id());
|
||||||
|
|
||||||
|
let blob = vec![42u8; 1000];
|
||||||
|
let id = [7u8; 32];
|
||||||
|
server
|
||||||
|
.transport
|
||||||
|
.serve_attachment(id, Arc::new(blob.clone()));
|
||||||
|
|
||||||
|
// Exact declared size: byte-identical round trip.
|
||||||
|
let got = fetcher
|
||||||
|
.transport
|
||||||
|
.fetch_blob(server_id, id, 1000)
|
||||||
|
.await
|
||||||
|
.expect("exact-size fetch succeeds");
|
||||||
|
assert_eq!(got, blob);
|
||||||
|
|
||||||
|
// Declared larger than served (short transfer): rejected, not cached as-is.
|
||||||
|
let err = fetcher
|
||||||
|
.transport
|
||||||
|
.fetch_blob(server_id, id, 2000)
|
||||||
|
.await
|
||||||
|
.expect_err("short transfer must fail");
|
||||||
|
assert!(
|
||||||
|
err.to_string().contains("incomplete transfer"),
|
||||||
|
"unexpected error: {err}"
|
||||||
|
);
|
||||||
|
|
||||||
|
// Declared smaller than served (overlong transfer): the bounded read
|
||||||
|
// rejects the stream rather than truncating it into a "valid" result.
|
||||||
|
let err = fetcher
|
||||||
|
.transport
|
||||||
|
.fetch_blob(server_id, id, 500)
|
||||||
|
.await
|
||||||
|
.expect_err("overlong transfer must fail");
|
||||||
|
assert!(
|
||||||
|
err.to_string().contains("read failed"),
|
||||||
|
"unexpected error: {err}"
|
||||||
|
);
|
||||||
|
|
||||||
|
// Unknown id: the empty body reads as the sender no longer having it.
|
||||||
|
let err = fetcher
|
||||||
|
.transport
|
||||||
|
.fetch_blob(server_id, [9u8; 32], 1000)
|
||||||
|
.await
|
||||||
|
.expect_err("unknown id must fail");
|
||||||
|
assert!(
|
||||||
|
err.to_string().contains("no longer has the file"),
|
||||||
|
"unexpected error: {err}"
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user