From f699899408a8ff2b89a6191e7972c0e8bc1bbd2f Mon Sep 17 00:00:00 2001 From: John O'Keefe Date: Sat, 25 Apr 2026 21:16:15 -0400 Subject: [PATCH] feat(sync): add ProgressService with merge, enrichment, and conflict detection Introduces a centralized ProgressService that handles all reading progress writes across web, KOReader, and Kobo clients. The service implements: - Merge strategy: reads existing progress first, then only overwrites non-nil fields from the incoming request. This fixes the data loss bug where partial updates (e.g., Kobo last-read-place sending only epubcfi and chapter) would NULL out percentage, character_offset, etc. - Enrichment: computes missing fields from available data: - character_offset from percentage + total_characters - current_page from percentage + total_pages - percentage from current_page + total_pages (reverse) - percentage from character_offset + total_characters (reverse) - Conflict detection: when a different source writes progress within 5 minutes with >1% difference, records a sync_conflicts row and broadcasts a WebSocket notification for real-time UI alerts. - Broadcast control: SaveProgressRequest.Broadcast flag lets Kobo last-read-place and SyncFromServer skip WebSocket broadcasts. - Pointer fields on SaveProgressRequest: nil means preserve existing, non-nil means overwrite. Eliminates ambiguity between zero values and not-provided fields. Also adds unit tests for buildProgressSnapshot helper function. --- internal/sync/progress.go | 289 +++++++++++++++++++++++++++++++++ internal/sync/progress_test.go | 40 +++++ 2 files changed, 329 insertions(+) diff --git a/internal/sync/progress.go b/internal/sync/progress.go index 63730e9..5631dfa 100644 --- a/internal/sync/progress.go +++ b/internal/sync/progress.go @@ -1,9 +1,17 @@ package sync import ( + "bookhoard/internal/database" + "context" "encoding/json" "fmt" + "log" "math" + "time" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgtype" ) // ProgressData represents universal progress data with multiple location references @@ -262,3 +270,284 @@ func EstimatedPages(totalCharacters int64) int { } return int(totalCharacters / CharsPerPage) } + +type ProgressService struct { + db *database.Queries + connManager *ConnectionManager +} + +func NewProgressService(db *database.Queries, connManager *ConnectionManager) *ProgressService { + return &ProgressService{db: db, connManager: connManager} +} + +type SaveProgressRequest struct { + MediaItemID pgtype.UUID + UserID pgtype.UUID + Source string + DeviceID pgtype.UUID + + Percentage *float64 + Epubcfi *string + CharacterOffset *int64 + Chapter *int + ChapterProgress *float64 + CurrentPage *int + TotalPages *int + ViewportX *float64 + ViewportY *float64 + ZoomLevel *float64 + ScrollX *float64 + ScrollY *float64 + PanelNumber *int + ReadingMode *string + + DeviceType string + DeviceName string + Broadcast bool +} + +func (s *ProgressService) SaveProgress(ctx context.Context, req SaveProgressRequest) (database.ReadingProgress, error) { + if req.Source == "" { + req.Source = "unknown" + } + + mediaItem, err := s.db.GetMediaItem(ctx, req.MediaItemID) + if err != nil { + return database.ReadingProgress{}, fmt.Errorf("media item not found: %w", err) + } + + existing, err := s.db.GetReadingProgress(ctx, database.GetReadingProgressParams{ + MediaItemID: req.MediaItemID, + UserID: req.UserID, + }) + hasExisting := err == nil + if err != nil && err != pgx.ErrNoRows { + return database.ReadingProgress{}, fmt.Errorf("failed to read existing progress: %w", err) + } + + params := database.UpdateUniversalProgressParams{ + MediaItemID: req.MediaItemID, + UserID: req.UserID, + } + + if hasExisting { + params.Percentage = existing.Percentage + params.CharacterOffset = existing.CharacterOffset + params.Epubcfi = existing.Epubcfi + params.Chapter = existing.Chapter + params.ChapterProgress = existing.ChapterProgress + params.ViewportX = existing.ViewportX + params.ViewportY = existing.ViewportY + params.ZoomLevel = existing.ZoomLevel + params.ScrollPositionX = existing.ScrollPositionX + params.ScrollPositionY = existing.ScrollPositionY + params.PanelNumber = existing.PanelNumber + params.ReadingMode = existing.ReadingMode + params.CurrentPage = existing.CurrentPage + params.TotalPages = existing.TotalPages + params.LastSyncDevice = existing.LastSyncDevice + params.LastSyncSource = existing.LastSyncSource + } + + if req.Percentage != nil { + params.Percentage = pgtype.Float8{Float64: *req.Percentage, Valid: true} + } + if req.Epubcfi != nil { + params.Epubcfi = pgtype.Text{String: *req.Epubcfi, Valid: *req.Epubcfi != ""} + } + if req.CharacterOffset != nil { + params.CharacterOffset = pgtype.Int8{Int64: *req.CharacterOffset, Valid: true} + } + if req.Chapter != nil { + params.Chapter = pgtype.Int4{Int32: int32(*req.Chapter), Valid: true} + } + if req.ChapterProgress != nil { + params.ChapterProgress = pgtype.Float8{Float64: *req.ChapterProgress, Valid: true} + } + if req.CurrentPage != nil { + params.CurrentPage = pgtype.Int4{Int32: int32(*req.CurrentPage), Valid: true} + } + if req.TotalPages != nil { + params.TotalPages = pgtype.Int4{Int32: int32(*req.TotalPages), Valid: true} + } + if req.ViewportX != nil { + params.ViewportX = pgtype.Float8{Float64: *req.ViewportX, Valid: true} + } + if req.ViewportY != nil { + params.ViewportY = pgtype.Float8{Float64: *req.ViewportY, Valid: true} + } + if req.ZoomLevel != nil { + params.ZoomLevel = pgtype.Float8{Float64: *req.ZoomLevel, Valid: true} + } + if req.ScrollX != nil { + params.ScrollPositionX = pgtype.Float8{Float64: *req.ScrollX, Valid: true} + } + if req.ScrollY != nil { + params.ScrollPositionY = pgtype.Float8{Float64: *req.ScrollY, Valid: true} + } + if req.PanelNumber != nil { + params.PanelNumber = pgtype.Int4{Int32: int32(*req.PanelNumber), Valid: true} + } + if req.ReadingMode != nil { + params.ReadingMode = pgtype.Text{String: *req.ReadingMode, Valid: true} + } + + params.LastSyncDevice = pgtype.Text{String: req.Source, Valid: true} + params.LastSyncSource = pgtype.Text{String: req.Source, Valid: true} + + if params.Percentage.Valid && !params.CharacterOffset.Valid && mediaItem.TotalCharacters.Valid && mediaItem.TotalCharacters.Int64 > 0 { + charOff := PercentageToCharacter(params.Percentage.Float64, mediaItem.TotalCharacters.Int64) + params.CharacterOffset = pgtype.Int8{Int64: charOff, Valid: true} + } + if params.Percentage.Valid && !params.CurrentPage.Valid && params.TotalPages.Valid && params.TotalPages.Int32 > 0 { + page := PercentageToPage(params.Percentage.Float64, int(params.TotalPages.Int32)) + params.CurrentPage = pgtype.Int4{Int32: int32(page), Valid: true} + } + if params.CurrentPage.Valid && params.TotalPages.Valid && params.TotalPages.Int32 > 0 && !params.Percentage.Valid { + pct := PageToPercentage(int(params.CurrentPage.Int32), int(params.TotalPages.Int32)) + params.Percentage = pgtype.Float8{Float64: pct, Valid: true} + } + if params.CharacterOffset.Valid && mediaItem.TotalCharacters.Valid && mediaItem.TotalCharacters.Int64 > 0 && !params.Percentage.Valid { + pct := CharacterToPercentage(params.CharacterOffset.Int64, mediaItem.TotalCharacters.Int64) + params.Percentage = pgtype.Float8{Float64: pct, Valid: true} + } + + conflictDetected := false + if hasExisting && existing.LastSyncSource.Valid && existing.LastSyncSource.String != req.Source { + if existing.LastSyncTimestamp.Valid { + timeDiff := time.Since(existing.LastSyncTimestamp.Time) + if timeDiff < 5*time.Minute { + existingPct := 0.0 + if existing.Percentage.Valid { + existingPct = existing.Percentage.Float64 + } + newPct := 0.0 + if params.Percentage.Valid { + newPct = params.Percentage.Float64 + } + diff := newPct - existingPct + if diff < 0 { + diff = -diff + } + if diff > 0.01 { + conflictDetected = true + } + } + } + } + + result, err := s.db.UpdateUniversalProgress(ctx, params) + if err != nil { + return database.ReadingProgress{}, fmt.Errorf("failed to upsert progress: %w", err) + } + + if conflictDetected { + newData := map[string]interface{}{ + "source": req.Source, + "timestamp": time.Now().Format(time.RFC3339), + "data": buildProgressSnapshot(req), + } + existingData := map[string]interface{}{ + "source": existing.LastSyncSource.String, + "timestamp": existing.LastSyncTimestamp.Time.Format(time.RFC3339), + "data": map[string]interface{}{ + "percentage": float64Ptr(existing.Percentage), + "epubcfi": textPtr(existing.Epubcfi), + "chapter": int32Ptr(existing.Chapter), + "character": int64Ptr(existing.CharacterOffset), + "page": int32Ptr(existing.CurrentPage), + "total_pages": int32Ptr(existing.TotalPages), + }, + } + conflictData := map[string]interface{}{ + "new": newData, + "existing": existingData, + } + conflictJSON, _ := json.Marshal(conflictData) + _, err := s.db.CreateSyncConflict(ctx, database.CreateSyncConflictParams{ + MediaItemID: req.MediaItemID, + UserID: req.UserID, + ConflictType: "progress", + ConflictData: conflictJSON, + }) + if err != nil { + log.Printf("Failed to record sync conflict: %v", err) + } else if s.connManager != nil { + s.connManager.BroadcastConflictNotification(req.MediaItemID.Bytes, "detection", "") + } + } + + if s.connManager != nil && req.Broadcast { + deviceName := req.DeviceName + if deviceName == "" { + deviceName = req.Source + " Device" + } + pct := 0.0 + if result.Percentage.Valid { + pct = result.Percentage.Float64 + } + s.connManager.BroadcastProgressUpdate( + uuid.UUID(req.MediaItemID.Bytes), + pct, + SourceDevice{ + ID: uuid.UUID(req.DeviceID.Bytes).String(), + Name: deviceName, + Type: req.Source, + }, + ) + } + + return result, nil +} + +func buildProgressSnapshot(req SaveProgressRequest) map[string]interface{} { + data := map[string]interface{}{} + if req.Percentage != nil { + data["percentage"] = *req.Percentage + } + if req.Epubcfi != nil { + data["epubcfi"] = *req.Epubcfi + } + if req.Chapter != nil { + data["chapter"] = *req.Chapter + } + if req.CharacterOffset != nil { + data["character"] = *req.CharacterOffset + } + if req.CurrentPage != nil { + data["page"] = *req.CurrentPage + } + if req.TotalPages != nil { + data["total_pages"] = *req.TotalPages + } + return data +} + +func float64Ptr(v pgtype.Float8) *float64 { + if v.Valid { + return &v.Float64 + } + return nil +} + +func textPtr(v pgtype.Text) *string { + if v.Valid { + return &v.String + } + return nil +} + +func int32Ptr(v pgtype.Int4) *int { + if v.Valid { + val := int(v.Int32) + return &val + } + return nil +} + +func int64Ptr(v pgtype.Int8) *int64 { + if v.Valid { + return &v.Int64 + } + return nil +} diff --git a/internal/sync/progress_test.go b/internal/sync/progress_test.go index 8d04489..1be0b2f 100644 --- a/internal/sync/progress_test.go +++ b/internal/sync/progress_test.go @@ -456,3 +456,43 @@ func TestEdgeCases(t *testing.T) { assert.InDelta(t, 100, page, 1) }) } + +func TestBuildProgressSnapshot(t *testing.T) { + t.Run("nil fields are omitted", func(t *testing.T) { + pct := 0.5 + ch := 3 + req := SaveProgressRequest{ + Percentage: &pct, + Chapter: &ch, + } + snapshot := buildProgressSnapshot(req) + assert.InDelta(t, 0.5, snapshot["percentage"], 0.001) + assert.Equal(t, 3, snapshot["chapter"]) + _, hasEpubcfi := snapshot["epubcfi"] + assert.False(t, hasEpubcfi) + }) + + t.Run("all fields present", func(t *testing.T) { + pct := 0.75 + cfi := "epubcfi(/6/4/2:10)" + ch := 5 + charOff := int64(10000) + page := 150 + tp := 200 + req := SaveProgressRequest{ + Percentage: &pct, + Epubcfi: &cfi, + Chapter: &ch, + CharacterOffset: &charOff, + CurrentPage: &page, + TotalPages: &tp, + } + snapshot := buildProgressSnapshot(req) + assert.InDelta(t, 0.75, snapshot["percentage"], 0.001) + assert.Equal(t, "epubcfi(/6/4/2:10)", snapshot["epubcfi"]) + assert.Equal(t, 5, snapshot["chapter"]) + assert.Equal(t, int64(10000), snapshot["character"]) + assert.Equal(t, 150, snapshot["page"]) + assert.Equal(t, 200, snapshot["total_pages"]) + }) +}