Add sync conflict detection and resolution system
Implement conflict detection for concurrent reading progress updates from different devices. Adds conflict management endpoints for listing, viewing, and resolving conflicts. - Add ConflictHandler with CRUD endpoints for conflict management - Implement automatic conflict detection in KOReader progress updates - Add WebSocket broadcast for real-time conflict notifications - Add database query for listing user conflicts by status - Add integration tests and Bruno API test collection
This commit is contained in:
@@ -754,6 +754,13 @@ RETURNING *;
|
||||
-- name: DeleteSyncConflict :exec
|
||||
DELETE FROM sync_conflicts WHERE id = $1;
|
||||
|
||||
-- name: ListAllConflictsByUserAndStatus :many
|
||||
SELECT sc.*, mi.title, mi.author
|
||||
FROM sync_conflicts sc
|
||||
JOIN media_items mi ON sc.media_item_id = mi.id
|
||||
WHERE sc.user_id = $1 AND sc.resolution_status = $2
|
||||
ORDER BY sc.created_at DESC;
|
||||
|
||||
-- ============================================
|
||||
-- PHASE 3: KOREADER SYNC PROTOCOL (Weeks 7-9)
|
||||
-- ============================================
|
||||
|
||||
@@ -0,0 +1,403 @@
|
||||
package handlers
|
||||
|
||||
import (
|
||||
"bookmann/internal/database"
|
||||
wsync "bookmann/internal/sync"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
"github.com/labstack/echo/v4"
|
||||
)
|
||||
|
||||
type ConflictHandler struct {
|
||||
db *database.Queries
|
||||
connManager *wsync.ConnectionManager
|
||||
}
|
||||
|
||||
func NewConflictHandler(db *database.Queries, connManager *wsync.ConnectionManager) *ConflictHandler {
|
||||
return &ConflictHandler{
|
||||
db: db,
|
||||
connManager: connManager,
|
||||
}
|
||||
}
|
||||
|
||||
type ConflictResolutionRequest struct {
|
||||
Winner string `json:"winner" validate:"required,oneof=koreader kobo web manual"`
|
||||
ManualData map[string]interface{} `json:"manual_data"`
|
||||
ApplyToAll bool `json:"apply_to_all_future_conflicts"`
|
||||
Reason string `json:"reason"`
|
||||
}
|
||||
|
||||
type ConflictSourceData struct {
|
||||
Source string `json:"source"`
|
||||
Timestamp time.Time `json:"timestamp"`
|
||||
Data map[string]interface{} `json:"data"`
|
||||
}
|
||||
|
||||
type ConflictDetailResponse struct {
|
||||
ID string `json:"id"`
|
||||
MediaItemID string `json:"media_item_id"`
|
||||
MediaItemTitle string `json:"media_item_title"`
|
||||
ConflictType string `json:"conflict_type"`
|
||||
ConflictData map[string]ConflictSourceData `json:"conflict_data"`
|
||||
ResolutionStatus string `json:"resolution_status"`
|
||||
ResolutionData map[string]interface{} `json:"resolution_data,omitempty"`
|
||||
ResolvedBy string `json:"resolved_by,omitempty"`
|
||||
ResolvedAt *time.Time `json:"resolved_at,omitempty"`
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
}
|
||||
|
||||
type ConflictListResponse struct {
|
||||
Conflicts []ConflictDetailResponse `json:"conflicts"`
|
||||
Total int `json:"total"`
|
||||
Unresolved int `json:"unresolved"`
|
||||
}
|
||||
|
||||
type ConflictResolveResponse struct {
|
||||
ConflictResolved bool `json:"conflict_resolved"`
|
||||
AppliedTo map[string]bool `json:"applied_to"`
|
||||
DevicesSynced []string `json:"devices_synced"`
|
||||
}
|
||||
|
||||
func (h *ConflictHandler) ListConflicts(c echo.Context) error {
|
||||
user := c.Get("user").(database.Users)
|
||||
|
||||
status := c.QueryParam("status")
|
||||
if status == "" {
|
||||
status = "unresolved"
|
||||
}
|
||||
|
||||
ctx := context.Background()
|
||||
|
||||
conflicts, err := h.db.ListAllConflictsByUserAndStatus(ctx, database.ListAllConflictsByUserAndStatusParams{
|
||||
UserID: user.ID,
|
||||
ResolutionStatus: pgtype.Text{String: status, Valid: true},
|
||||
})
|
||||
|
||||
if err != nil && err != pgx.ErrNoRows {
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to list conflicts")
|
||||
}
|
||||
|
||||
response := ConflictListResponse{
|
||||
Conflicts: make([]ConflictDetailResponse, 0),
|
||||
Total: len(conflicts),
|
||||
Unresolved: 0,
|
||||
}
|
||||
|
||||
for _, conflict := range conflicts {
|
||||
var conflictData map[string]ConflictSourceData
|
||||
if err := json.Unmarshal(conflict.ConflictData, &conflictData); err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
detail := ConflictDetailResponse{
|
||||
ID: uuid.UUID(conflict.ID.Bytes).String(),
|
||||
MediaItemID: uuid.UUID(conflict.MediaItemID.Bytes).String(),
|
||||
MediaItemTitle: conflict.Title,
|
||||
ConflictType: conflict.ConflictType,
|
||||
ConflictData: conflictData,
|
||||
ResolutionStatus: conflict.ResolutionStatus.String,
|
||||
CreatedAt: conflict.CreatedAt.Time,
|
||||
}
|
||||
|
||||
if conflict.ResolvedBy.Valid {
|
||||
detail.ResolvedBy = uuid.UUID(conflict.ResolvedBy.Bytes).String()
|
||||
}
|
||||
if conflict.ResolvedAt.Valid {
|
||||
detail.ResolvedAt = &conflict.ResolvedAt.Time
|
||||
}
|
||||
if conflict.ResolutionData != nil {
|
||||
if err := json.Unmarshal(conflict.ResolutionData, &detail.ResolutionData); err == nil {
|
||||
}
|
||||
}
|
||||
|
||||
response.Conflicts = append(response.Conflicts, detail)
|
||||
if conflict.ResolutionStatus.String == "unresolved" {
|
||||
response.Unresolved++
|
||||
}
|
||||
}
|
||||
|
||||
return c.JSON(http.StatusOK, response)
|
||||
}
|
||||
|
||||
func (h *ConflictHandler) GetConflict(c echo.Context) error {
|
||||
user := c.Get("user").(database.Users)
|
||||
|
||||
conflictID, err := uuid.Parse(c.Param("id"))
|
||||
if err != nil {
|
||||
return echo.NewHTTPError(http.StatusBadRequest, "invalid conflict ID")
|
||||
}
|
||||
|
||||
conflictUUID := pgtype.UUID{Bytes: [16]byte(conflictID), Valid: true}
|
||||
conflict, err := h.db.GetSyncConflict(context.Background(), conflictUUID)
|
||||
if err != nil {
|
||||
if err == pgx.ErrNoRows {
|
||||
return echo.NewHTTPError(http.StatusNotFound, "conflict not found")
|
||||
}
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to get conflict")
|
||||
}
|
||||
|
||||
if conflict.UserID.Bytes != user.ID.Bytes {
|
||||
return echo.NewHTTPError(http.StatusForbidden, "access denied")
|
||||
}
|
||||
|
||||
mediaItem, err := h.db.GetMediaItem(context.Background(), conflict.MediaItemID)
|
||||
if err != nil {
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to get media item")
|
||||
}
|
||||
|
||||
var conflictData map[string]ConflictSourceData
|
||||
if err := json.Unmarshal(conflict.ConflictData, &conflictData); err != nil {
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to parse conflict data")
|
||||
}
|
||||
|
||||
detail := ConflictDetailResponse{
|
||||
ID: uuid.UUID(conflict.ID.Bytes).String(),
|
||||
MediaItemID: uuid.UUID(conflict.MediaItemID.Bytes).String(),
|
||||
MediaItemTitle: mediaItem.Title,
|
||||
ConflictType: conflict.ConflictType,
|
||||
ConflictData: conflictData,
|
||||
ResolutionStatus: conflict.ResolutionStatus.String,
|
||||
CreatedAt: conflict.CreatedAt.Time,
|
||||
}
|
||||
|
||||
if conflict.ResolvedBy.Valid {
|
||||
detail.ResolvedBy = uuid.UUID(conflict.ResolvedBy.Bytes).String()
|
||||
}
|
||||
if conflict.ResolvedAt.Valid {
|
||||
detail.ResolvedAt = &conflict.ResolvedAt.Time
|
||||
}
|
||||
if conflict.ResolutionData != nil {
|
||||
if err := json.Unmarshal(conflict.ResolutionData, &detail.ResolutionData); err == nil {
|
||||
}
|
||||
}
|
||||
|
||||
return c.JSON(http.StatusOK, detail)
|
||||
}
|
||||
|
||||
func (h *ConflictHandler) ResolveConflict(c echo.Context) error {
|
||||
user := c.Get("user").(database.Users)
|
||||
|
||||
conflictID, err := uuid.Parse(c.Param("id"))
|
||||
if err != nil {
|
||||
return echo.NewHTTPError(http.StatusBadRequest, "invalid conflict ID")
|
||||
}
|
||||
|
||||
var req ConflictResolutionRequest
|
||||
if err := c.Bind(&req); err != nil {
|
||||
return echo.NewHTTPError(http.StatusBadRequest, "invalid request body")
|
||||
}
|
||||
|
||||
if req.Winner == "manual" && req.ManualData == nil {
|
||||
return echo.NewHTTPError(http.StatusBadRequest, "manual_data required when winner is manual")
|
||||
}
|
||||
|
||||
conflictUUID := pgtype.UUID{Bytes: [16]byte(conflictID), Valid: true}
|
||||
conflict, err := h.db.GetSyncConflict(context.Background(), conflictUUID)
|
||||
if err != nil {
|
||||
if err == pgx.ErrNoRows {
|
||||
return echo.NewHTTPError(http.StatusNotFound, "conflict not found")
|
||||
}
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to get conflict")
|
||||
}
|
||||
|
||||
if conflict.UserID.Bytes != user.ID.Bytes {
|
||||
return echo.NewHTTPError(http.StatusForbidden, "access denied")
|
||||
}
|
||||
|
||||
if conflict.ResolutionStatus.String != "unresolved" {
|
||||
return echo.NewHTTPError(http.StatusBadRequest, "conflict already resolved")
|
||||
}
|
||||
|
||||
var conflictData map[string]ConflictSourceData
|
||||
if err := json.Unmarshal(conflict.ConflictData, &conflictData); err != nil {
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to parse conflict data")
|
||||
}
|
||||
|
||||
winnerData := map[string]interface{}{}
|
||||
if req.Winner == "manual" {
|
||||
winnerData = req.ManualData
|
||||
} else {
|
||||
source, ok := conflictData[req.Winner]
|
||||
if !ok {
|
||||
return echo.NewHTTPError(http.StatusBadRequest, "invalid winner source")
|
||||
}
|
||||
winnerData = source.Data
|
||||
}
|
||||
|
||||
appliedTo := map[string]bool{
|
||||
"progress": false,
|
||||
"annotations": false,
|
||||
}
|
||||
|
||||
if conflict.ConflictType == "progress" {
|
||||
if err := h.applyProgressResolution(conflict.MediaItemID, conflict.UserID, winnerData); err == nil {
|
||||
appliedTo["progress"] = true
|
||||
}
|
||||
}
|
||||
|
||||
resolutionData := map[string]interface{}{
|
||||
"winner": req.Winner,
|
||||
"applied_to": appliedTo,
|
||||
"reason": req.Reason,
|
||||
"resolved_at": time.Now(),
|
||||
}
|
||||
|
||||
resolutionDataJSON, _ := json.Marshal(resolutionData)
|
||||
|
||||
_, err = h.db.ResolveSyncConflict(context.Background(), database.ResolveSyncConflictParams{
|
||||
ID: conflictUUID,
|
||||
ResolutionStatus: pgtype.Text{String: "user_resolved", Valid: true},
|
||||
ResolutionData: resolutionDataJSON,
|
||||
ResolvedBy: pgtype.UUID{Bytes: user.ID.Bytes, Valid: true},
|
||||
})
|
||||
if err != nil {
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to resolve conflict")
|
||||
}
|
||||
|
||||
devicesSynced := h.notifyDevicesOfResolution(conflict.MediaItemID, winnerData)
|
||||
|
||||
response := ConflictResolveResponse{
|
||||
ConflictResolved: true,
|
||||
AppliedTo: appliedTo,
|
||||
DevicesSynced: devicesSynced,
|
||||
}
|
||||
|
||||
return c.JSON(http.StatusOK, response)
|
||||
}
|
||||
|
||||
func (h *ConflictHandler) applyProgressResolution(mediaItemID pgtype.UUID, userID pgtype.UUID, data map[string]interface{}) error {
|
||||
ctx := context.Background()
|
||||
|
||||
existingProgress, err := h.db.GetReadingProgress(ctx, database.GetReadingProgressParams{
|
||||
MediaItemID: mediaItemID,
|
||||
UserID: userID,
|
||||
})
|
||||
if err != nil && err != pgx.ErrNoRows {
|
||||
return err
|
||||
}
|
||||
|
||||
percentage := 0.0
|
||||
if p, ok := data["percentage"].(float64); ok {
|
||||
percentage = p
|
||||
}
|
||||
|
||||
var epubcfi pgtype.Text
|
||||
if e, ok := data["epubcfi"].(string); ok {
|
||||
epubcfi = pgtype.Text{String: e, Valid: true}
|
||||
}
|
||||
|
||||
var chapter pgtype.Int4
|
||||
if c, ok := data["chapter"].(float64); ok {
|
||||
chapter = pgtype.Int4{Int32: int32(c), Valid: true}
|
||||
}
|
||||
|
||||
var characterOffset pgtype.Int8
|
||||
if c, ok := data["character"].(float64); ok {
|
||||
characterOffset = pgtype.Int8{Int64: int64(c), Valid: true}
|
||||
}
|
||||
|
||||
currentPage := existingProgress.CurrentPage
|
||||
totalPages := existingProgress.TotalPages
|
||||
|
||||
if p, ok := data["page"].(float64); ok {
|
||||
currentPage = pgtype.Int4{Int32: int32(p), Valid: true}
|
||||
}
|
||||
if p, ok := data["total_pages"].(float64); ok {
|
||||
totalPages = pgtype.Int4{Int32: int32(p), Valid: true}
|
||||
}
|
||||
|
||||
_, err = h.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: "conflict_resolution", Valid: true},
|
||||
LastSyncSource: pgtype.Text{String: "manual", Valid: true},
|
||||
ViewportY: pgtype.Float8{},
|
||||
ScrollPositionX: pgtype.Float8{},
|
||||
ScrollPositionY: pgtype.Float8{},
|
||||
PanelNumber: pgtype.Int4{},
|
||||
ReadingMode: pgtype.Text{},
|
||||
ZoomLevel: pgtype.Float8{},
|
||||
})
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
func (h *ConflictHandler) notifyDevicesOfResolution(mediaItemID pgtype.UUID, data map[string]interface{}) []string {
|
||||
devices, err := h.db.ListDevicesByType(context.Background(), "koreader")
|
||||
if err != nil {
|
||||
return []string{}
|
||||
}
|
||||
|
||||
synced := []string{}
|
||||
for _, device := range devices {
|
||||
if device.SyncEnabled.Bool {
|
||||
synced = append(synced, uuid.UUID(device.ID.Bytes).String())
|
||||
}
|
||||
}
|
||||
|
||||
return synced
|
||||
}
|
||||
|
||||
func (h *ConflictHandler) DeleteConflict(c echo.Context) error {
|
||||
user := c.Get("user").(database.Users)
|
||||
|
||||
conflictID, err := uuid.Parse(c.Param("id"))
|
||||
if err != nil {
|
||||
return echo.NewHTTPError(http.StatusBadRequest, "invalid conflict ID")
|
||||
}
|
||||
|
||||
conflictUUID := pgtype.UUID{Bytes: [16]byte(conflictID), Valid: true}
|
||||
conflict, err := h.db.GetSyncConflict(context.Background(), conflictUUID)
|
||||
if err != nil {
|
||||
if err == pgx.ErrNoRows {
|
||||
return echo.NewHTTPError(http.StatusNotFound, "conflict not found")
|
||||
}
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to get conflict")
|
||||
}
|
||||
|
||||
if conflict.UserID.Bytes != user.ID.Bytes {
|
||||
return echo.NewHTTPError(http.StatusForbidden, "access denied")
|
||||
}
|
||||
|
||||
if err := h.db.DeleteSyncConflict(context.Background(), conflictUUID); err != nil {
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to delete conflict")
|
||||
}
|
||||
|
||||
return c.NoContent(http.StatusNoContent)
|
||||
}
|
||||
|
||||
func (h *ConflictHandler) DismissAllResolved(c echo.Context) error {
|
||||
user := c.Get("user").(database.Users)
|
||||
|
||||
conflicts, err := h.db.ListAllConflictsByUserAndStatus(context.Background(), database.ListAllConflictsByUserAndStatusParams{
|
||||
UserID: user.ID,
|
||||
ResolutionStatus: pgtype.Text{String: "user_resolved", Valid: true},
|
||||
})
|
||||
if err != nil {
|
||||
return echo.NewHTTPError(http.StatusInternalServerError, "failed to list conflicts")
|
||||
}
|
||||
|
||||
deleted := 0
|
||||
for _, conflict := range conflicts {
|
||||
if err := h.db.DeleteSyncConflict(context.Background(), conflict.ID); err == nil {
|
||||
deleted++
|
||||
}
|
||||
}
|
||||
|
||||
return c.JSON(http.StatusOK, map[string]interface{}{
|
||||
"deleted": deleted,
|
||||
})
|
||||
}
|
||||
+122
-14
@@ -3,11 +3,13 @@ package handlers
|
||||
import (
|
||||
"bookmann/internal/database"
|
||||
wsync "bookmann/internal/sync"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/jackc/pgx/v5"
|
||||
"github.com/jackc/pgx/v5/pgtype"
|
||||
"github.com/labstack/echo/v4"
|
||||
)
|
||||
@@ -240,6 +242,35 @@ func (h *KOReaderHandler) SyncProgress(c echo.Context) error {
|
||||
}
|
||||
|
||||
func (h *KOReaderHandler) updateProgressForBook(c echo.Context, userID pgtype.UUID, mediaItemID pgtype.UUID, book KOReaderBookProgress) error {
|
||||
ctx := c.Request().Context()
|
||||
|
||||
existingProgress, err := h.db.GetReadingProgress(ctx, database.GetReadingProgressParams{
|
||||
MediaItemID: mediaItemID,
|
||||
UserID: userID,
|
||||
})
|
||||
|
||||
if err != nil && err != pgx.ErrNoRows {
|
||||
return err
|
||||
}
|
||||
|
||||
hasExistingProgress := err != pgx.ErrNoRows
|
||||
conflictDetected := false
|
||||
|
||||
if hasExistingProgress && existingProgress.LastSyncSource.Valid {
|
||||
if existingProgress.LastSyncSource.String != "koreader" && existingProgress.LastSyncTimestamp.Valid {
|
||||
timeDiff := time.Since(existingProgress.LastSyncTimestamp.Time)
|
||||
if timeDiff < 5*time.Minute {
|
||||
percentageDiff := book.Percentage - existingProgress.Percentage.Float64
|
||||
if percentageDiff < 0 {
|
||||
percentageDiff = -percentageDiff
|
||||
}
|
||||
if percentageDiff > 0.01 {
|
||||
conflictDetected = true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
var epubcfi pgtype.Text
|
||||
var chapter pgtype.Int4
|
||||
var characterOffset pgtype.Int8
|
||||
@@ -262,7 +293,7 @@ func (h *KOReaderHandler) updateProgressForBook(c echo.Context, userID pgtype.UU
|
||||
totalPages = pgtype.Int4{Int32: int32(*book.TotalPages), Valid: true}
|
||||
}
|
||||
|
||||
_, err := h.db.UpdateUniversalProgress(c.Request().Context(), database.UpdateUniversalProgressParams{
|
||||
_, err = h.db.UpdateUniversalProgress(ctx, database.UpdateUniversalProgressParams{
|
||||
MediaItemID: mediaItemID,
|
||||
UserID: userID,
|
||||
Percentage: pgtype.Float8{Float64: book.Percentage, Valid: true},
|
||||
@@ -274,26 +305,103 @@ func (h *KOReaderHandler) updateProgressForBook(c echo.Context, userID pgtype.UU
|
||||
TotalPages: totalPages,
|
||||
LastSyncDevice: pgtype.Text{String: "koreader", Valid: true},
|
||||
LastSyncSource: pgtype.Text{String: "koreader", Valid: true},
|
||||
ViewportY: pgtype.Float8{},
|
||||
ScrollPositionX: pgtype.Float8{},
|
||||
ScrollPositionY: pgtype.Float8{},
|
||||
PanelNumber: pgtype.Int4{},
|
||||
ReadingMode: pgtype.Text{},
|
||||
ZoomLevel: pgtype.Float8{},
|
||||
})
|
||||
|
||||
if err == nil {
|
||||
// Broadcast progress update to all connected WebSocket clients
|
||||
deviceInfo := book.DeviceInfo
|
||||
if deviceInfo.DeviceModel == "" {
|
||||
deviceInfo.DeviceModel = "KOReader Device"
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if conflictDetected {
|
||||
koreaderData := map[string]interface{}{
|
||||
"source": "koreader",
|
||||
"timestamp": time.Now(),
|
||||
"data": map[string]interface{}{
|
||||
"percentage": book.Percentage,
|
||||
},
|
||||
}
|
||||
if book.Epubcfi != nil {
|
||||
koreaderData["data"].(map[string]interface{})["epubcfi"] = *book.Epubcfi
|
||||
}
|
||||
if book.Chapter != nil {
|
||||
koreaderData["data"].(map[string]interface{})["chapter"] = *book.Chapter
|
||||
}
|
||||
if book.Character != nil {
|
||||
koreaderData["data"].(map[string]interface{})["character"] = *book.Character
|
||||
}
|
||||
if book.Page != nil {
|
||||
koreaderData["data"].(map[string]interface{})["page"] = *book.Page
|
||||
}
|
||||
if book.TotalPages != nil {
|
||||
koreaderData["data"].(map[string]interface{})["total_pages"] = *book.TotalPages
|
||||
}
|
||||
|
||||
h.connManager.BroadcastProgressUpdate(
|
||||
uuid.UUID(mediaItemID.Bytes),
|
||||
book.Percentage,
|
||||
wsync.SourceDevice{
|
||||
ID: uuid.UUID(userID.Bytes).String(),
|
||||
Name: deviceInfo.DeviceModel,
|
||||
Type: "koreader",
|
||||
existingData := map[string]interface{}{
|
||||
"source": existingProgress.LastSyncSource.String,
|
||||
"timestamp": existingProgress.LastSyncTimestamp.Time,
|
||||
"data": map[string]interface{}{
|
||||
"percentage": existingProgress.Percentage.Float64,
|
||||
},
|
||||
)
|
||||
}
|
||||
if existingProgress.Epubcfi.Valid {
|
||||
existingData["data"].(map[string]interface{})["epubcfi"] = existingProgress.Epubcfi.String
|
||||
}
|
||||
if existingProgress.Chapter.Valid {
|
||||
existingData["data"].(map[string]interface{})["chapter"] = existingProgress.Chapter.Int32
|
||||
}
|
||||
if existingProgress.CharacterOffset.Valid {
|
||||
existingData["data"].(map[string]interface{})["character"] = existingProgress.CharacterOffset.Int64
|
||||
}
|
||||
if existingProgress.CurrentPage.Valid {
|
||||
existingData["data"].(map[string]interface{})["page"] = existingProgress.CurrentPage.Int32
|
||||
}
|
||||
if existingProgress.TotalPages.Valid {
|
||||
existingData["data"].(map[string]interface{})["total_pages"] = existingProgress.TotalPages.Int32
|
||||
}
|
||||
|
||||
conflictData := map[string]interface{}{
|
||||
"koreader": koreaderData,
|
||||
"existing": existingData,
|
||||
}
|
||||
conflictDataJSON, _ := json.Marshal(conflictData)
|
||||
|
||||
_, err := h.db.CreateSyncConflict(ctx, database.CreateSyncConflictParams{
|
||||
MediaItemID: mediaItemID,
|
||||
UserID: userID,
|
||||
ConflictType: "progress",
|
||||
ConflictData: conflictDataJSON,
|
||||
})
|
||||
if err == nil {
|
||||
h.connManager.BroadcastConflictNotification(
|
||||
mediaItemID.Bytes,
|
||||
"detection",
|
||||
"",
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
deviceInfo := book.DeviceInfo
|
||||
if deviceInfo.DeviceModel == "" {
|
||||
deviceInfo.DeviceModel = "KOReader Device"
|
||||
}
|
||||
|
||||
h.connManager.BroadcastProgressUpdate(
|
||||
uuid.UUID(mediaItemID.Bytes),
|
||||
book.Percentage,
|
||||
wsync.SourceDevice{
|
||||
ID: uuid.UUID(userID.Bytes).String(),
|
||||
Name: deviceInfo.DeviceModel,
|
||||
Type: "koreader",
|
||||
},
|
||||
)
|
||||
|
||||
_, err = h.db.UpdateDeviceLastSync(ctx, pgtype.UUID{Bytes: [16]byte{}, Valid: false})
|
||||
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
@@ -121,6 +121,20 @@ func (m *ConnectionManager) BroadcastAnnotationUpdate(bookID uuid.UUID, annotati
|
||||
m.Broadcast(msg)
|
||||
}
|
||||
|
||||
// BroadcastConflictNotification broadcasts a conflict notification to all connected clients
|
||||
func (m *ConnectionManager) BroadcastConflictNotification(bookID [16]byte, notificationType string, conflictID string) {
|
||||
msg := BroadcastMessage{
|
||||
Type: MessageTypeConflict,
|
||||
Timestamp: time.Now().Format(time.RFC3339),
|
||||
Data: map[string]interface{}{
|
||||
"book_id": uuid.UUID(bookID).String(),
|
||||
"notification_type": notificationType,
|
||||
"conflict_id": conflictID,
|
||||
},
|
||||
}
|
||||
m.Broadcast(msg)
|
||||
}
|
||||
|
||||
// AddConnection adds a new WebSocket connection
|
||||
func (m *ConnectionManager) AddConnection(conn *DeviceConnection) {
|
||||
m.mu.Lock()
|
||||
|
||||
Reference in New Issue
Block a user