garmin_activity_id uniqueness becomes per-user (UNIQUE(user_id, garmin_activity_id), added in Task 1) so two users' Garmin accounts can never collide even in the unlikely event their activity ids coincide.
275 lines
11 KiB
Go
275 lines
11 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
|
|
}
|