238 lines
7.1 KiB
Go
238 lines
7.1 KiB
Go
package cron
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"log"
|
|
"time"
|
|
|
|
"trip-planner/internal/cache"
|
|
"trip-planner/internal/yandex"
|
|
)
|
|
|
|
// StationMonitor tracks the status and consecutive zero-trip days for a station.
|
|
type StationMonitor struct {
|
|
ID string
|
|
Yandex *yandex.Client
|
|
Cache cache.Cache
|
|
|
|
// ScheduleFunc is the function used to check a station's schedule.
|
|
// Defaults to checkStationSchedule if not set.
|
|
ScheduleFunc func(context.Context, string) (int, error)
|
|
|
|
// ZeroSince is the timestamp when the current zero-trip streak began.
|
|
// Zero if the station is not in a zero-trip streak.
|
|
ZeroSince time.Time
|
|
|
|
// LastSeenFlight is the timestamp of the last successful schedule check.
|
|
LastSeenFlight time.Time
|
|
}
|
|
|
|
// Status represents the current status of a station.
|
|
type Status string
|
|
|
|
const (
|
|
// StatusActive means the station has trips and is operating normally.
|
|
StatusActive Status = "active"
|
|
// StatusClosed means the station has had N consecutive days of zero trips.
|
|
StatusClosed Status = "closed"
|
|
)
|
|
|
|
// stationStatusKey returns the Redis key for station status.
|
|
func stationStatusKey(id string) *cache.CacheKey {
|
|
return &cache.CacheKey{
|
|
Kind: "station",
|
|
Code: id,
|
|
}
|
|
}
|
|
|
|
// zeroDaysKey returns the Redis key for tracking consecutive zero-trip days.
|
|
func zeroDaysKey(id string) *cache.CacheKey {
|
|
return &cache.CacheKey{
|
|
Kind: "station_zero_days",
|
|
Code: id,
|
|
}
|
|
}
|
|
|
|
// zeroSinceKey returns the Redis key for tracking the zero-trip streak start timestamp.
|
|
func zeroSinceKey(id string) *cache.CacheKey {
|
|
return &cache.CacheKey{
|
|
Kind: "station_zero_since",
|
|
Code: id,
|
|
}
|
|
}
|
|
|
|
// lastSeenFlightKey returns the Redis key for tracking the last seen flight timestamp.
|
|
func lastSeenFlightKey(id string) *cache.CacheKey {
|
|
return &cache.CacheKey{
|
|
Kind: "station_last_seen_flight",
|
|
Code: id,
|
|
}
|
|
}
|
|
|
|
// checkStationSchedule queries the Yandex /schedule endpoint for a station
|
|
// and returns the number of trips found.
|
|
func checkStationSchedule(ctx context.Context, yc *yandex.Client, stationID string) (int, error) {
|
|
// The Yandex Do method handles the API request with rate limiting,
|
|
// circuit breaking, and retry. It returns a Response with the
|
|
// schedule data including interval segments.
|
|
resp, err := yc.Do(ctx, "GET", "/station/"+stationID, map[string]string{
|
|
"date": time.Now().Format("2006-01-02"),
|
|
})
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
// The response contains Segments which represent trips/intervals
|
|
tripCount := len(resp.Segments)
|
|
|
|
return tripCount, nil
|
|
}
|
|
|
|
// updateStationStatus updates the station's status in cache based on trip count.
|
|
// It returns the new status. Writes status, zero-days count, zero-since timestamp,
|
|
// and last-seen-flight timestamp separately; partial failures may leave cache
|
|
// inconsistent but do not lose the core state.
|
|
func (sm *StationMonitor) updateStationStatus(ctx context.Context, tripCount int) (Status, error) {
|
|
cacheKey := stationStatusKey(sm.ID)
|
|
zeroDaysKey := zeroDaysKey(sm.ID)
|
|
zeroSinceKey := zeroSinceKey(sm.ID)
|
|
lastSeenFlightKey := lastSeenFlightKey(sm.ID)
|
|
|
|
// Get current status from cache
|
|
data, err := sm.Cache.Get(ctx, cacheKey)
|
|
var currentStatus Status
|
|
if err != nil {
|
|
currentStatus = StatusActive
|
|
} else if data != nil {
|
|
statusStr := string(data)
|
|
if statusStr == string(StatusClosed) {
|
|
currentStatus = StatusClosed
|
|
} else {
|
|
currentStatus = StatusActive
|
|
}
|
|
} else {
|
|
currentStatus = StatusActive
|
|
}
|
|
|
|
// Get current zero-trip day count
|
|
zeroDaysData, err := sm.Cache.Get(ctx, zeroDaysKey)
|
|
var zeroDays int
|
|
if err != nil {
|
|
zeroDays = 0
|
|
} else if zeroDaysData != nil {
|
|
var n int
|
|
_, err := fmt.Sscanf(string(zeroDaysData), "%d", &n)
|
|
if err == nil {
|
|
zeroDays = n
|
|
}
|
|
}
|
|
|
|
// Get current zero-since timestamp
|
|
zeroSinceData, err := sm.Cache.Get(ctx, zeroSinceKey)
|
|
var zeroSince time.Time
|
|
if err == nil && zeroSinceData != nil {
|
|
var zeroSinceUnix int64
|
|
_, err := fmt.Sscanf(string(zeroSinceData), "%d", &zeroSinceUnix)
|
|
if err == nil {
|
|
zeroSince = time.Unix(zeroSinceUnix, 0)
|
|
} else {
|
|
zeroSince = time.Time{}
|
|
}
|
|
}
|
|
|
|
// Get current last-seen-flight timestamp
|
|
lastSeenFlightData, err := sm.Cache.Get(ctx, lastSeenFlightKey)
|
|
var lastSeenFlight time.Time
|
|
if err == nil && lastSeenFlightData != nil {
|
|
var lastSeenFlightUnix int64
|
|
_, err := fmt.Sscanf(string(lastSeenFlightData), "%d", &lastSeenFlightUnix)
|
|
if err == nil {
|
|
lastSeenFlight = time.Unix(lastSeenFlightUnix, 0)
|
|
} else {
|
|
lastSeenFlight = time.Time{}
|
|
}
|
|
}
|
|
|
|
// Update status based on trip count
|
|
var newStatus Status
|
|
|
|
if tripCount > 0 {
|
|
newStatus = StatusActive
|
|
zeroDays = 0
|
|
zeroSince = time.Time{}
|
|
lastSeenFlight = time.Now()
|
|
} else {
|
|
zeroDays++
|
|
if zeroSince.IsZero() {
|
|
zeroSince = time.Now()
|
|
}
|
|
if zeroDays >= 3 {
|
|
newStatus = StatusClosed
|
|
} else {
|
|
newStatus = currentStatus
|
|
}
|
|
}
|
|
|
|
// Write updated status to cache with 24h TTL
|
|
if err := sm.Cache.Set(ctx, cacheKey, []byte(newStatus), 24*time.Hour); err != nil {
|
|
return "", fmt.Errorf("cache set status: %w", err)
|
|
}
|
|
|
|
// Write updated zero days count to cache with 24h TTL
|
|
if err := sm.Cache.Set(ctx, zeroDaysKey, []byte(fmt.Sprintf("%d", zeroDays)), 24*time.Hour); err != nil {
|
|
return "", fmt.Errorf("cache set zero days: %w", err)
|
|
}
|
|
|
|
// Write updated zero-since timestamp to cache with 24h TTL
|
|
if err := sm.Cache.Set(ctx, zeroSinceKey, []byte(fmt.Sprintf("%d", zeroSince.Unix())), 24*time.Hour); err != nil {
|
|
return "", fmt.Errorf("cache set zero since: %w", err)
|
|
}
|
|
|
|
// Write updated last-seen-flight timestamp to cache with 24h TTL
|
|
if err := sm.Cache.Set(ctx, lastSeenFlightKey, []byte(fmt.Sprintf("%d", lastSeenFlight.Unix())), 24*time.Hour); err != nil {
|
|
return "", fmt.Errorf("cache set last seen flight: %w", err)
|
|
}
|
|
|
|
return newStatus, nil
|
|
}
|
|
|
|
// ProcessStation checks a single station's schedule and updates its status.
|
|
// This function is designed to be called by a cron job or scheduler.
|
|
func ProcessStation(ctx context.Context, monitor *StationMonitor) error {
|
|
// Use the injected ScheduleFunc or the default checkStationSchedule
|
|
tripCount := 0
|
|
var err error
|
|
|
|
if monitor.ScheduleFunc != nil {
|
|
tripCount, err = monitor.ScheduleFunc(ctx, monitor.ID)
|
|
} else {
|
|
tripCount, err = checkStationSchedule(ctx, monitor.Yandex, monitor.ID)
|
|
}
|
|
if err != nil {
|
|
log.Printf("WARNING: failed to check schedule for station %s: %v", monitor.ID, err)
|
|
// If API fails, don't change the status - keep current
|
|
return nil
|
|
}
|
|
|
|
newStatus, err := monitor.updateStationStatus(ctx, tripCount)
|
|
if err != nil {
|
|
log.Printf("WARNING: failed to update status for station %s: %v", monitor.ID, err)
|
|
return err
|
|
}
|
|
|
|
log.Printf("INFO: station %s status updated to %s (trips today: %d)", monitor.ID, newStatus, tripCount)
|
|
return nil
|
|
}
|
|
|
|
// ProcessAllStations checks all monitored stations and updates their statuses.
|
|
// monitors is a list of StationMonitor instances for each station to check.
|
|
// This is the main function that a cron job would call.
|
|
func ProcessAllStations(ctx context.Context, monitors []*StationMonitor) error {
|
|
for _, monitor := range monitors {
|
|
if err := ProcessStation(ctx, monitor); err != nil {
|
|
log.Printf("ERROR: failed to process station %s: %v", monitor.ID, err)
|
|
}
|
|
}
|
|
return nil
|
|
}
|