Simplify GetDeviceStatus return by using inline new() instead of assigning to a local variable and returning its address.
255 lines
6.1 KiB
Go
255 lines
6.1 KiB
Go
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
|
|
}
|