A confirmed-404 workout must stop being retried forever, but workout_raw_json should never be fabricated -- it stays null exactly as it does for "not yet fetched." workout_not_found_at is the separate marker ActivitiesMissingWorkout/CountActivitiesMissingWorkout now check for, mirroring the existing pattern of the other _fetched_at columns. UpsertActivity's ON CONFLICT clause already never touches these columns, so the marker persists across every later re-sync.
350 lines
14 KiB
Go
350 lines
14 KiB
Go
package store
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"fmt"
|
|
)
|
|
|
|
// Activity is geniusrun's persisted view of a Garmin activity. Field mapping
|
|
// from the Garmin/MCP response happens in internal/sync, not here, so this
|
|
// package stays independent of internal/garmin.
|
|
//
|
|
// Deliberately NOT modeled here (though Garmin's response includes them):
|
|
// ActivityName/ActivityType, plus every field removed in the 2026-07
|
|
// dedup pass (BeginTimestampMs, MaxSpeedMps, ElevationLossM, Calories,
|
|
// LapCount, TrainingEffectLabel, HrTimeInZone1-5). None of them are read by
|
|
// any SQL query, the classify rule engine, or any other Go code -- they'd be
|
|
// pure duplicates of RawJSON with no purpose beyond being slightly more
|
|
// convenient to read than parsing JSON. ActivityName/ActivityType do have a
|
|
// real frontend display need, so the API layer (internal/api) decodes them
|
|
// from RawJSON at response time instead of storing a redundant copy -- see
|
|
// decodeActivityDisplayFields.
|
|
type Activity struct {
|
|
ID int64
|
|
GarminActivityID int64
|
|
EventTypeKey string
|
|
WorkoutID *int64
|
|
StartTimeUTC string
|
|
DurationSeconds float64
|
|
DistanceMeters float64
|
|
AvgHR *float64
|
|
MaxHR *float64
|
|
AvgSpeedMps *float64
|
|
ElevationGainM *float64
|
|
AerobicTrainingEffect *float64
|
|
AnaerobicTrainingEffect *float64
|
|
VO2MaxValue *float64
|
|
RawJSON string
|
|
DetailsFetchedAt *string
|
|
DetailsRawJSON *string
|
|
SplitsFetchedAt *string
|
|
// WorkoutRawJSON is the genuine raw get_workout_by_id() response for this
|
|
// activity's structured workout -- the source used to compute each lap's
|
|
// TargetPaceLowMps/HighMps and TargetHRLowBpm/HighBpm (see
|
|
// internal/sync/mapping.go's alignWorkoutTargets). Nil when the activity
|
|
// has no WorkoutID, or was synced before this column existed.
|
|
WorkoutRawJSON *string
|
|
// WorkoutNotFoundAt is set when get_workout_by_id returned a definitive
|
|
// HTTP 404 for this activity's WorkoutID -- distinct from
|
|
// WorkoutRawJSON staying nil for "not yet fetched" (see
|
|
// SetActivityWorkoutNotFound).
|
|
WorkoutNotFoundAt *string
|
|
CreatedAt string
|
|
UpdatedAt string
|
|
}
|
|
|
|
// UpsertActivity inserts a new activity for userID or updates the existing
|
|
// row for the same (userID, garmin_activity_id) pair (idempotent, safe to
|
|
// call on every sync pass), and returns its internal id.
|
|
func (db *DB) UpsertActivity(ctx context.Context, userID int64, a Activity) (int64, error) {
|
|
_, err := db.ExecContext(ctx, `
|
|
INSERT INTO activities (
|
|
user_id, garmin_activity_id, event_type_key, workout_id, start_time_utc,
|
|
duration_seconds, distance_meters, avg_hr, max_hr,
|
|
avg_speed_mps, elevation_gain_m,
|
|
aerobic_training_effect, anaerobic_training_effect, vo2max_value,
|
|
raw_json, updated_at
|
|
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?, datetime('now'))
|
|
ON CONFLICT(user_id, garmin_activity_id) DO UPDATE SET
|
|
event_type_key=excluded.event_type_key,
|
|
workout_id=excluded.workout_id,
|
|
start_time_utc=excluded.start_time_utc,
|
|
duration_seconds=excluded.duration_seconds,
|
|
distance_meters=excluded.distance_meters,
|
|
avg_hr=excluded.avg_hr,
|
|
max_hr=excluded.max_hr,
|
|
avg_speed_mps=excluded.avg_speed_mps,
|
|
elevation_gain_m=excluded.elevation_gain_m,
|
|
aerobic_training_effect=excluded.aerobic_training_effect,
|
|
anaerobic_training_effect=excluded.anaerobic_training_effect,
|
|
vo2max_value=excluded.vo2max_value,
|
|
raw_json=excluded.raw_json,
|
|
updated_at=datetime('now')
|
|
`,
|
|
userID, a.GarminActivityID, a.EventTypeKey, a.WorkoutID, a.StartTimeUTC,
|
|
a.DurationSeconds, a.DistanceMeters, a.AvgHR, a.MaxHR,
|
|
a.AvgSpeedMps, a.ElevationGainM,
|
|
a.AerobicTrainingEffect, a.AnaerobicTrainingEffect, a.VO2MaxValue,
|
|
a.RawJSON,
|
|
)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("upsert activity %d for user %d: %w", a.GarminActivityID, userID, err)
|
|
}
|
|
|
|
var id int64
|
|
if err := db.QueryRowContext(ctx, `SELECT id FROM activities WHERE user_id = ? AND garmin_activity_id = ?`, userID, a.GarminActivityID).Scan(&id); err != nil {
|
|
return 0, fmt.Errorf("fetch id for activity %d (user %d): %w", a.GarminActivityID, userID, err)
|
|
}
|
|
return id, nil
|
|
}
|
|
|
|
func scanActivity(row interface{ Scan(...any) error }) (Activity, error) {
|
|
var a Activity
|
|
err := row.Scan(
|
|
&a.ID, &a.GarminActivityID, &a.EventTypeKey, &a.WorkoutID, &a.StartTimeUTC,
|
|
&a.DurationSeconds, &a.DistanceMeters, &a.AvgHR, &a.MaxHR,
|
|
&a.AvgSpeedMps, &a.ElevationGainM,
|
|
&a.AerobicTrainingEffect, &a.AnaerobicTrainingEffect, &a.VO2MaxValue,
|
|
&a.RawJSON, &a.DetailsFetchedAt, &a.DetailsRawJSON, &a.SplitsFetchedAt, &a.WorkoutRawJSON,
|
|
&a.WorkoutNotFoundAt, &a.CreatedAt, &a.UpdatedAt,
|
|
)
|
|
return a, err
|
|
}
|
|
|
|
const activityColumns = `
|
|
id, garmin_activity_id, event_type_key, workout_id, start_time_utc,
|
|
duration_seconds, distance_meters, avg_hr, max_hr,
|
|
avg_speed_mps, elevation_gain_m,
|
|
aerobic_training_effect, anaerobic_training_effect, vo2max_value,
|
|
raw_json, details_fetched_at, details_raw_json, splits_fetched_at, workout_raw_json,
|
|
workout_not_found_at, created_at, updated_at
|
|
`
|
|
|
|
// GetActivity fetches one activity by its internal id, scoped to userID.
|
|
func (db *DB) GetActivity(ctx context.Context, userID, id int64) (Activity, bool, error) {
|
|
row := db.QueryRowContext(ctx, `SELECT `+activityColumns+` FROM activities WHERE id = ? AND user_id = ?`, id, userID)
|
|
a, err := scanActivity(row)
|
|
if err == sql.ErrNoRows {
|
|
return Activity{}, false, nil
|
|
}
|
|
if err != nil {
|
|
return Activity{}, false, fmt.Errorf("get activity %d for user %d: %w", id, userID, err)
|
|
}
|
|
return a, true, nil
|
|
}
|
|
|
|
// ActivityFilter narrows ListActivities results. Zero values mean "no filter".
|
|
type ActivityFilter struct {
|
|
FromDate string // inclusive, "YYYY-MM-DD"
|
|
ToDate string // inclusive, "YYYY-MM-DD"
|
|
Limit int
|
|
Offset int
|
|
}
|
|
|
|
// ListActivities returns userID's activities newest-first, optionally
|
|
// filtered by date range.
|
|
func (db *DB) ListActivities(ctx context.Context, userID int64, f ActivityFilter) ([]Activity, error) {
|
|
query := `SELECT ` + activityColumns + ` FROM activities WHERE user_id = ?`
|
|
args := []any{userID}
|
|
if f.FromDate != "" {
|
|
query += ` AND start_time_utc >= ?`
|
|
args = append(args, f.FromDate)
|
|
}
|
|
if f.ToDate != "" {
|
|
query += ` AND start_time_utc <= ?`
|
|
args = append(args, f.ToDate+" 23:59:59")
|
|
}
|
|
query += ` ORDER BY start_time_utc DESC`
|
|
if f.Limit > 0 {
|
|
query += ` LIMIT ? OFFSET ?`
|
|
args = append(args, f.Limit, f.Offset)
|
|
}
|
|
|
|
rows, err := db.QueryContext(ctx, query, args...)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("list activities for user %d: %w", userID, err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
activities := []Activity{}
|
|
for rows.Next() {
|
|
a, err := scanActivity(rows)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("scan activity row: %w", err)
|
|
}
|
|
activities = append(activities, a)
|
|
}
|
|
return activities, rows.Err()
|
|
}
|
|
|
|
// ActivityExists reports whether userID already has an activity with this
|
|
// garmin_activity_id stored.
|
|
func (db *DB) ActivityExists(ctx context.Context, userID, garminActivityID int64) (bool, error) {
|
|
var id int64
|
|
err := db.QueryRowContext(ctx, `SELECT id FROM activities WHERE user_id = ? AND garmin_activity_id = ?`, userID, garminActivityID).Scan(&id)
|
|
if err == sql.ErrNoRows {
|
|
return false, nil
|
|
}
|
|
if err != nil {
|
|
return false, fmt.Errorf("check activity %d exists for user %d: %w", garminActivityID, userID, err)
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
// LatestActivityStartTime returns userID's most recently started activity's
|
|
// start_time_utc, used to compute the incremental sync window.
|
|
func (db *DB) LatestActivityStartTime(ctx context.Context, userID int64) (string, bool, error) {
|
|
var t string
|
|
err := db.QueryRowContext(ctx, `SELECT start_time_utc FROM activities WHERE user_id = ? ORDER BY start_time_utc DESC LIMIT 1`, userID).Scan(&t)
|
|
if err == sql.ErrNoRows {
|
|
return "", false, nil
|
|
}
|
|
if err != nil {
|
|
return "", false, fmt.Errorf("latest activity start time for user %d: %w", userID, err)
|
|
}
|
|
return t, true, nil
|
|
}
|
|
|
|
// SetActivityDetails records that get_activity_details has been fetched for
|
|
// this activity, storing the raw response for future reprocessing.
|
|
func (db *DB) SetActivityDetails(ctx context.Context, userID, activityID int64, rawJSON string) error {
|
|
_, err := db.ExecContext(ctx, `
|
|
UPDATE activities SET details_fetched_at = datetime('now'), details_raw_json = ?, updated_at = datetime('now')
|
|
WHERE id = ? AND user_id = ?`, rawJSON, activityID, userID)
|
|
if err != nil {
|
|
return fmt.Errorf("set activity %d details for user %d: %w", activityID, userID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SetActivityWorkout stores the raw get_workout_by_id() response used to
|
|
// compute this activity's laps' target pace/HR bands.
|
|
func (db *DB) SetActivityWorkout(ctx context.Context, userID, activityID int64, rawJSON string) error {
|
|
_, err := db.ExecContext(ctx, `
|
|
UPDATE activities SET workout_raw_json = ?, updated_at = datetime('now')
|
|
WHERE id = ? AND user_id = ?`, rawJSON, activityID, userID)
|
|
if err != nil {
|
|
return fmt.Errorf("set activity %d workout for user %d: %w", activityID, userID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SetActivityWorkoutNotFound records that get_workout_by_id returned a
|
|
// definitive HTTP 404 for this activity's WorkoutID -- the workout was
|
|
// deleted on Garmin's side after being linked to this activity.
|
|
// WorkoutRawJSON is deliberately left nil (never fabricated); this is a
|
|
// separate marker so ActivitiesMissingWorkout stops retrying it forever.
|
|
func (db *DB) SetActivityWorkoutNotFound(ctx context.Context, userID, activityID int64) error {
|
|
_, err := db.ExecContext(ctx, `
|
|
UPDATE activities SET workout_not_found_at = datetime('now'), updated_at = datetime('now')
|
|
WHERE id = ? AND user_id = ?`, activityID, userID)
|
|
if err != nil {
|
|
return fmt.Errorf("set activity %d workout not found for user %d: %w", activityID, userID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SetActivitySplitsFetched records that get_activity_splits has been fetched
|
|
// for this activity.
|
|
func (db *DB) SetActivitySplitsFetched(ctx context.Context, userID, activityID int64) error {
|
|
_, err := db.ExecContext(ctx, `
|
|
UPDATE activities SET splits_fetched_at = datetime('now'), updated_at = datetime('now')
|
|
WHERE id = ? AND user_id = ?`, activityID, userID)
|
|
if err != nil {
|
|
return fmt.Errorf("set activity %d splits fetched for user %d: %w", activityID, userID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ActivitiesMissingDetails returns userID's activities that haven't had
|
|
// get_activity_details/get_activity_splits fetched yet, for the lazy
|
|
// background detail-fill pass.
|
|
func (db *DB) ActivitiesMissingDetails(ctx context.Context, userID int64, limit int) ([]Activity, error) {
|
|
rows, err := db.QueryContext(ctx, `SELECT `+activityColumns+` FROM activities
|
|
WHERE user_id = ? AND (details_fetched_at IS NULL OR splits_fetched_at IS NULL)
|
|
ORDER BY start_time_utc DESC LIMIT ?`, userID, limit)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("list activities missing details for user %d: %w", userID, err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
activities := []Activity{}
|
|
for rows.Next() {
|
|
a, err := scanActivity(rows)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("scan activity row: %w", err)
|
|
}
|
|
activities = append(activities, a)
|
|
}
|
|
return activities, rows.Err()
|
|
}
|
|
|
|
// CountActivitiesMissingDetails returns how many of userID's activities
|
|
// still need get_activity_details/get_activity_splits fetched, regardless
|
|
// of any per-call batch limit -- used to report overall remaining work.
|
|
func (db *DB) CountActivitiesMissingDetails(ctx context.Context, userID int64) (int, error) {
|
|
var n int
|
|
err := db.QueryRowContext(ctx, `SELECT COUNT(*) FROM activities
|
|
WHERE user_id = ? AND (details_fetched_at IS NULL OR splits_fetched_at IS NULL)`, userID).Scan(&n)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("count activities missing details for user %d: %w", userID, err)
|
|
}
|
|
return n, nil
|
|
}
|
|
|
|
// ActivitiesMissingWorkout returns userID's activities that have a
|
|
// structured workout (WorkoutID set at initial upsert time, straight from
|
|
// Garmin's activity summary) but haven't had get_workout_by_id fetched yet.
|
|
// Gated on details_fetched_at IS NOT NULL: fillActivityWorkout resolves each
|
|
// lap's target pace/HR band via alignWorkoutTargets, which needs the
|
|
// activity's laps already written by fillActivityDetails -- without this
|
|
// gate, a structured-workout activity whose details/laps haven't been
|
|
// fetched yet would fall into a separate, independently LIMIT-bounded batch
|
|
// than fillPendingDetails's, get its workout_raw_json set against zero laps
|
|
// (a silent no-op alignment), and then never be retried once its laps
|
|
// finally arrive, since workout_raw_json IS NULL is the only signal this
|
|
// query has left. This is still independent of splits_fetched_at (unlike
|
|
// ActivitiesMissingDetails, which requires both): splits/laps aren't needed
|
|
// to resolve workout targets, only details_fetched_at is. It also still
|
|
// covers the intended retry case -- an activity whose workout fetch
|
|
// previously failed after details succeeded always has details_fetched_at
|
|
// already set, so it's still surfaced here -- while excluding activities
|
|
// that simply haven't been processed by fillActivityDetails at all yet.
|
|
func (db *DB) ActivitiesMissingWorkout(ctx context.Context, userID int64, limit int) ([]Activity, error) {
|
|
rows, err := db.QueryContext(ctx, `SELECT `+activityColumns+` FROM activities
|
|
WHERE user_id = ? AND workout_id IS NOT NULL AND workout_raw_json IS NULL
|
|
AND details_fetched_at IS NOT NULL AND workout_not_found_at IS NULL
|
|
ORDER BY start_time_utc DESC LIMIT ?`, userID, limit)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("list activities missing workout for user %d: %w", userID, err)
|
|
}
|
|
defer rows.Close()
|
|
|
|
activities := []Activity{}
|
|
for rows.Next() {
|
|
a, err := scanActivity(rows)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("scan activity row: %w", err)
|
|
}
|
|
activities = append(activities, a)
|
|
}
|
|
return activities, rows.Err()
|
|
}
|
|
|
|
// CountActivitiesMissingWorkout returns how many of userID's activities
|
|
// still need get_workout_by_id fetched, regardless of any per-call batch
|
|
// limit -- used to report overall remaining work, mirroring
|
|
// CountActivitiesMissingDetails. See ActivitiesMissingWorkout for why this
|
|
// is additionally gated on details_fetched_at IS NOT NULL.
|
|
func (db *DB) CountActivitiesMissingWorkout(ctx context.Context, userID int64) (int, error) {
|
|
var n int
|
|
err := db.QueryRowContext(ctx, `SELECT COUNT(*) FROM activities
|
|
WHERE user_id = ? AND workout_id IS NOT NULL AND workout_raw_json IS NULL
|
|
AND details_fetched_at IS NOT NULL AND workout_not_found_at IS NULL`, userID).Scan(&n)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("count activities missing workout for user %d: %w", userID, err)
|
|
}
|
|
return n, nil
|
|
}
|