use super::{ClientContext, PeerRegistry, UploadJob};
use crate::types::{UploadInfo, UploadStatus};
use std::path::PathBuf;
const MAX_UPLOAD_EVENTS: usize = 256;
#[derive(Debug, Clone)]
pub struct UploadOffer {
pub requester_key: String,
pub token: u32,
pub virtual_path: String,
pub size: u64,
}
#[derive(Debug, Clone)]
pub struct QueuedUpload {
pub requester_key: String,
pub downloader: String,
pub virtual_path: String,
pub real_path: PathBuf,
pub size: u64,
pub privileged: bool,
pub seq: u64,
}
const fn rank(job: &QueuedUpload) -> (bool, u64) {
(!job.privileged, job.seq)
}
#[must_use]
pub fn serve_order(queue: &[QueuedUpload]) -> Vec<usize> {
let mut order: Vec<usize> = (0..queue.len()).collect();
order.sort_by_key(|&i| rank(&queue[i]));
order
}
#[must_use]
pub fn next_to_serve(queue: &[QueuedUpload]) -> Option<usize> {
queue
.iter()
.enumerate()
.min_by_key(|(_, job)| rank(job))
.map(|(index, _)| index)
}
fn as_place(zero_based: usize) -> u32 {
u32::try_from(zero_based + 1).unwrap_or(u32::MAX)
}
#[must_use]
pub fn place_in_queue(
queue: &[QueuedUpload],
downloader: &str,
filename: &str,
) -> Option<u32> {
serve_order(queue)
.into_iter()
.position(|i| {
queue[i].downloader == downloader
&& queue[i].virtual_path == filename
})
.map(as_place)
}
impl ClientContext {
pub fn enqueue_upload(
&mut self,
requester_key: &str,
downloader: &str,
filename: &str,
real_path: std::path::PathBuf,
size: u64,
) {
self.upload_seq += 1;
self.upload_queue.push(QueuedUpload {
requester_key: requester_key.to_string(),
downloader: downloader.to_string(),
virtual_path: filename.to_string(),
real_path,
size,
privileged: self.privileged_users.contains(downloader),
seq: self.upload_seq,
});
}
pub fn pump_uploads(
&mut self,
mut next_token: impl FnMut() -> u32,
) -> (Option<PeerRegistry>, Vec<UploadOffer>) {
let mut offers = Vec::new();
while self.uploads_in_flight() < self.upload_slots {
let Some(index) = next_to_serve(&self.upload_queue) else {
break;
};
let job = self.upload_queue.remove(index);
let token = next_token();
offers.push(UploadOffer {
requester_key: job.requester_key,
token,
virtual_path: job.virtual_path.clone(),
size: job.size,
});
self.uploads.insert(
token,
UploadJob {
downloader: job.downloader,
real_path: job.real_path,
virtual_path: job.virtual_path,
size: job.size,
},
);
}
for waiting in self.queued_uploads() {
let already = self.upload_events.iter().any(|seen| {
seen.username == waiting.username
&& seen.filename == waiting.filename
});
if !already {
self.upload_events.push(waiting);
}
}
if self.upload_events.len() > MAX_UPLOAD_EVENTS {
let excess = self.upload_events.len() - MAX_UPLOAD_EVENTS;
self.upload_events.drain(..excess);
}
(self.peer_registry.clone(), offers)
}
pub fn take_upload_events(&mut self) -> Vec<UploadInfo> {
std::mem::take(&mut self.upload_events)
}
fn uploads_in_flight(&self) -> usize {
self.uploads.len()
+ self
.active_uploads
.values()
.filter(|upload| {
matches!(upload.status, UploadStatus::InProgress)
})
.count()
}
pub(super) fn queued_uploads(&self) -> Vec<UploadInfo> {
serve_order(&self.upload_queue)
.into_iter()
.enumerate()
.map(|(place, index)| {
let job = &self.upload_queue[index];
UploadInfo {
username: job.downloader.clone(),
filename: job.virtual_path.clone(),
size: job.size,
bytes_sent: 0,
speed_bytes_per_sec: 0.0,
status: UploadStatus::Queued(as_place(place)),
}
})
.collect()
}
#[must_use]
pub fn place_in_queue(
&self,
downloader: &str,
filename: &str,
) -> Option<u32> {
place_in_queue(&self.upload_queue, downloader, filename)
}
pub fn set_privileged_users(&mut self, users: Vec<String>) {
self.privileged_users = users.into_iter().collect();
for job in &mut self.upload_queue {
job.privileged = self.privileged_users.contains(&job.downloader);
}
}
#[must_use]
pub fn is_privileged(&self, username: &str) -> bool {
self.privileged_users.contains(username)
}
pub fn release_upload_slots(&mut self, downloader: &str) -> bool {
let before = self.upload_queue.len() + self.uploads.len();
self.upload_queue.retain(|job| job.downloader != downloader);
self.uploads.retain(|_, job| job.downloader != downloader);
before != self.upload_queue.len() + self.uploads.len()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn queued(user: &str, privileged: bool, seq: u64) -> QueuedUpload {
QueuedUpload {
requester_key: user.to_string(),
downloader: user.to_string(),
virtual_path: format!("@@share\\{user}.mp3"),
real_path: PathBuf::from("/tmp/x.mp3"),
size: 1,
privileged,
seq,
}
}
#[test]
fn an_empty_queue_serves_nobody() {
assert_eq!(next_to_serve(&[]), None);
}
#[test]
fn equals_are_served_in_the_order_they_asked() {
let queue = [queued("amy", false, 1), queued("bob", false, 2)];
assert_eq!(next_to_serve(&queue), Some(0));
assert_eq!(serve_order(&queue), [0, 1]);
}
#[test]
fn a_privileged_peer_is_served_before_one_that_asked_first() {
let queue = [
queued("early_plain", false, 1),
queued("late_donor", true, 2),
];
assert_eq!(
next_to_serve(&queue),
Some(1),
"the donor jumps the plain user who was already waiting"
);
}
#[test]
fn privileged_peers_keep_first_come_order_among_themselves() {
let queue = [
queued("donor_late", true, 9),
queued("plain", false, 2),
queued("donor_early", true, 5),
];
let order: Vec<&str> = serve_order(&queue)
.into_iter()
.map(|i| queue[i].downloader.as_str())
.collect();
assert_eq!(order, ["donor_early", "donor_late", "plain"]);
}
#[test]
fn the_reported_place_counts_from_one() {
let queue = [queued("amy", false, 1), queued("bob", false, 2)];
assert_eq!(place_in_queue(&queue, "amy", "@@share\\amy.mp3"), Some(1));
assert_eq!(place_in_queue(&queue, "bob", "@@share\\bob.mp3"), Some(2));
}
#[test]
fn the_reported_place_reflects_being_overtaken() {
let mut queue = vec![queued("plain", false, 1)];
assert_eq!(
place_in_queue(&queue, "plain", "@@share\\plain.mp3"),
Some(1)
);
queue.push(queued("donor", true, 2));
assert_eq!(
place_in_queue(&queue, "plain", "@@share\\plain.mp3"),
Some(2),
"a peer that got overtaken must be told so, not left on 1"
);
}
#[test]
fn a_file_nobody_queued_has_no_place() {
let queue = [queued("amy", false, 1)];
assert_eq!(place_in_queue(&queue, "amy", "@@share\\other.mp3"), None);
assert_eq!(place_in_queue(&queue, "zoe", "@@share\\amy.mp3"), None);
}
}
#[cfg(test)]
mod context_tests {
use crate::client::{ClientContext, DEFAULT_UPLOAD_SLOTS};
use crate::types::UploadStatus;
use std::path::PathBuf;
fn context(slots: usize) -> (ClientContext, impl FnMut() -> u32) {
let mut ctx = ClientContext::new();
ctx.upload_slots = slots;
let mut next = 0;
(ctx, move || {
next += 1;
next
})
}
fn ask(ctx: &mut ClientContext, user: &str, file: &str) {
ctx.enqueue_upload(user, user, file, PathBuf::from("/tmp/x"), 4096);
}
fn offered(
ctx: &mut ClientContext,
token: &mut impl FnMut() -> u32,
) -> Vec<String> {
ctx.pump_uploads(token)
.1
.into_iter()
.map(|offer| offer.requester_key)
.collect()
}
#[test]
fn a_fresh_context_starts_with_the_default_cap() {
assert_eq!(ClientContext::new().upload_slots, DEFAULT_UPLOAD_SLOTS);
}
#[test]
fn the_pump_fills_the_slots_and_leaves_the_rest_queued() {
let (mut ctx, mut token) = context(2);
for user in ["a", "b", "c"] {
ask(&mut ctx, user, "f.mp3");
}
assert_eq!(offered(&mut ctx, &mut token), ["a", "b"]);
assert_eq!(
ctx.place_in_queue("c", "f.mp3"),
Some(1),
"the third asker waits, now at the front of the queue"
);
assert!(
offered(&mut ctx, &mut token).is_empty(),
"pumping again must not oversubscribe the slots"
);
}
#[test]
fn a_privileged_list_arriving_late_re_ranks_who_is_waiting() {
let (mut ctx, mut token) = context(1);
ask(&mut ctx, "blocker", "f.mp3");
ask(&mut ctx, "plain", "f.mp3");
ask(&mut ctx, "donor", "f.mp3");
assert_eq!(offered(&mut ctx, &mut token), ["blocker"]);
assert_eq!(ctx.place_in_queue("plain", "f.mp3"), Some(1));
assert_eq!(ctx.place_in_queue("donor", "f.mp3"), Some(2));
ctx.set_privileged_users(vec!["donor".to_string()]);
assert!(ctx.is_privileged("donor"));
assert!(!ctx.is_privileged("plain"));
assert_eq!(
ctx.place_in_queue("donor", "f.mp3"),
Some(1),
"the donor should overtake the peer that asked first"
);
assert_eq!(
ctx.place_in_queue("plain", "f.mp3"),
Some(2),
"and the plain peer should be told they were overtaken"
);
}
#[test]
fn the_privileged_peer_is_the_one_served_next() {
let (mut ctx, mut token) = context(1);
ctx.set_privileged_users(vec!["donor".to_string()]);
ask(&mut ctx, "blocker", "f.mp3");
assert_eq!(offered(&mut ctx, &mut token), ["blocker"]);
ask(&mut ctx, "plain", "f.mp3");
ask(&mut ctx, "donor", "f.mp3");
ctx.uploads.clear(); assert_eq!(
offered(&mut ctx, &mut token),
["donor"],
"the freed slot goes to the donor, not the longer-waiting peer"
);
}
#[test]
fn a_finished_transfer_frees_its_slot_for_the_next_in_line() {
let (mut ctx, mut token) = context(1);
ask(&mut ctx, "first", "f.mp3");
ask(&mut ctx, "second", "f.mp3");
assert_eq!(offered(&mut ctx, &mut token), ["first"]);
ctx.uploads.clear();
assert_eq!(offered(&mut ctx, &mut token), ["second"]);
assert_eq!(ctx.place_in_queue("second", "f.mp3"), None);
}
#[test]
fn completed_transfers_do_not_keep_occupying_slots() {
let (mut ctx, mut token) = context(1);
ask(&mut ctx, "first", "f.mp3");
offered(&mut ctx, &mut token);
let job = ctx.uploads.remove(&1).expect("the offer was recorded");
ctx.active_uploads.insert(
1,
crate::client::ActiveUpload {
username: job.downloader,
filename: job.virtual_path,
size: job.size,
bytes_sent: std::sync::Arc::new(
std::sync::atomic::AtomicU64::new(job.size),
),
cancel: std::sync::Arc::new(
std::sync::atomic::AtomicBool::new(false),
),
status: UploadStatus::Completed,
started: std::time::Instant::now(),
},
);
ask(&mut ctx, "second", "f.mp3");
assert_eq!(
offered(&mut ctx, &mut token),
["second"],
"a completed upload must not hold its slot"
);
}
#[test]
fn queued_uploads_are_reported_in_the_order_they_will_be_served() {
let (mut ctx, mut token) = context(1);
ctx.set_privileged_users(vec!["donor".to_string()]);
ask(&mut ctx, "blocker", "f.mp3");
offered(&mut ctx, &mut token);
ask(&mut ctx, "plain", "f.mp3");
ask(&mut ctx, "donor", "f.mp3");
let reported: Vec<(String, UploadStatus)> = ctx
.queued_uploads()
.into_iter()
.map(|upload| (upload.username, upload.status))
.collect();
assert_eq!(
reported,
[
("donor".to_string(), UploadStatus::Queued(1)),
("plain".to_string(), UploadStatus::Queued(2)),
]
);
}
#[test]
fn a_queued_upload_is_recorded_even_after_it_leaves_the_queue() {
let (mut ctx, mut token) = context(1);
ask(&mut ctx, "amy", "@@share\\amy.mp3");
ask(&mut ctx, "bob", "@@share\\bob.mp3");
let _ = ctx.pump_uploads(&mut token);
assert!(
ctx.queued_uploads().iter().any(|u| u.username == "bob"),
"bob should be waiting while amy holds the only slot"
);
ctx.uploads.clear();
ctx.active_uploads.clear();
let _ = ctx.pump_uploads(&mut token);
assert!(
ctx.queued_uploads().is_empty(),
"the queue should have drained"
);
let events = ctx.take_upload_events();
assert!(
events.iter().any(|u| u.username == "bob"
&& matches!(u.status, UploadStatus::Queued(_))),
"bob's wait should still be reportable, got {events:?}"
);
assert!(
ctx.take_upload_events().is_empty(),
"draining twice must not repeat the same events"
);
}
}