zc2 0.0.28

P2P compute broker with credit-based billing, WAL, and broker mesh support
//! Policy pre-selection ahead of the broadcast race (marketplace task 5).
//!
//! Pure selector: given cached per-peer adverts (price + aggregated
//! resources) and an observed latency for each peer, pick the single best
//! candidate broker to offer a task to directly, before falling back to the
//! legacy broadcast race in `server::dispatch_to_peers`.
//!
//! Admission: a peer is only a candidate when its advert is present, fresh
//! (`now - epoch <= max_advert_age_secs`), and its aggregated resources
//! satisfy the request (cpus/memory/gpus available >= required AND at least
//! one healthy worker). Ranking (v1 policy): lowest observed latency first;
//! ties broken by lower price, then by lower epoch-age (fresher wins).

use super::discovery::Advert;
use super::router::ResourceRequirements;

/// A peer broker admitted and ranked as the best candidate for direct offer.
#[derive(Debug, Clone, PartialEq)]
pub struct AdmittedBroker {
    pub url: String,
    pub fp: String,
    pub price_per_hour: f64,
    pub latency_ms: f64,
}

/// Select the best executor broker among `peers` for `req`, or `None` if no
/// peer is admitted (caller should fall back to the broadcast race).
///
/// `peers` is `(url, advert, latency_ms)` per known peer. `now` is the
/// caller's current monotonic/wall-clock seconds, compared against each
/// advert's `epoch` for staleness.
pub fn select_executor_broker(
    peers: &[(String, Option<Advert>, f64)],
    req: &ResourceRequirements,
    max_advert_age_secs: u64,
    now: u64,
) -> Option<AdmittedBroker> {
    let mut best: Option<(AdmittedBroker, u64)> = None; // (candidate, epoch-age)

    for (url, advert, latency_ms) in peers {
        let advert = match advert {
            Some(a) => a,
            None => continue,
        };
        let age = now.saturating_sub(advert.epoch);
        if age > max_advert_age_secs {
            continue; // stale
        }
        let res = &advert.resources;
        if (res.cpus_available as f64) < req.cpus {
            continue;
        }
        if res.memory_available < req.memory_bytes {
            continue;
        }
        if res.gpus_available < req.gpus {
            continue;
        }
        if res.healthy_workers == 0 {
            continue;
        }

        let candidate = AdmittedBroker {
            url: url.clone(),
            fp: advert.fp.clone(),
            price_per_hour: advert.price_per_hour,
            latency_ms: *latency_ms,
        };

        best = match best {
            None => Some((candidate, age)),
            Some((ref cur, cur_age)) => {
                let better = if candidate.latency_ms != cur.latency_ms {
                    candidate.latency_ms < cur.latency_ms
                } else if candidate.price_per_hour != cur.price_per_hour {
                    candidate.price_per_hour < cur.price_per_hour
                } else {
                    age < cur_age
                };
                if better {
                    Some((candidate, age))
                } else {
                    best
                }
            }
        };
    }

    best.map(|(c, _)| c)
}

#[cfg(test)]
mod tests {
    use super::super::worker::BrokerResources;
    use super::*;

    fn advert(fp: &str, price: f64, resources: BrokerResources, epoch: u64) -> Advert {
        Advert {
            fp: fp.to_string(),
            price_per_hour: price,
            resources,
            epoch,
        }
    }

    #[test]
    fn select_prefers_lowest_latency_among_resource_sufficient_fresh_peers() {
        let now = 1000;
        let big = BrokerResources {
            cpus_available: 16,
            memory_available: 64_000_000_000,
            gpus_available: 2,
            healthy_workers: 3,
        };
        let small = BrokerResources {
            cpus_available: 1,
            memory_available: 1_000_000_000,
            gpus_available: 0,
            healthy_workers: 1,
        };
        let req = ResourceRequirements {
            cpus: 4.0,
            memory_bytes: 8_000_000_000,
            gpus: 1,
            ..Default::default()
        };
        let peers = vec![
            (
                "http://a:9000".into(),
                Some(advert("fpa", 5.0, big.clone(), now)),
                40.0,
            ), // admitted, 40ms
            (
                "http://b:9000".into(),
                Some(advert("fpb", 3.0, big.clone(), now)),
                12.0,
            ), // admitted, 12ms  <-- winner
            (
                "http://c:9000".into(),
                Some(advert("fpc", 1.0, small, now)),
                5.0,
            ), // rejected: too few resources
            (
                "http://d:9000".into(),
                Some(advert("fpd", 1.0, big, now.saturating_sub(9999))),
                1.0,
            ), // rejected: stale
        ];
        let sel = select_executor_broker(&peers, &req, 120, now).unwrap();
        assert_eq!(sel.url, "http://b:9000");
    }

    #[test]
    fn select_none_when_no_admitted_peer() {
        let req = ResourceRequirements {
            cpus: 999.0,
            ..Default::default()
        };
        assert!(select_executor_broker(&[], &req, 120, 0).is_none());
    }
}