Mirrors ActivitiesMissingDetails/CountActivitiesMissingDetails, but queries workout_id IS NOT NULL AND workout_raw_json IS NULL -- independent of whether details/splits were ever fetched, since workout_id is known at initial upsert time. This also fixes a latent gap: an activity whose details were fetched successfully but whose workout fetch failed in the same run previously had no way to ever be retried, since ActivitiesMissingDetails stops returning it the moment details_fetched_at/splits_fetched_at are set.
319 lines
13 KiB
Go
319 lines
13 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
|
|
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.
|
|
// Independently queryable from ActivitiesMissingDetails: workout_id is
|
|
// known well before any detail fetch, and workout_raw_json is only ever
|
|
// set by SetActivityWorkout, so this also picks up an activity whose
|
|
// details were fetched successfully in some prior run but whose workout
|
|
// fetch failed back then -- ActivitiesMissingDetails would never surface
|
|
// that activity again (details_fetched_at/splits_fetched_at are already
|
|
// set), silently losing its target pace/HR bands forever without this.
|
|
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
|
|
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.
|
|
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`, userID).Scan(&n)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("count activities missing workout for user %d: %w", userID, err)
|
|
}
|
|
return n, nil
|
|
}
|