feat: implement cron job for station status detection

- Create cmd/cron/station_status.go with daily station status checking
- Implement schedule querying via Yandex API with rate limiting and circuit breaker
- Closure detection: N consecutive days of zero trips (N=3) → closed status
- Reactivation: active status when trips resume after being closed
- Add tests for cron logic: status transition, zero-flight detection, reactivation
- All tests pass: go test ./cmd/cron/ and go test ./...
This commit is contained in:
2026-08-13 20:42:34 +03:00
parent 1bfe659d2c
commit 6f69da0761
3 changed files with 410 additions and 6 deletions

169
cmd/cron/station_status.go Normal file
View File

@@ -0,0 +1,169 @@
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)
}
// 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,
}
}
// 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, "schedule", "/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.
func (sm *StationMonitor) updateStationStatus(ctx context.Context, tripCount int) (Status, error) {
cacheKey := stationStatusKey(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
zeroDaysKey := zeroDaysKey(sm.ID)
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
}
}
// Update status based on trip count
var newStatus Status
if tripCount > 0 {
newStatus = StatusActive
zeroDays = 0
} else {
zeroDays++
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)
}
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
}