cargagep-v2/src/core/controller/linka.rs
Antoine Pelletier 2d6ef8de6c feat: add linka
2026-08-25 12:10:53 +02:00

385 lines
14 KiB
Rust

//! Keeping this app and the Linka Go platform in step.
//!
//! Three things happen on every tick, in this order:
//!
//! 1. the fleet is read — lock state, battery, and who is riding what — and
//! kept as one snapshot the admin page reads without ever waiting on the
//! platform;
//! 2. a bike the platform reports out of service is taken out of service here
//! too (the platform is the truth: the mechanic works there, and the same
//! lock refuses to open either way);
//! 3. the access list is reconciled: everybody a live reservation entitles is
//! on it, and nobody else.
//!
//! The reconciliation replaces the flags the old system kept on each booking.
//! Flags could not survive a booking being edited, an address being added after
//! approval, or a call that failed once — a set difference survives all three,
//! and a failed call is simply retried at the next tick.
use std::{
collections::{HashSet, VecDeque},
sync::{Mutex, OnceLock},
};
use chrono::Utc;
use tracing::{debug, error, info, warn};
use crate::{
core::{
controller::{AnonAppController, ControllerError},
models::{
bike::{Bike, BikeStatus},
linka::{BikeLive, FleetLive, LiveBooking, Rider, UsageAlert},
},
},
services::{
linka::{self, LinkaError, locks, rentals, whitelist},
telegram::{self, Notification},
},
};
/// How much of the access list a pass goes over.
///
/// The platform cannot be read back, so this app works from its own record of
/// what it has posted. That record is only ever an assumption: somebody may
/// change the list on the platform, and an entry can go missing without this
/// app hearing about it. A [`Sweep::Full`] is what repairs that.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Sweep {
/// Only the difference with what this app believes it has already posted.
/// Cheap enough to run every minute.
Diff,
/// Every entitled address is asserted again, whatever the record says, and
/// everybody else is still taken off. Costs one call per live address.
Full,
}
/// Riders are let in half an hour before their reservation starts and taken off
/// the list half an hour after it ends — the time to walk to the bike, and to
/// bring it back.
const GRACE_MINUTES: i32 = 30;
impl AnonAppController {
/// One pass of the synchronisation. Never fails the caller: a platform that
/// is down leaves the app exactly as it was, and the next tick tries again.
pub async fn sync_linka(&self, sweep: Sweep) {
if !linka::configured() {
debug!("[LINKA] not configured, nothing to synchronise");
return;
}
self.sync_fleet().await;
self.sync_access_list(sweep).await;
}
/// Reads the fleet, follows the platform's service state, and raises the
/// alerts.
async fn sync_fleet(&self) {
let bikes = match self.db.get_bikes().await {
Ok(bikes) => bikes,
Err(err) => return error!("[LINKA] cannot read the fleet: {err}"),
};
let locks = match locks::fetch().await {
Ok(locks) => locks,
Err(err) => return warn!("[LINKA] cannot read the locks: {err}"),
};
let rides = match rentals::ongoing().await {
Ok(rides) => rides,
Err(err) => return warn!("[LINKA] cannot read the ongoing rides: {err}"),
};
let mut live = Vec::new();
for lock in &locks {
let Some(bike) = bikes.iter().find(|bike| bike.name == lock.number) else {
debug!("[LINKA] lock {} matches no bike here", lock.number);
continue;
};
// Whoever is on this very bike, if anybody
let rider = rides
.iter()
.find(|ride| ride.bikes.contains(&lock.number))
.map(|ride| Rider {
name: ride.rider.clone(),
email: ride.email.clone(),
since: ride.since,
});
live.push(BikeLive {
bike: bike.id,
lock_state: lock.state,
battery: lock.battery,
out_of_service: lock.out_of_service,
rider,
});
self.follow_service_state(bike, lock.out_of_service).await;
}
let alerts = match self.db.live_bookings().await {
Ok(bookings) => alerts(&rides, &bikes, &bookings),
Err(err) => {
error!("[LINKA] cannot read the live reservations: {err}");
Vec::new()
}
};
announce(&alerts);
linka::store(FleetLive {
updated_at: Some(Utc::now()),
bikes: live,
alerts,
});
}
/// The platform is where a bike is taken out of service for real; this app
/// follows. The other direction is pushed as it happens, in
/// `set_bike_status`.
async fn follow_service_state(&self, bike: &Bike, out_of_service: bool) {
let wanted = if out_of_service {
BikeStatus::OutOfService
} else {
BikeStatus::InService
};
if bike.status == wanted {
return;
}
info!(
"[LINKA] {} is {} on the platform: following",
bike.name,
if out_of_service {
"out of service"
} else {
"back in service"
}
);
if let Err(err) = self.db.set_bike_status(bike.id, wanted).await {
error!(
"[LINKA] cannot follow the service state of {}: {err}",
bike.name
);
}
}
/// Puts everybody a live reservation entitles on the platform's access
/// list, and takes off everybody else.
pub async fn sync_access_list(&self, sweep: Sweep) {
if !linka::configured() {
return;
}
let (wanted, current) = match (
self.db.emails_to_allow(GRACE_MINUTES).await,
self.db.allowed_emails().await,
) {
(Ok(wanted), Ok(current)) => (wanted, current),
(Err(err), _) | (_, Err(err)) => {
return error!("[LINKA] cannot work out the access list: {err}");
}
};
let wanted: HashSet<String> = wanted.into_iter().collect();
let current: HashSet<String> = current.into_iter().collect();
// A full sweep asks for everybody again, including those the record
// already counts as posted: an address the platform lost — taken off
// there, or an answer this app misread — is put back rather than
// missing until the reservation ends.
let to_allow: Vec<&String> = match sweep {
Sweep::Diff => wanted.difference(&current).collect(),
Sweep::Full => wanted.iter().collect(),
};
for email in to_allow {
match whitelist::allow(email).await {
// Recorded only once the platform has taken it: a failure is
// retried at the next tick rather than forgotten
Ok(()) => match self.db.record_allowed(email).await {
Ok(()) => info!("[LINKA] {email} may now unlock the bikes"),
Err(err) => error!("[LINKA] cannot record {email}: {err}"),
},
Err(LinkaError::DryRun) => {}
Err(err) => warn!("[LINKA] cannot allow {email}: {err}"),
}
}
for email in current.difference(&wanted) {
match whitelist::revoke(email).await {
Ok(()) => match self.db.forget_allowed(email).await {
Ok(()) => info!("[LINKA] {email} may no longer unlock the bikes"),
Err(err) => error!("[LINKA] cannot forget {email}: {err}"),
},
Err(LinkaError::DryRun) => {}
Err(err) => warn!("[LINKA] cannot revoke {email}: {err}"),
}
}
}
/// The last picture of the fleet, as read by the admin page
pub fn fleet_live(&self) -> Result<FleetLive, ControllerError> {
Ok(linka::snapshot())
}
}
/// Every bike being ridden by somebody no live reservation entitles to it.
///
/// Two shapes of trouble, told apart by whether the rider has a reservation
/// running at all — the group is told which, because the answer is not the
/// same: a stranger on a bike, or somebody on the wrong one.
fn alerts(rides: &[rentals::Rental], bikes: &[Bike], bookings: &[LiveBooking]) -> Vec<UsageAlert> {
let mut alerts = Vec::new();
for ride in rides {
let theirs: Vec<&LiveBooking> = bookings
.iter()
.filter(|booking| booking.covers(&ride.email))
.collect();
for number in &ride.bikes {
if theirs.iter().any(|booking| booking.allows(number)) {
continue;
}
let Some(bike) = bikes.iter().find(|bike| bike.name == *number) else {
continue;
};
alerts.push(UsageAlert {
bike: bike.id,
bike_name: bike.name.clone(),
rider: Rider {
name: ride.rider.clone(),
email: ride.email.clone(),
since: ride.since,
},
// Named when there is exactly one: "your reservation is for the
// 1000" only makes sense when there is one to point at
reservation: theirs.first().map(|booking| booking.reservation),
allowed: theirs
.iter()
.flat_map(|booking| booking.bikes.clone())
.collect(),
});
}
}
alerts
}
/// Tells the group about the alerts it has not been told about yet.
///
/// An episode is one rider on one bike: the message goes out when it starts,
/// and again only if it stops and starts anew. Without this the group would be
/// told once a minute for as long as the ride lasts.
fn announce(alerts: &[UsageAlert]) {
static ANNOUNCED: OnceLock<Mutex<VecDeque<String>>> = OnceLock::new();
/// Enough to remember the rides of a busy day; the oldest keys fall out
const REMEMBERED: usize = 64;
let mut announced = ANNOUNCED
.get_or_init(|| Mutex::new(VecDeque::new()))
.lock()
.expect("the announced alerts lock is poisoned");
let live: HashSet<String> = alerts.iter().map(UsageAlert::key).collect();
// An episode that is over is forgotten, so the next one is announced
announced.retain(|key| live.contains(key));
for alert in alerts {
let key = alert.key();
if announced.contains(&key) {
continue;
}
announced.push_back(key);
while announced.len() > REMEMBERED {
announced.pop_front();
}
telegram::notify(match alert.reservation {
Some(_) => Notification::BikeRiddenOutsideReservation(alert.clone()),
None => Notification::BikeRiddenWithoutReservation(alert.clone()),
});
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::core::models::bike::{BikeSize, BikeStatus};
fn bike(id: i32, name: &str) -> Bike {
Bike {
id,
name: name.to_owned(),
key_number: None,
key_quantity: 0,
drivetrain: None,
battery: None,
size: BikeSize::Large,
status: BikeStatus::InService,
}
}
fn ride(email: &str, bikes: &[&str]) -> rentals::Rental {
rentals::Rental {
rider: "Edgar Wolff".to_owned(),
email: email.to_owned(),
bikes: bikes.iter().map(|name| (*name).to_owned()).collect(),
since: None,
}
}
fn booking(id: i32, emails: &[&str], bikes: &[&str]) -> LiveBooking {
LiveBooking {
reservation: id,
emails: emails.iter().map(|email| (*email).to_owned()).collect(),
bikes: bikes.iter().map(|name| (*name).to_owned()).collect(),
}
}
#[test]
fn a_ride_covered_by_a_reservation_raises_nothing() {
let alerts = alerts(
&[ride("edgar@epfl.ch", &["3000"])],
&[bike(3, "3000")],
&[booking(7, &["Edgar@epfl.ch"], &["3000"])],
);
assert!(alerts.is_empty(), "{alerts:?}");
}
#[test]
fn a_ride_by_a_stranger_names_no_reservation() {
let alerts = alerts(
&[ride("nobody@epfl.ch", &["3000"])],
&[bike(3, "3000")],
&[booking(7, &["edgar@epfl.ch"], &["3000"])],
);
assert_eq!(alerts.len(), 1);
assert_eq!(alerts[0].bike, 3);
assert_eq!(alerts[0].reservation, None);
assert!(alerts[0].allowed.is_empty());
}
#[test]
fn a_ride_on_a_bike_the_reservation_does_not_hold_names_it() {
let alerts = alerts(
&[ride("edgar@epfl.ch", &["3000"])],
&[bike(1, "1000"), bike(3, "3000")],
&[booking(7, &["edgar@epfl.ch"], &["1000"])],
);
assert_eq!(alerts.len(), 1);
assert_eq!(alerts[0].bike_name, "3000");
assert_eq!(alerts[0].reservation, Some(7));
assert_eq!(alerts[0].allowed, vec!["1000"]);
}
#[test]
fn a_bike_this_app_does_not_know_is_ignored() {
let alerts = alerts(&[ride("a@epfl.ch", &["9000"])], &[bike(1, "1000")], &[]);
assert!(alerts.is_empty());
}
#[test]
fn two_reservations_of_the_same_rider_are_both_honoured() {
let alerts = alerts(
&[ride("edgar@epfl.ch", &["1000", "3000"])],
&[bike(1, "1000"), bike(3, "3000")],
&[
booking(7, &["edgar@epfl.ch"], &["1000"]),
booking(8, &["edgar@epfl.ch"], &["3000"]),
],
);
assert!(alerts.is_empty(), "{alerts:?}");
}
}