Files
bookhoard/internal/sync/queue.go
john-okeefe 1cda4e5191 feat(handlers): integrate ProgressService into media, koreader, kobo, and queue
All four progress write paths now delegate to ProgressService.SaveProgress:

- MediaHandler: UpdateMediaReadingProgress uses ProgressService for web
  saves with richer request body (reading_mode, zoom_level, scroll). GET
  now uses GetUniversalProgress query that JOINs media_items for
  format_group, total_characters, chapter_count.

- KOReaderHandler: updateProgressForBook delegates to ProgressService.
  Fixed device ID bug (was using userID, now uses deviceID). Removed
  duplicate UpdateDeviceLastSync with zero UUID. Added pgtype helper
  functions (textPtrToPgText, intPtrToPgInt4, int64PtrToPgInt8).

- KoboHandler: all four progress write points (Markup ReadingSync, Markup
  last-read-place, AnalyticsGettests, SyncFromServer) delegate to
  ProgressService. Fixed empty epubcfi string now correctly set to
  Valid: false. SyncFromServer preserves last_sync_source=bookhoard
  and Broadcast: false.

- QueueProcessor: syncProgress delegates to ProgressService.

- main.go: creates ProgressService after ConnectionManager, injects via
  SetProgressService() on all handlers and queue processor.

Handler tests cover pgtype conversion helpers (textPtrToPgText, etc.)
and device icon mapping.
2026-04-25 21:16:29 -04:00

468 lines
12 KiB
Go

package sync
import (
"bookhoard/internal/database"
"context"
"encoding/json"
"fmt"
"log"
"time"
"github.com/jackc/pgx/v5/pgtype"
)
const (
SyncTypeProgress = "progress"
SyncTypeNote = "note"
SyncTypeHighlight = "highlight"
SyncTypeBookmark = "bookmark"
)
const (
SyncStatusPending = "pending"
SyncStatusProcessing = "processing"
SyncStatusCompleted = "completed"
SyncStatusFailed = "failed"
)
const (
PriorityUserInitiated = 1
PriorityBookCompletion = 2
PriorityCriticalNote = 3
PriorityPageTurn = 5
PriorityCheckpoint = 7
PriorityBackgroundSync = 10
)
type SyncQueueProcessor struct {
db *database.Queries
progressSvc *ProgressService
progressChan chan *ProgressUpdate
interval time.Duration
batchSize int
}
type ProgressUpdate struct {
DeviceID pgtype.UUID
MediaItemID pgtype.UUID
UserID pgtype.UUID
Percentage float64
Epubcfi *string
Chapter *int
Character *int64
Page *int
TotalPages *int
Source string
SyncMode string
}
type SyncQueueItem struct {
ID pgtype.UUID
DeviceID pgtype.UUID
MediaItemID pgtype.UUID
SyncType string
SyncData []byte
Priority int32
Attempts int32
MaxAttempts int32
Status string
ErrorMessage *string
CreatedAt time.Time
ProcessedAt *time.Time
}
func NewSyncQueueProcessor(db *database.Queries) *SyncQueueProcessor {
return &SyncQueueProcessor{
db: db,
progressChan: make(chan *ProgressUpdate, 100),
interval: 5 * time.Second,
batchSize: 50,
}
}
func (p *SyncQueueProcessor) SetProgressService(svc *ProgressService) {
p.progressSvc = svc
}
func (p *SyncQueueProcessor) Start(ctx context.Context) {
log.Printf("Starting sync queue processor (interval: %v, batch: %d)", p.interval, p.batchSize)
ticker := time.NewTicker(p.interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
log.Println("Sync queue processor stopped")
return
case update := <-p.progressChan:
p.enqueueProgressUpdate(ctx, update)
case <-ticker.C:
p.processQueue(ctx)
}
}
}
func (p *SyncQueueProcessor) EnqueueProgress(update *ProgressUpdate) error {
select {
case p.progressChan <- update:
return nil
default:
return fmt.Errorf("progress channel is full")
}
}
func (p *SyncQueueProcessor) enqueueProgressUpdate(ctx context.Context, update *ProgressUpdate) {
syncData := map[string]interface{}{
"percentage": update.Percentage,
"source": update.Source,
"timestamp": time.Now().Format(time.RFC3339),
}
if update.Epubcfi != nil {
syncData["epubcfi"] = *update.Epubcfi
}
if update.Chapter != nil {
syncData["chapter"] = *update.Chapter
}
if update.Character != nil {
syncData["character"] = *update.Character
}
if update.Page != nil {
syncData["page"] = *update.Page
}
if update.TotalPages != nil {
syncData["total_pages"] = *update.TotalPages
}
syncDataJSON, err := json.Marshal(syncData)
if err != nil {
log.Printf("Failed to marshal sync data: %v", err)
return
}
priority := PriorityPageTurn
if update.SyncMode == "checkpoint" {
priority = PriorityCheckpoint
} else if update.Percentage >= 0.99 {
priority = PriorityBookCompletion
}
maxAttempts := int32(3)
if update.SyncMode == "immediate" {
maxAttempts = 5
}
_, err = p.db.CreateSyncQueueItem(ctx, database.CreateSyncQueueItemParams{
DeviceID: update.DeviceID,
MediaItemID: update.MediaItemID,
SyncType: SyncTypeProgress,
SyncData: syncDataJSON,
Priority: pgtype.Int4{Int32: int32(priority), Valid: true},
MaxAttempts: pgtype.Int4{Int32: maxAttempts, Valid: true},
Status: pgtype.Text{String: SyncStatusPending, Valid: true},
})
if err != nil {
log.Printf("Failed to create sync queue item: %v", err)
}
}
func (p *SyncQueueProcessor) processQueue(ctx context.Context) {
items := p.getPendingItems(ctx)
if len(items) == 0 {
p.checkStuckItems(ctx)
return
}
log.Printf("Processing %d sync queue items", len(items))
for _, item := range items {
select {
case <-ctx.Done():
return
default:
p.processItem(ctx, item)
}
}
}
func (p *SyncQueueProcessor) getPendingItems(ctx context.Context) []SyncQueueItem {
var allItems []SyncQueueItem
var offset int32 = 0
batchSize := int32(p.batchSize)
for {
dbItems, err := p.db.ListPendingSyncQueueItems(ctx, database.ListPendingSyncQueueItemsParams{
DeviceID: pgtype.UUID{Valid: false},
Limit: batchSize,
})
if err != nil {
log.Printf("Error listing pending items: %v", err)
break
}
if len(dbItems) == 0 {
break
}
for _, item := range dbItems {
if item.Status.String == SyncStatusPending {
allItems = append(allItems, dbItemToSyncQueueItem(item))
}
}
if int32(len(dbItems)) < batchSize {
break
}
offset += batchSize
}
return allItems
}
func (p *SyncQueueProcessor) checkStuckItems(ctx context.Context) {
items, err := p.db.GetStuckSyncQueueItems(ctx)
if err != nil {
return
}
for _, item := range items {
if item.Attempts.Int32 < item.MaxAttempts.Int32 {
log.Printf("Resetting stuck queue item %s", item.ID.Bytes)
_, err := p.db.UpdateSyncQueueItemStatus(ctx, database.UpdateSyncQueueItemStatusParams{
ID: item.ID,
Status: pgtype.Text{String: SyncStatusPending, Valid: true},
ErrorMessage: pgtype.Text{String: "Item was stuck in processing state", Valid: true},
})
if err != nil {
log.Printf("Failed to reset stuck item: %v", err)
}
}
}
}
func (p *SyncQueueProcessor) processItem(ctx context.Context, item SyncQueueItem) {
_, err := p.db.UpdateSyncQueueItemStatus(ctx, database.UpdateSyncQueueItemStatusParams{
ID: item.ID,
Status: pgtype.Text{String: SyncStatusProcessing, Valid: true},
ErrorMessage: pgtype.Text{},
})
if err != nil {
log.Printf("Failed to update item status to processing: %v", err)
return
}
var syncData map[string]interface{}
err = json.Unmarshal(item.SyncData, &syncData)
if err != nil {
log.Printf("Failed to unmarshal sync data: %v", err)
p.markItemFailed(ctx, item, fmt.Sprintf("Invalid sync data: %v", err))
return
}
err = p.executeSync(ctx, item, syncData)
if err != nil {
log.Printf("Sync failed for item %s: %v", item.ID.Bytes, err)
if item.Attempts >= item.MaxAttempts {
p.markItemFailed(ctx, item, fmt.Sprintf("Max attempts reached: %v", err))
} else {
nextRetry := p.calculateNextRetry(item.Attempts)
log.Printf("Scheduling retry for item %s at %v (attempt %d/%d)",
item.ID.Bytes, nextRetry, item.Attempts+1, item.MaxAttempts)
p.markItemPending(ctx, item, fmt.Sprintf("Will retry: %v", err))
}
return
}
_, err = p.db.UpdateSyncQueueItemStatus(ctx, database.UpdateSyncQueueItemStatusParams{
ID: item.ID,
Status: pgtype.Text{String: SyncStatusCompleted, Valid: true},
ErrorMessage: pgtype.Text{},
})
if err != nil {
log.Printf("Failed to mark item as completed: %v", err)
}
log.Printf("Successfully processed sync queue item %s", item.ID.Bytes)
}
func (p *SyncQueueProcessor) executeSync(ctx context.Context, item SyncQueueItem, syncData map[string]interface{}) error {
device, err := p.db.GetDevice(ctx, item.DeviceID)
if err != nil {
return fmt.Errorf("device not found: %w", err)
}
switch item.SyncType {
case SyncTypeProgress:
return p.syncProgress(ctx, device.UserID, item.MediaItemID, syncData)
case SyncTypeNote:
return p.syncNote(ctx, device.UserID, item.MediaItemID, syncData)
case SyncTypeHighlight:
return p.syncHighlight(ctx, device.UserID, item.MediaItemID, syncData)
default:
return fmt.Errorf("unsupported sync type: %s", item.SyncType)
}
}
func (p *SyncQueueProcessor) syncProgress(ctx context.Context, userID pgtype.UUID, mediaItemID pgtype.UUID, syncData map[string]interface{}) error {
percentage, ok := syncData["percentage"].(float64)
if !ok {
return fmt.Errorf("missing percentage in sync data")
}
source := "queue"
if v, ok := syncData["source"].(string); ok {
source = v
}
if p.progressSvc != nil {
req := SaveProgressRequest{
MediaItemID: mediaItemID,
UserID: userID,
Source: source,
Percentage: &percentage,
Broadcast: false,
}
if v, ok := syncData["epubcfi"].(string); ok {
req.Epubcfi = &v
}
if v, ok := syncData["chapter"].(float64); ok {
ch := int(v)
req.Chapter = &ch
}
if v, ok := syncData["character"].(float64); ok {
co := int64(v)
req.CharacterOffset = &co
}
if v, ok := syncData["page"].(float64); ok {
pg := int(v)
req.CurrentPage = &pg
}
if v, ok := syncData["total_pages"].(float64); ok {
tp := int(v)
req.TotalPages = &tp
}
_, err := p.progressSvc.SaveProgress(ctx, req)
return err
}
var epubcfi pgtype.Text
if v, ok := syncData["epubcfi"].(string); ok {
epubcfi = pgtype.Text{String: v, Valid: true}
}
var chapter pgtype.Int4
if v, ok := syncData["chapter"].(float64); ok {
chapter = pgtype.Int4{Int32: int32(v), Valid: true}
}
var characterOffset pgtype.Int8
if v, ok := syncData["character"].(float64); ok {
characterOffset = pgtype.Int8{Int64: int64(v), Valid: true}
}
var currentPage pgtype.Int4
if v, ok := syncData["page"].(float64); ok {
currentPage = pgtype.Int4{Int32: int32(v), Valid: true}
}
var totalPages pgtype.Int4
if v, ok := syncData["total_pages"].(float64); ok {
totalPages = pgtype.Int4{Int32: int32(v), Valid: true}
}
_, err := p.db.UpdateUniversalProgress(ctx, database.UpdateUniversalProgressParams{
MediaItemID: mediaItemID,
UserID: userID,
Percentage: pgtype.Float8{Float64: percentage, Valid: true},
Epubcfi: epubcfi,
Chapter: chapter,
ChapterProgress: pgtype.Float8{Float64: percentage, Valid: true},
CharacterOffset: characterOffset,
CurrentPage: currentPage,
TotalPages: totalPages,
LastSyncDevice: pgtype.Text{String: source, Valid: true},
LastSyncSource: pgtype.Text{String: source, Valid: true},
ViewportX: pgtype.Float8{},
ScrollPositionX: pgtype.Float8{},
ScrollPositionY: pgtype.Float8{},
PanelNumber: pgtype.Int4{},
ReadingMode: pgtype.Text{},
ZoomLevel: pgtype.Float8{},
})
return err
}
func (p *SyncQueueProcessor) syncNote(ctx context.Context, userID pgtype.UUID, mediaItemID pgtype.UUID, syncData map[string]interface{}) error {
return fmt.Errorf("note sync not yet implemented")
}
func (p *SyncQueueProcessor) syncHighlight(ctx context.Context, userID pgtype.UUID, mediaItemID pgtype.UUID, syncData map[string]interface{}) error {
return fmt.Errorf("highlight sync not yet implemented")
}
func (p *SyncQueueProcessor) markItemFailed(ctx context.Context, item SyncQueueItem, errMsg string) {
_, err := p.db.UpdateSyncQueueItemStatus(ctx, database.UpdateSyncQueueItemStatusParams{
ID: item.ID,
Status: pgtype.Text{String: SyncStatusFailed, Valid: true},
ErrorMessage: pgtype.Text{String: errMsg, Valid: true},
})
if err != nil {
log.Printf("Failed to mark item as failed: %v", err)
}
}
func (p *SyncQueueProcessor) markItemPending(ctx context.Context, item SyncQueueItem, errMsg string) {
_, err := p.db.UpdateSyncQueueItemStatus(ctx, database.UpdateSyncQueueItemStatusParams{
ID: item.ID,
Status: pgtype.Text{String: SyncStatusPending, Valid: true},
ErrorMessage: pgtype.Text{String: errMsg, Valid: true},
})
if err != nil {
log.Printf("Failed to mark item as pending: %v", err)
}
}
func (p *SyncQueueProcessor) calculateNextRetry(attempts int32) time.Time {
switch attempts {
case 0:
return time.Now()
case 1:
return time.Now().Add(1 * time.Minute)
case 2:
return time.Now().Add(5 * time.Minute)
case 3:
return time.Now().Add(15 * time.Minute)
case 4:
return time.Now().Add(1 * time.Hour)
default:
return time.Now().Add(24 * time.Hour)
}
}
func (p *SyncQueueProcessor) GetQueueStats(ctx context.Context, deviceID pgtype.UUID) (database.GetSyncQueueStatsRow, error) {
return p.db.GetSyncQueueStats(ctx, deviceID)
}
func dbItemToSyncQueueItem(item database.SyncQueue) SyncQueueItem {
return SyncQueueItem{
ID: item.ID,
DeviceID: item.DeviceID,
MediaItemID: item.MediaItemID,
SyncType: item.SyncType,
Priority: item.Priority.Int32,
Attempts: item.Attempts.Int32,
MaxAttempts: item.MaxAttempts.Int32,
Status: item.Status.String,
ErrorMessage: &item.ErrorMessage.String,
CreatedAt: item.CreatedAt.Time,
ProcessedAt: &item.ProcessedAt.Time,
}
}