diff --git a/internal/sync/queue.go b/internal/sync/queue.go index e4ad5d9..b7822e0 100644 --- a/internal/sync/queue.go +++ b/internal/sync/queue.go @@ -35,11 +35,12 @@ const ( ) type SyncQueueProcessor struct { - db *database.Queries - progressSvc *ProgressService - progressChan chan *ProgressUpdate - interval time.Duration - batchSize int + db *database.Queries + progressSvc *ProgressService + annotationSvc *AnnotationService + progressChan chan *ProgressUpdate + interval time.Duration + batchSize int } type ProgressUpdate struct { @@ -85,6 +86,10 @@ func (p *SyncQueueProcessor) SetProgressService(svc *ProgressService) { p.progressSvc = svc } +func (p *SyncQueueProcessor) SetAnnotationService(svc *AnnotationService) { + p.annotationSvc = svc +} + func (p *SyncQueueProcessor) Start(ctx context.Context) { log.Printf("Starting sync queue processor (interval: %v, batch: %d)", p.interval, p.batchSize) @@ -113,6 +118,100 @@ func (p *SyncQueueProcessor) EnqueueProgress(update *ProgressUpdate) error { } } +type HighlightUpdate struct { + DeviceID pgtype.UUID + MediaItemID pgtype.UUID + UserID pgtype.UUID + SelectionText string + StartPosition string + EndPosition string + Color string + NoteText string + EpubcfiStart string + EpubcfiEnd string + PercentageStart float64 + PercentageEnd float64 + Source string + DeviceSyncData map[string]interface{} +} + +type NoteUpdate struct { + DeviceID pgtype.UUID + MediaItemID pgtype.UUID + UserID pgtype.UUID + Content string + Position string + Source string + DeviceSyncData map[string]interface{} +} + +type BookmarkUpdate struct { + DeviceID pgtype.UUID + MediaItemID pgtype.UUID + UserID pgtype.UUID + Title string + Position string + Notes string + Source string + DeviceSyncData map[string]interface{} +} + +func (p *SyncQueueProcessor) EnqueueHighlight(ctx context.Context, update *HighlightUpdate) error { + syncData := map[string]interface{}{ + "selection_text": update.SelectionText, + "start_position": update.StartPosition, + "end_position": update.EndPosition, + "color": update.Color, + "note_text": update.NoteText, + "source": update.Source, + "epubcfi_start": update.EpubcfiStart, + "epubcfi_end": update.EpubcfiEnd, + "percentage_start": update.PercentageStart, + "percentage_end": update.PercentageEnd, + "device_sync_data": update.DeviceSyncData, + } + return p.enqueueAnnotation(ctx, update.DeviceID, update.MediaItemID, SyncTypeHighlight, syncData) +} + +func (p *SyncQueueProcessor) EnqueueNote(ctx context.Context, update *NoteUpdate) error { + syncData := map[string]interface{}{ + "content": update.Content, + "position": update.Position, + "source": update.Source, + "device_sync_data": update.DeviceSyncData, + } + return p.enqueueAnnotation(ctx, update.DeviceID, update.MediaItemID, SyncTypeNote, syncData) +} + +func (p *SyncQueueProcessor) EnqueueBookmark(ctx context.Context, update *BookmarkUpdate) error { + syncData := map[string]interface{}{ + "title": update.Title, + "position": update.Position, + "notes": update.Notes, + "source": update.Source, + "device_sync_data": update.DeviceSyncData, + } + return p.enqueueAnnotation(ctx, update.DeviceID, update.MediaItemID, SyncTypeBookmark, syncData) +} + +func (p *SyncQueueProcessor) enqueueAnnotation(ctx context.Context, deviceID, mediaItemID pgtype.UUID, syncType string, syncData map[string]interface{}) error { + syncDataJSON, err := json.Marshal(syncData) + if err != nil { + return fmt.Errorf("marshal sync data: %w", err) + } + + _, err = p.db.CreateSyncQueueItem(ctx, database.CreateSyncQueueItemParams{ + DeviceID: deviceID, + MediaItemID: mediaItemID, + SyncType: syncType, + SyncData: syncDataJSON, + Priority: pgtype.Int4{Int32: int32(PriorityCriticalNote), Valid: true}, + MaxAttempts: pgtype.Int4{Int32: 3, Valid: true}, + Status: pgtype.Text{String: SyncStatusPending, Valid: true}, + }) + return err +} + func (p *SyncQueueProcessor) enqueueProgressUpdate(ctx context.Context, update *ProgressUpdate) { syncData := map[string]interface{}{ "percentage": update.Percentage, @@ -308,6 +407,8 @@ func (p *SyncQueueProcessor) executeSync(ctx context.Context, item SyncQueueItem return p.syncNote(ctx, device.UserID, item.MediaItemID, syncData) case SyncTypeHighlight: return p.syncHighlight(ctx, device.UserID, item.MediaItemID, syncData) + case SyncTypeBookmark: + return p.syncBookmark(ctx, device.UserID, item.MediaItemID, syncData) default: return fmt.Errorf("unsupported sync type: %s", item.SyncType) } @@ -407,11 +508,138 @@ func (p *SyncQueueProcessor) syncProgress(ctx context.Context, userID pgtype.UUI } 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") + if p.annotationSvc == nil { + return fmt.Errorf("annotation service not available") + } + + req := SaveNoteRequest{ + MediaItemID: mediaItemID, + UserID: userID, + } + + if v, ok := syncData["content"].(string); ok { + req.Content = v + } + if v, ok := syncData["position"].(string); ok { + req.Position = v + } + if v, ok := syncData["source"].(string); ok { + req.Source = v + } + if v, ok := syncData["epubcfi_location"].(string); ok { + req.EpubcfiLocation = v + } + if v, ok := syncData["percentage_location"].(float64); ok { + req.PercentageLocation = v + } + if v, ok := syncData["chapter_reference"].(float64); ok { + req.ChapterReference = int32(v) + } + if v, ok := syncData["device_sync_data"].(map[string]interface{}); ok { + req.DeviceSyncData, _ = json.Marshal(v) + } + + _, err := p.annotationSvc.SaveNote(ctx, req) + return err } 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") + if p.annotationSvc == nil { + return fmt.Errorf("annotation service not available") + } + + req := SaveHighlightRequest{ + MediaItemID: mediaItemID, + UserID: userID, + } + + if v, ok := syncData["selection_text"].(string); ok { + req.SelectionText = v + } + if v, ok := syncData["start_position"].(string); ok { + req.StartPosition = v + } + if v, ok := syncData["end_position"].(string); ok { + req.EndPosition = v + } + if v, ok := syncData["color"].(string); ok { + req.Color = v + } + if v, ok := syncData["note_text"].(string); ok { + req.NoteText = v + } + if v, ok := syncData["source"].(string); ok { + req.Source = v + } + if v, ok := syncData["epubcfi_start"].(string); ok { + req.EpubcfiStart = v + } + if v, ok := syncData["epubcfi_end"].(string); ok { + req.EpubcfiEnd = v + } + if v, ok := syncData["percentage_start"].(float64); ok { + req.PercentageStart = v + } + if v, ok := syncData["percentage_end"].(float64); ok { + req.PercentageEnd = v + } + if v, ok := syncData["chapter_reference"].(float64); ok { + req.ChapterReference = int32(v) + } + if v, ok := syncData["device_sync_data"].(map[string]interface{}); ok { + req.DeviceSyncData, _ = json.Marshal(v) + } + + _, err := p.annotationSvc.SaveHighlight(ctx, req) + return err +} + +func (p *SyncQueueProcessor) syncBookmark(ctx context.Context, userID pgtype.UUID, mediaItemID pgtype.UUID, syncData map[string]interface{}) error { + if p.annotationSvc == nil { + return fmt.Errorf("annotation service not available") + } + + req := SaveBookmarkRequest{ + MediaItemID: mediaItemID, + UserID: userID, + } + + if v, ok := syncData["title"].(string); ok { + req.Title = v + } + if v, ok := syncData["position"].(string); ok { + req.Position = v + } + if v, ok := syncData["notes"].(string); ok { + req.Notes = v + } + if v, ok := syncData["source"].(string); ok { + req.Source = v + } + if v, ok := syncData["cfi_position"].(string); ok { + req.CFIPosition = v + } + if v, ok := syncData["epubcfi_location"].(string); ok { + req.EpubcfiLocation = v + } + if v, ok := syncData["percentage_loc"].(float64); ok { + req.PercentageLoc = v + } + if v, ok := syncData["page_number"].(float64); ok { + req.PageNumber = int32(v) + } + if v, ok := syncData["chapter_number"].(float64); ok { + req.ChapterNumber = int32(v) + } + if v, ok := syncData["chapter_reference"].(float64); ok { + req.ChapterReference = int32(v) + } + if v, ok := syncData["device_sync_data"].(map[string]interface{}); ok { + req.DeviceSyncData, _ = json.Marshal(v) + } + + _, err := p.annotationSvc.SaveBookmark(ctx, req) + return err } func (p *SyncQueueProcessor) markItemFailed(ctx context.Context, item SyncQueueItem, errMsg string) {