package queue import ( "bytes" "context" "log/slog" "sync" "time" "gobuz/internal/downloader" "gobuz/internal/log" "gobuz/internal/models" "gobuz/internal/web/pages" ) type UserQueue struct { userID int64 mu sync.RWMutex items map[string]*models.QueueItem order []string broker *SSEBroker dl *downloader.Downloader lockRegistry *AlbumLockRegistry pages *pages.Pages currentAlbumID string trackToAlbumMap map[string]string isProcessing bool logger *slog.Logger onComplete func() ctx context.Context cancel context.CancelFunc } func newUserQueue(userID int64, broker *SSEBroker, dl *downloader.Downloader, lockRegistry *AlbumLockRegistry, pgs *pages.Pages, logger *slog.Logger, onComplete func()) *UserQueue { ctx, cancel := context.WithCancel(context.Background()) return &UserQueue{ userID: userID, items: make(map[string]*models.QueueItem), order: make([]string, 0), broker: broker, dl: dl, lockRegistry: lockRegistry, pages: pgs, trackToAlbumMap: make(map[string]string), logger: logger, onComplete: onComplete, ctx: ctx, cancel: cancel, } } func (u *UserQueue) setOnComplete(fn func()) { u.mu.Lock() defer u.mu.Unlock() u.onComplete = fn } func (u *UserQueue) Stop() { u.cancel() } func (u *UserQueue) IsProcessing() bool { u.mu.RLock() defer u.mu.RUnlock() return u.isProcessing } func (u *UserQueue) SetProcessing(val bool) { u.mu.Lock() defer u.mu.Unlock() u.isProcessing = val } func (u *UserQueue) OnAlbum(msg downloader.MsgAlbum) { u.mu.Lock() defer u.mu.Unlock() if !u.isProcessing { return } albumID := msg.AlbumID if albumID != "" { if item, exists := u.items[albumID]; exists { u.currentAlbumID = albumID item.Status = models.StatusDownloading u.broadcastProgressLocked() } } } func (u *UserQueue) OnRegisterTrack(msg downloader.MsgRegisterTrack) { u.mu.Lock() defer u.mu.Unlock() if !u.isProcessing { return } albumID := msg.AlbumID if albumID == "" { albumID = u.currentAlbumID } if albumID == "" { return } album, exists := u.items[albumID] if !exists { album = &models.QueueItem{ Release: models.Release{ ID: albumID, Title: msg.Name, }, Status: models.StatusDownloading, } u.items[albumID] = album u.order = append(u.order, albumID) } album.Status = models.StatusDownloading u.trackToAlbumMap[msg.ID] = albumID var track *models.TrackProgress for _, t := range album.Tracks { if t.ID == msg.ID { track = t break } } if track == nil { track = &models.TrackProgress{ ID: msg.ID, TrackNumber: len(album.Tracks) + 1, Title: msg.Name, Status: models.StatusDownloading, Counter: msg.Counter, } album.Tracks = append(album.Tracks, track) } else { track.Title = msg.Name track.Status = models.StatusDownloading track.Counter = msg.Counter } u.broadcastProgressLocked() } func (u *UserQueue) OnSetTotal(msg downloader.MsgSetTotal) { u.mu.Lock() defer u.mu.Unlock() if albumID, ok := u.trackToAlbumMap[msg.ID]; ok { if album, exists := u.items[albumID]; exists { for _, t := range album.Tracks { if t.ID == msg.ID { t.TotalBytes = msg.Total u.broadcastProgressLocked() break } } } } } func (u *UserQueue) OnDone(msg downloader.MsgDone) { u.mu.Lock() defer u.mu.Unlock() if albumID, ok := u.trackToAlbumMap[msg.ID]; ok { if album, exists := u.items[albumID]; exists { allDone := true for _, t := range album.Tracks { if t.ID == msg.ID { t.Status = models.StatusDone if t.TotalBytes > 0 { t.DownloadedBytes = t.TotalBytes t.Percent = 100 } } if t.Status != models.StatusDone { allDone = false } } if len(album.Tracks) > 0 && allDone { album.Status = models.StatusDone album.Percent = 100 } u.broadcastProgressLocked() } } } func (u *UserQueue) OnFailed(msg downloader.MsgFailed) { u.mu.Lock() defer u.mu.Unlock() if albumID, ok := u.trackToAlbumMap[msg.ID]; ok { if album, exists := u.items[albumID]; exists { for _, t := range album.Tracks { if t.ID == msg.ID { t.Status = models.StatusFailed break } } album.Status = models.StatusFailed u.broadcastProgressLocked() } } } func (u *UserQueue) tick() { u.mu.Lock() defer u.mu.Unlock() changed := false for _, album := range u.items { if album.Status == models.StatusDownloading { var sumDownloaded int64 var sumTotal int64 var sumPercent int trackCount := len(album.Tracks) for _, t := range album.Tracks { if t.Status == models.StatusDownloading && t.Counter != nil { downloaded := t.Counter.Load() if downloaded != t.DownloadedBytes { t.DownloadedBytes = downloaded if t.TotalBytes > 0 { pct := int((float64(downloaded) / float64(t.TotalBytes)) * 100) if pct > 100 { pct = 100 } t.Percent = pct } changed = true } } sumDownloaded += t.DownloadedBytes sumTotal += t.TotalBytes sumPercent += t.Percent } if trackCount > 0 { album.DownloadedBytes = sumDownloaded album.TotalBytes = sumTotal album.Percent = sumPercent / trackCount } } } if changed { u.broadcastProgressLocked() } } func (u *UserQueue) IsQueued(id string) bool { if id == "" { return false } u.mu.RLock() defer u.mu.RUnlock() _, exists := u.items[id] return exists } func (u *UserQueue) StageRelease(item models.QueueItem) { u.mu.Lock() if item.ID == "" { u.mu.Unlock() return } if _, exists := u.items[item.ID]; !exists { item.Status = models.StatusStaged u.items[item.ID] = &item u.order = append(u.order, item.ID) } u.mu.Unlock() } func (u *UserQueue) RemoveRelease(id string) { if id == "" { return } u.mu.Lock() delete(u.items, id) newOrder := make([]string, 0, len(u.order)) for _, oid := range u.order { if oid != id { newOrder = append(newOrder, oid) } } u.order = newOrder u.mu.Unlock() } func (u *UserQueue) Clear() { u.mu.Lock() if !u.isProcessing { u.items = make(map[string]*models.QueueItem) u.order = nil u.trackToAlbumMap = make(map[string]string) } else { newItems := make(map[string]*models.QueueItem) var newOrder []string for _, id := range u.order { if item, ok := u.items[id]; ok && (item.Status == models.StatusDownloading || item.Status == models.StatusPending) { newItems[id] = item newOrder = append(newOrder, id) } } u.items = newItems u.order = newOrder } u.mu.Unlock() } func (u *UserQueue) GetItem(id string) *models.QueueItem { if id == "" { return nil } u.mu.RLock() defer u.mu.RUnlock() if item, exists := u.items[id]; exists { cpy := *item return &cpy } return nil } func (u *UserQueue) StartDownloads() { u.mu.Lock() for _, id := range u.order { if item, ok := u.items[id]; ok && item.Status == models.StatusStaged { item.Status = models.StatusPending } } u.broadcastProgressLocked() if u.isProcessing { u.mu.Unlock() return } u.isProcessing = true u.mu.Unlock() go u.processQueueLoop() } func (u *UserQueue) processQueueLoop() { defer func() { u.mu.Lock() u.isProcessing = false u.broadcastProgressLocked() u.mu.Unlock() }() stopTicker := make(chan struct{}) go func() { ticker := time.NewTicker(500 * time.Millisecond) defer ticker.Stop() for { select { case <-ticker.C: u.tick() case <-stopTicker: return } } }() defer close(stopTicker) for { select { case <-u.ctx.Done(): return default: } u.mu.Lock() var nextID string var requestedBy string for _, id := range u.order { if item, ok := u.items[id]; ok && item.Status == models.StatusPending { nextID = id requestedBy = item.RequestedBy item.Status = models.StatusDownloading u.currentAlbumID = id break } } u.broadcastProgressLocked() u.mu.Unlock() if nextID == "" { break } u.logger.Info("starting download", "id", nextID, "requested_by", requestedBy) l := u.lockRegistry.GetLock(nextID) l.Lock() downloadCtx := downloader.WithRequestedBy(u.ctx, requestedBy) downloadCtx = downloader.WithProgressListener(downloadCtx, u) downloadCtx = log.IntoContext(downloadCtx, u.logger) if err := u.dl.DownloadAlbum(downloadCtx, nextID); err != nil { u.logger.Error("album download failed", "id", nextID, "err", err) } l.Unlock() var onComplete func() u.mu.Lock() if item, ok := u.items[nextID]; ok { allDone := true hasFailed := false if len(item.Tracks) == 0 && item.TracksCount > 0 { allDone = false hasFailed = true } else { for _, t := range item.Tracks { if t.Status == models.StatusFailed { hasFailed = true } if t.Status != models.StatusDone { allDone = false } } } if hasFailed || !allDone || (item.TracksCount > 0 && len(item.Tracks) < item.TracksCount) { item.Status = models.StatusFailed u.logger.Warn("download failed", "id", nextID, "requested_by", requestedBy) } else { item.Status = models.StatusDone item.Percent = 100 for _, t := range item.Tracks { t.Status = models.StatusDone t.Percent = 100 } u.logger.Info("download completed", "id", nextID, "requested_by", requestedBy, "tracks", len(item.Tracks)) onComplete = u.onComplete } } u.broadcastProgressLocked() u.mu.Unlock() if onComplete != nil { onComplete() } } } func (u *UserQueue) getViewModelLocked() pages.QueueViewModel { items := make([]*models.QueueItem, len(u.order)) stagedCount := 0 isDownloading := false for i, id := range u.order { item := u.items[id] items[i] = item switch item.Status { case models.StatusStaged: stagedCount++ case models.StatusDownloading, models.StatusPending: isDownloading = true } } return pages.QueueViewModel{ Items: items, StagedCount: stagedCount, Downloading: isDownloading, } } func (u *UserQueue) GetViewModel() pages.QueueViewModel { u.mu.RLock() defer u.mu.RUnlock() return u.getViewModelLocked() } func (u *UserQueue) broadcastProgressLocked() { vm := u.getViewModelLocked() var buf bytes.Buffer if err := u.pages.RenderQueue(&buf, vm); err == nil { u.broker.Broadcast(u.userID, "queue-progress", buf.String()) } }