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.
This commit is contained in:
2026-04-25 21:16:15 -04:00
parent e87f481988
commit f699899408
2 changed files with 329 additions and 0 deletions
+289
View File
@@ -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
}