5 Commits

5 changed files with 85 additions and 16 deletions
+4
View File
@@ -8,10 +8,14 @@ host = "x86_64-unknown-linux-gnu"
chrono = "0.4" chrono = "0.4"
csv = "1.4" csv = "1.4"
gtfs-structures = "0.47" gtfs-structures = "0.47"
gtfs-realtime = "0.2"
log = "0.4" log = "0.4"
prost = "0.14"
prost-types = "0.14"
reqwest = { version = "0.11", features = ["blocking"] } reqwest = { version = "0.11", features = ["blocking"] }
sdl3 = {version = "0.17", features = ["ttf"]} sdl3 = {version = "0.17", features = ["ttf"]}
serde = { version = "1.0", features = ["derive"]} serde = { version = "1.0", features = ["derive"]}
serde_json = "1.0"
time-format = "1.2" time-format = "1.2"
yaml_serde = "0.10" yaml_serde = "0.10"
zip = "8.3" zip = "8.3"
+13 -2
View File
@@ -5,7 +5,7 @@ mod refresher;
mod realtime; mod realtime;
pub mod structs; pub mod structs;
use chrono::{DateTime, Local, Timelike}; use chrono::{DateTime, Local, Timelike};
use log::{debug, trace}; use log::{debug, trace, warn};
use std::{collections::{HashMap, HashSet}, fs::File }; use std::{collections::{HashMap, HashSet}, fs::File };
use gtfs_structures::{Exception, RawTrip}; use gtfs_structures::{Exception, RawTrip};
use crate::gtfs::{loader::load_gtfs, structs::{Arrival, Gtfs, Preferences, Error}}; use crate::gtfs::{loader::load_gtfs, structs::{Arrival, Gtfs, Preferences, Error}};
@@ -18,12 +18,14 @@ impl Gtfs<'_> {
let target_date = naive_target.date(); let target_date = naive_target.date();
// Find which calendars apply // Find which calendars apply
debug!("Looking for calendars that apply to date {:#?}", target_date);
let mut active_service_ids: HashSet<String> = HashSet::new(); let mut active_service_ids: HashSet<String> = HashSet::new();
for (id, calendar) in self.calendar.iter() { for (id, calendar) in self.calendar.iter() {
if calendar.valid_weekday(target_date) if calendar.valid_weekday(target_date)
&& calendar.start_date <= target_date && calendar.start_date <= target_date
&& calendar.end_date > target_date { && calendar.end_date > target_date {
active_service_ids.insert(id.to_string()); active_service_ids.insert(id.to_string());
debug!("Matched calendar: {:#?}", calendar);
} }
} }
debug!("Found {} services active today", active_service_ids.len()); debug!("Found {} services active today", active_service_ids.len());
@@ -43,7 +45,7 @@ impl Gtfs<'_> {
} }
} }
} }
debug!("After exceptions, there are now {} services active", active_service_ids.len()); debug!("After exceptions, there are now {} services active: {:#?}", active_service_ids.len(), active_service_ids);
// Find the trips happening on these calendars // Find the trips happening on these calendars
let mut trips: HashMap<&String, &RawTrip> = HashMap::new(); let mut trips: HashMap<&String, &RawTrip> = HashMap::new();
@@ -82,9 +84,18 @@ impl Gtfs<'_> {
} }
} }
} }
arrivals.sort(); arrivals.sort();
debug!("Found {} arrivals", arrivals.len()); debug!("Found {} arrivals", arrivals.len());
// Update real-time deltas
let realtime_result = self.realtime_update(&mut arrivals);
if realtime_result.is_err() {
warn!("Unable to update realtime arrivals: {:#?}", realtime_result.err().unwrap()._message)
}
return Some(arrivals); return Some(arrivals);
} }
+60 -12
View File
@@ -1,15 +1,14 @@
use log::debug;
/*** /***
* Implementation of GTFS-R polling * Implementation of GTFS-R polling
*/ */
use reqwest::{StatusCode, blocking::Client};
use std::{collections::HashMap, time::Duration};
use std::{str::FromStr, time::Duration}; use crate::gtfs::{self, structs::{Arrival, Error, Gtfs}};
use reqwest::{blocking::Client, header::{HeaderName, HeaderValue}};
use crate::gtfs::structs::{Arrival, Error, Gtfs};
impl Gtfs<'_> { impl Gtfs<'_> {
fn realtime_update (&self, arrivals: Vec<Arrival<'_>>) -> Result<bool, Error>{ pub(crate) fn realtime_update (&self, arrivals: &mut Vec<Arrival<'_>>) -> Result<(), gtfs::structs::Error>{
// Poll GTFS-R API // Poll GTFS-R API
let client = Client::builder() let client = Client::builder()
@@ -17,20 +16,69 @@ impl Gtfs<'_> {
.connect_timeout(Duration::from_secs(5)) .connect_timeout(Duration::from_secs(5))
.build()?; .build()?;
let mut response = client let response = client
.get(self.preferences.realtime_url) .get(&self.preferences.realtime_url)
.header("x-api-key", self.preferences.realtime_api_key) .header("x-api-key", &self.preferences.realtime_api_key)
.send()?; .send()?;
if response.status() != StatusCode::OK {
return Err(Error {
_message : format!("HTTP Respomse: {:#?}. Payload: \n{:#?}\n", response.status(), response.text().unwrap())
})
}
// Parse response // Parse response
let response = json::rdr let response_bytes= response.bytes()?;
let data: Result<gtfs_realtime::FeedMessage, prost::DecodeError> = prost::Message::decode(response_bytes.as_ref());
if data.is_err() {
return Err(Error {
_message : format! ("Error loading realtime prtobuf: {:#?}", data.err().unwrap())
})
}
// Build a map of (trip, stop) -> Arrival for faster lookup
let mut lookup: HashMap<(String, String), usize> = HashMap::new();
let num_arrivals = arrivals.len();
for i in 0..num_arrivals {
let arrival = arrivals.get(i).unwrap();
lookup.insert((arrival.trip.id.clone(), arrival.stop.id.clone()), i);
}
// Match deltas to existing arrivals // Match deltas to existing arrivals
let entities = data.unwrap().entity;
debug!("Loaded {} entities from realtime update.", entities.len());
for entity in entities {
if entity.trip_update.is_some() && entity.stop.is_some() {
let trip_update = entity.trip_update.unwrap();
// Update the arrival times // Look up the entry in arrivals for the current
for stop_update in trip_update.stop_time_update {
if stop_update.stop_id.is_none() {
continue;
}
let trip_id = trip_update.trip.trip_id.clone().unwrap();
let stop_id = stop_update.stop_id.unwrap();
let key = (trip_id, stop_id);
if lookup.contains_key(&key) {
// Update the arrival time
let arrival_index = lookup.get(&key);
let arrival: &mut Arrival<'_> = arrivals.get_mut(*arrival_index.unwrap()).unwrap();
if let Some(update) = stop_update.departure.or(stop_update.arrival) {
let new_time = (((arrival.departure_time as i64)
+ (update.delay.unwrap_or(0) as i64))) as u32;
debug!("Updated arrival time of {:#?} by {}s", &key, update.delay.unwrap());
arrival.departure_time = new_time;
}
}
}
}
}
return Ok(());
}
} }
}
+6
View File
@@ -28,6 +28,12 @@ impl From<ZipError> for Error {
} }
} }
impl From<serde_json::Error> for Error {
fn from(value: serde_json::Error) -> Self {
return Error { _message: value.to_string() }
}
}
// This is to store the preferences for the GTFS(-R) side of the code. // This is to store the preferences for the GTFS(-R) side of the code.
#[derive(Debug, PartialEq, Serialize, Deserialize)] #[derive(Debug, PartialEq, Serialize, Deserialize)]
+1 -1
View File
@@ -36,7 +36,7 @@ impl Screen<'_> {
fn format_due_for(&self, due_in_mins: i32, departure_time: u32) -> String { fn format_due_for(&self, due_in_mins: i32, departure_time: u32) -> String {
trace!("Due in mins: {:02}", due_in_mins); trace!("Due in mins: {:02}", due_in_mins);
if due_in_mins <= 1 { if due_in_mins <= 1 {
return String::from("due"); return String::from("Due");
} }
if due_in_mins < 60 { if due_in_mins < 60 {
return due_in_mins.to_string() + "min"; return due_in_mins.to_string() + "min";