diff --git a/server/src/subway_gleam/server/state/gtfs_store.gleam b/server/src/subway_gleam/server/state/gtfs_store.gleam index 939a7a2..95346ea 100644 --- a/server/src/subway_gleam/server/state/gtfs_store.gleam +++ b/server/src/subway_gleam/server/state/gtfs_store.gleam @@ -1,9 +1,11 @@ import gleam/dict import gleam/erlang/atom import gleam/erlang/process +import gleam/httpc import gleam/int import gleam/list import gleam/result +import gleam/string import gleam/time/duration import gleam/time/timestamp.{type Timestamp} @@ -71,42 +73,58 @@ pub fn update(store: GtfsStore) -> Nil { // TODO: is this OTP compliant? // If this process crashes it shouldn't take down the main process process.spawn_unlinked(fn() { - use #(rt.Data(arrivals:, final_stops:, trips:, alerts:), last_updated) <- result.map( - fetch_all_rt_feeds(), - ) - let arrivals = dict.to_list(arrivals) - let final_stops = dict.to_list(final_stops) - let trips = dict.to_list(trips) - let alerts = [#(Nil, alerts)] - - // Update data - // TODO: this is non-atomic. data will update before last-updated is set. - // Maybe make one table and make a larger Key type? - table.replace_all_items(in: store.data.arrivals, with: arrivals) - table.replace_all_items(in: store.data.final_stops, with: final_stops) - table.replace_all_items(in: store.data.trips, with: trips) - table.replace_all_items(in: store.data.alerts, with: alerts) - table.insert(last_updated, into: store.data.last_updated, for: Nil) - - // Notify watchers - store.watchers - |> table.each(fn(subject, _: Nil) { process.send(subject, Nil) }) - - let end_time = timestamp.system_time() - let duration = timestamp.difference(start_time, end_time) - log.debug( - "Finished gtfs rt update.", - with: log.context([ - #( - "duration", - duration.to_milliseconds(duration) |> int.to_string <> "ms", - ), - ]), - ) + let res = do_update(store) + + let duration = { + let end_time = timestamp.system_time() + let duration = timestamp.difference(start_time, end_time) + #( + "duration", + duration |> duration.to_milliseconds |> int.to_string <> "ms", + ) + } + + case res { + Ok(Nil) -> + log.debug("Finished gtfs rt update.", with: log.context([duration])) + Error(error) -> { + // TODO: don't use string.inspect for this + let error = string.inspect(error) + log.error( + "Could not fetch gtfs rt update.", + with: log.context([duration, #("error", error)]), + ) + } + } }) Nil } +fn do_update(store: GtfsStore) -> Result(Nil, rt.FetchGtfsError) { + use #(rt.Data(arrivals:, final_stops:, trips:, alerts:), last_updated) <- result.map( + fetch_all_rt_feeds(), + ) + let arrivals = dict.to_list(arrivals) + let final_stops = dict.to_list(final_stops) + let trips = dict.to_list(trips) + let alerts = [#(Nil, alerts)] + + // Update data + // TODO: this is non-atomic. data will update before last-updated is set. + // Maybe make one table and make a larger Key type? + table.replace_all_items(in: store.data.arrivals, with: arrivals) + table.replace_all_items(in: store.data.final_stops, with: final_stops) + table.replace_all_items(in: store.data.trips, with: trips) + table.replace_all_items(in: store.data.alerts, with: alerts) + table.insert(last_updated, into: store.data.last_updated, for: Nil) + + // Notify watchers + store.watchers + |> table.each(fn(subject, _: Nil) { process.send(subject, Nil) }) + + Nil +} + pub fn subscribe_watcher( to store: GtfsStore, watcher watcher: process.Subject(Nil),