From a7f2b83bddf9507ba46a58436c9c80ca4186da71 Mon Sep 17 00:00:00 2001 From: John O'Keefe Date: Sat, 31 Jan 2026 13:06:56 -0500 Subject: [PATCH] Integrate sync queue system and device cap API - Start queue processor as background goroutine - Initialize and register queue handler - Add queue management routes (7 endpoints) - Update KOReader handler to use checkpoint sync mode - Add device cap management route (PUT /api/auth/users/:id/max-devices) - Register all new endpoints with proper middleware --- cmd/server/main.go | 20 ++++++- internal/handlers/koreader.go | 100 +++++++++++++++++++++++++++++++++- 2 files changed, 116 insertions(+), 4 deletions(-) diff --git a/cmd/server/main.go b/cmd/server/main.go index 01a7f21..5c39eab 100644 --- a/cmd/server/main.go +++ b/cmd/server/main.go @@ -57,9 +57,14 @@ func main() { connManager := sync.NewConnectionManager() connManager.StartCleanupTask() - koreaderHandler := handlers.NewKOReaderHandler(queries, connManager) + // Create sync queue processor + queueProcessor := sync.NewSyncQueueProcessor(queries) + go queueProcessor.Start(context.Background()) + + koreaderHandler := handlers.NewKOReaderHandler(queries, connManager, queueProcessor) wsHandler := handlers.NewWSHandler(queries, connManager, cfg.JWTSecret, deviceAuthMiddleware) conflictHandler := handlers.NewConflictHandler(queries, connManager) + queueHandler := handlers.NewQueueHandler(queries, queueProcessor) e := echo.New() @@ -147,6 +152,7 @@ func main() { // Admin-only routes for user and folder management admin := protected.Group("/auth", handlers.AdminMiddleware) admin.GET("/users", authHandler.ListUsers) + admin.PUT("/users/:id/max-devices", authHandler.UpdateUserMaxDevices) // Library management routes library := protected.Group("/libraries") @@ -238,6 +244,18 @@ func main() { conflicts.DELETE("/:id", conflictHandler.DeleteConflict) conflicts.POST("/dismiss-all", conflictHandler.DismissAllResolved) + // Sync queue management routes (protected - require user auth) + queue := protected.Group("/queue") + queue.GET("/devices/:device_id/stats", queueHandler.GetDeviceQueueStats) + queue.GET("/devices/:device_id/items", queueHandler.ListDeviceQueueItems) + queue.POST("/items/:item_id/retry", queueHandler.RetryQueueItem) + queue.DELETE("/items/:item_id", queueHandler.DeleteQueueItem) + queue.DELETE("/devices/:device_id/clear", queueHandler.ClearDeviceQueue) + + // Admin-only queue routes + adminQueue := queue.Group("", handlers.AdminMiddleware) + adminQueue.GET("/items", queueHandler.ListAllQueueItems) + // WebSocket endpoint for real-time sync e.GET("/ws/sync", wsHandler.HandleWebSocket) diff --git a/internal/handlers/koreader.go b/internal/handlers/koreader.go index 13f14da..feb1410 100644 --- a/internal/handlers/koreader.go +++ b/internal/handlers/koreader.go @@ -17,10 +17,11 @@ import ( type KOReaderHandler struct { db *database.Queries connManager *wsync.ConnectionManager + queue *wsync.SyncQueueProcessor } -func NewKOReaderHandler(db *database.Queries, connManager *wsync.ConnectionManager) *KOReaderHandler { - return &KOReaderHandler{db: db, connManager: connManager} +func NewKOReaderHandler(db *database.Queries, connManager *wsync.ConnectionManager, queue *wsync.SyncQueueProcessor) *KOReaderHandler { + return &KOReaderHandler{db: db, connManager: connManager, queue: queue} } type KOReaderProgressRequest struct { @@ -168,6 +169,15 @@ func (h *KOReaderHandler) SyncProgress(c echo.Context) error { pgUserID := pgtype.UUID{Bytes: userID, Valid: true} + syncMode := req.SyncMode + if syncMode == "" { + syncMode = "immediate" + } + + if syncMode == "checkpoint" { + return h.handleCheckpointSync(c, device, pgUserID, req) + } + booksSynced := 0 conflicts := []KOReaderConflict{} @@ -222,7 +232,7 @@ func (h *KOReaderHandler) SyncProgress(c echo.Context) error { }) } - if req.SyncMode == "immediate" { + if syncMode == "immediate" { return c.JSON(http.StatusAccepted, KOReaderSyncResponse{ SyncStatus: "accepted", BooksSynced: booksSynced, @@ -241,6 +251,90 @@ func (h *KOReaderHandler) SyncProgress(c echo.Context) error { }) } +func (h *KOReaderHandler) handleCheckpointSync(c echo.Context, device database.Devices, userID pgtype.UUID, req KOReaderProgressRequest) error { + booksEnqueued := 0 + + for _, book := range req.Books { + var mediaUUID uuid.UUID + var err error + + if book.UUID != "" { + mediaUUID, err = uuid.Parse(book.UUID) + if err != nil { + continue + } + + mediaItem, err := h.db.GetMediaItem(c.Request().Context(), pgtype.UUID{Bytes: mediaUUID, Valid: true}) + if err == nil { + err = h.enqueueProgressForBook(c, device.ID, userID, mediaItem.ID, book) + if err == nil { + booksEnqueued++ + } + } + } else if book.FilePath != "" { + mediaItem, err := h.db.GetMediaItemByFilePath(c.Request().Context(), book.FilePath) + if err == nil { + err = h.enqueueProgressForBook(c, device.ID, userID, mediaItem.ID, book) + if err == nil { + booksEnqueued++ + } + } + } else if book.Title != "" { + mediaItems, err := h.db.ListMediaItems(c.Request().Context(), database.ListMediaItemsParams{ + Limit: 100, + Offset: 0, + }) + if err == nil { + for _, mi := range mediaItems { + if mi.Title == book.Title && (book.Authors == nil || mi.Author.String == book.Authors[0]) { + err = h.enqueueProgressForBook(c, device.ID, userID, mi.ID, book) + if err == nil { + booksEnqueued++ + } + break + } + } + } + } + } + + _, err := h.db.UpdateDeviceLastSync(c.Request().Context(), device.ID) + if err != nil { + return c.JSON(http.StatusInternalServerError, map[string]string{ + "error": "failed to update device timestamp", + }) + } + + return c.JSON(http.StatusAccepted, map[string]interface{}{ + "sync_status": "checkpoint_enqueued", + "books_enqueued": booksEnqueued, + "message": "Sync will be processed in the background", + "timestamp": time.Now().Format(time.RFC3339), + }) +} + +func (h *KOReaderHandler) enqueueProgressForBook(c echo.Context, deviceID pgtype.UUID, userID pgtype.UUID, mediaItemID pgtype.UUID, book KOReaderBookProgress) error { + if h.queue == nil { + return fmt.Errorf("sync queue not available") + } + + update := &wsync.ProgressUpdate{ + DeviceID: deviceID, + MediaItemID: mediaItemID, + UserID: userID, + Percentage: book.Percentage, + Epubcfi: book.Epubcfi, + Chapter: book.Chapter, + Character: book.Character, + Page: book.Page, + TotalPages: book.TotalPages, + Source: "koreader", + SyncMode: "checkpoint", + } + + return h.queue.EnqueueProgress(update) +} + func (h *KOReaderHandler) updateProgressForBook(c echo.Context, userID pgtype.UUID, mediaItemID pgtype.UUID, book KOReaderBookProgress) error { ctx := c.Request().Context()