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 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.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, 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 } // 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 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`, userID).Scan(&n) if err != nil { return 0, fmt.Errorf("count activities missing workout for user %d: %w", userID, err) } return n, nil }