package sync import ( "bookhoard/internal/database" "context" "log" "time" "github.com/jackc/pgx/v5/pgtype" ) const ( OnlineThreshold = 5 * time.Minute OfflineCheckInterval = 2 * time.Minute ) type OfflineDetector struct { db *database.Queries queue *SyncQueueProcessor checkInterval time.Duration onlineThreshold time.Duration } type DeviceOnlineStatus struct { DeviceID pgtype.UUID DeviceName string DeviceType string IsOnline bool LastSeen time.Time UserID pgtype.UUID } func NewOfflineDetector(db *database.Queries, queue *SyncQueueProcessor) *OfflineDetector { return &OfflineDetector{ db: db, queue: queue, checkInterval: OfflineCheckInterval, onlineThreshold: OnlineThreshold, } } func (d *OfflineDetector) Start(ctx context.Context) { log.Printf("Starting offline detector (threshold: %v)", d.onlineThreshold) ticker := time.NewTicker(d.checkInterval) defer ticker.Stop() for { select { case <-ctx.Done(): log.Println("Offline detector stopped") return case <-ticker.C: d.checkDeviceStatus(ctx) } } } func (d *OfflineDetector) checkDeviceStatus(ctx context.Context) { users, err := d.db.ListUsers(ctx) if err != nil { log.Printf("Failed to list users for offline check: %v", err) return } for _, user := range users { devices, err := d.db.ListDevicesByUser(ctx, user.ID) if err != nil { continue } for _, device := range devices { if !device.SyncEnabled.Bool { continue } status := d.getDeviceStatus(device) if !status.IsOnline { d.handleOfflineDevice(ctx, status) } else if d.wasOffline(ctx, status.DeviceID) { d.handleReconnectedDevice(ctx, status) } d.updateDeviceStatus(ctx, status) } } } func (d *OfflineDetector) getDeviceStatus(device database.Devices) DeviceOnlineStatus { lastSeen := time.Now() if device.LastSeen.Valid { lastSeen = device.LastSeen.Time } isOnline := time.Since(lastSeen) < d.onlineThreshold return DeviceOnlineStatus{ DeviceID: device.ID, DeviceName: device.DeviceName, DeviceType: device.DeviceType, IsOnline: isOnline, LastSeen: lastSeen, UserID: device.UserID, } } func (d *OfflineDetector) wasOffline(ctx context.Context, deviceID pgtype.UUID) bool { device, err := d.db.GetDevice(ctx, deviceID) if err != nil { return false } if !device.LastSeen.Valid { return false } timeSinceLastSeen := time.Since(device.LastSeen.Time) return timeSinceLastSeen >= d.onlineThreshold && timeSinceLastSeen < d.onlineThreshold*2 } func (d *OfflineDetector) handleOfflineDevice(ctx context.Context, status DeviceOnlineStatus) { log.Printf("Device offline detected: %s (%s)", status.DeviceName, status.DeviceType) d.setOfflineMode(ctx, status.DeviceID) } func (d *OfflineDetector) handleReconnectedDevice(ctx context.Context, status DeviceOnlineStatus) { log.Printf("Device reconnected: %s (%s)", status.DeviceName, status.DeviceType) d.setOnlineMode(ctx, status.DeviceID) d.processQueuedItems(ctx, status.DeviceID) } func (d *OfflineDetector) setOfflineMode(ctx context.Context, deviceID pgtype.UUID) { _, err := d.db.UpdateDevice(ctx, database.UpdateDeviceParams{ ID: deviceID, DeviceName: "", SyncEnabled: pgtype.Bool{Bool: false, Valid: true}, AutoSync: pgtype.Bool{}, }) if err != nil { log.Printf("Failed to set offline mode for device: %v", err) } } func (d *OfflineDetector) setOnlineMode(ctx context.Context, deviceID pgtype.UUID) { _, err := d.db.UpdateDevice(ctx, database.UpdateDeviceParams{ ID: deviceID, DeviceName: "", SyncEnabled: pgtype.Bool{Bool: true, Valid: true}, AutoSync: pgtype.Bool{}, }) if err != nil { log.Printf("Failed to set online mode for device: %v", err) } } func (d *OfflineDetector) updateDeviceStatus(ctx context.Context, status DeviceOnlineStatus) { _, err := d.db.UpdateDeviceLastSeen(ctx, status.DeviceID) if err != nil { log.Printf("Failed to update device last seen: %v", err) } } func (d *OfflineDetector) processQueuedItems(ctx context.Context, deviceID pgtype.UUID) { log.Printf("Processing queued items for reconnected device %s", deviceID.Bytes) items, err := d.db.ListPendingSyncQueueItems(ctx, database.ListPendingSyncQueueItemsParams{ DeviceID: deviceID, Limit: 50, }) if err != nil { log.Printf("Failed to list queued items for device: %v", err) return } log.Printf("Found %d queued items for device", len(items)) priorityItems := make([]database.SyncQueue, 0) for _, item := range items { if item.Priority.Int32 <= PriorityBookCompletion { priorityItems = append(priorityItems, item) } } if len(priorityItems) > 0 { log.Printf("Processing %d priority items for device", len(priorityItems)) for _, item := range priorityItems { _, err := d.db.UpdateSyncQueueItemStatus(ctx, database.UpdateSyncQueueItemStatusParams{ ID: item.ID, Status: pgtype.Text{String: "pending", Valid: true}, ErrorMessage: pgtype.Text{}, }) if err != nil { log.Printf("Failed to re-priority queue item: %v", err) } } } } func (d *OfflineDetector) ForceReconnectDevice(ctx context.Context, deviceID pgtype.UUID) error { device, err := d.db.GetDevice(ctx, deviceID) if err != nil { return err } status := d.getDeviceStatus(device) status.IsOnline = true status.LastSeen = time.Now() d.handleReconnectedDevice(ctx, status) d.updateDeviceStatus(ctx, status) return nil } func (d *OfflineDetector) GetDeviceStatus(ctx context.Context, deviceID pgtype.UUID) (*DeviceOnlineStatus, error) { device, err := d.db.GetDevice(ctx, deviceID) if err != nil { return nil, err } return new(d.getDeviceStatus(device)), nil } func (d *OfflineDetector) ListOfflineDevices(ctx context.Context) ([]DeviceOnlineStatus, error) { users, err := d.db.ListUsers(ctx) if err != nil { return nil, err } var offlineDevices []DeviceOnlineStatus for _, user := range users { devices, err := d.db.ListDevicesByUser(ctx, user.ID) if err != nil { continue } for _, device := range devices { status := d.getDeviceStatus(device) if !status.IsOnline { offlineDevices = append(offlineDevices, status) } } } return offlineDevices, nil }