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
This commit is contained in:
2026-01-31 13:06:56 -05:00
parent 50c632babf
commit a7f2b83bdd
2 changed files with 116 additions and 4 deletions
+19 -1
View File
@@ -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)
+97 -3
View File
@@ -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()