Part of per-user profile isolation: profile rows are no longer a global singleton, so every read/write requires the caller's userID. This also requires scoping ListActivities, ListWorkoutKinds, UpsertActivity, CreateWorkoutKind, GetSyncState, UpdateSyncState, and ResetAllSyncedData to userID, plus updating all related tests in the store package.
278 lines
10 KiB
Go
278 lines
10 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 or updates the existing row for the
|
|
// same garmin_activity_id (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(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 garmin_activity_id = ? AND user_id = ?`, a.GarminActivityID, userID).Scan(&id); err != nil {
|
|
return 0, fmt.Errorf("fetch id for activity %d: %w", a.GarminActivityID, 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.
|
|
func (db *DB) GetActivity(ctx context.Context, id int64) (Activity, bool, error) {
|
|
row := db.QueryRowContext(ctx, `SELECT `+activityColumns+` FROM activities WHERE id = ?`, id)
|
|
a, err := scanActivity(row)
|
|
if err == sql.ErrNoRows {
|
|
return Activity{}, false, nil
|
|
}
|
|
if err != nil {
|
|
return Activity{}, false, fmt.Errorf("get activity %d: %w", id, 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 activities for userID 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 an activity with this garmin_activity_id is
|
|
// already stored, checked before each upsert during a sync pass so the
|
|
// reported "activities fetched" count reflects genuinely new activities, not
|
|
// every activity Garmin's API happens to return for the queried date range
|
|
// (which, thanks to the incremental overlap window and backfill's
|
|
// already-covered history, is almost always a re-listing of known ones).
|
|
func (db *DB) ActivityExists(ctx context.Context, garminActivityID int64) (bool, error) {
|
|
var id int64
|
|
err := db.QueryRowContext(ctx, `SELECT id FROM activities WHERE garmin_activity_id = ?`, garminActivityID).Scan(&id)
|
|
if err == sql.ErrNoRows {
|
|
return false, nil
|
|
}
|
|
if err != nil {
|
|
return false, fmt.Errorf("check activity %d exists: %w", garminActivityID, err)
|
|
}
|
|
return true, nil
|
|
}
|
|
|
|
// LatestActivityStartTime returns the start_time_utc of the most recently
|
|
// started activity we have, used to compute the incremental sync window.
|
|
func (db *DB) LatestActivityStartTime(ctx context.Context) (string, bool, error) {
|
|
var t string
|
|
err := db.QueryRowContext(ctx, `SELECT start_time_utc FROM activities ORDER BY start_time_utc DESC LIMIT 1`).Scan(&t)
|
|
if err == sql.ErrNoRows {
|
|
return "", false, nil
|
|
}
|
|
if err != nil {
|
|
return "", false, fmt.Errorf("latest activity start time: %w", 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, 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 = ?`, rawJSON, activityID)
|
|
if err != nil {
|
|
return fmt.Errorf("set activity %d details: %w", activityID, 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, activityID int64, rawJSON string) error {
|
|
_, err := db.ExecContext(ctx, `
|
|
UPDATE activities SET workout_raw_json = ?, updated_at = datetime('now')
|
|
WHERE id = ?`, rawJSON, activityID)
|
|
if err != nil {
|
|
return fmt.Errorf("set activity %d workout: %w", activityID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// SetActivitySplitsFetched records that get_activity_splits has been fetched
|
|
// for this activity.
|
|
func (db *DB) SetActivitySplitsFetched(ctx context.Context, activityID int64) error {
|
|
_, err := db.ExecContext(ctx, `
|
|
UPDATE activities SET splits_fetched_at = datetime('now'), updated_at = datetime('now')
|
|
WHERE id = ?`, activityID)
|
|
if err != nil {
|
|
return fmt.Errorf("set activity %d splits fetched: %w", activityID, err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ActivitiesMissingDetails returns 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, limit int) ([]Activity, error) {
|
|
rows, err := db.QueryContext(ctx, `SELECT `+activityColumns+` FROM activities
|
|
WHERE details_fetched_at IS NULL OR splits_fetched_at IS NULL
|
|
ORDER BY start_time_utc DESC LIMIT ?`, limit)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("list activities missing details: %w", 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 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) (int, error) {
|
|
var n int
|
|
err := db.QueryRowContext(ctx, `SELECT COUNT(*) FROM activities
|
|
WHERE details_fetched_at IS NULL OR splits_fetched_at IS NULL`).Scan(&n)
|
|
if err != nil {
|
|
return 0, fmt.Errorf("count activities missing details: %w", err)
|
|
}
|
|
return n, nil
|
|
}
|