package service import ( "strings" "sync" ) type LiveDownload struct { Logs strings.Builder flushed int // length of Logs already persisted to the DB } type ProgressCache struct { mu sync.RWMutex data map[int64]*LiveDownload } func NewProgressCache() *ProgressCache { return &ProgressCache{ data: make(map[int64]*LiveDownload), } } func (c *ProgressCache) Set(id int64, d *LiveDownload) { c.mu.Lock() defer c.mu.Unlock() c.data[id] = d } func (c *ProgressCache) AppendLog(id int64, line string) { c.mu.Lock() defer c.mu.Unlock() if d, ok := c.data[id]; ok { d.Logs.WriteString(line) d.Logs.WriteByte('\n') } } // Snapshot returns the full in-memory log buffer for a live download, read under // the lock so it never races the writer in AppendLog (strings.Builder is not // safe for concurrent read/write). Empty string if the download isn't live. func (c *ProgressCache) Snapshot(id int64) string { c.mu.RLock() defer c.mu.RUnlock() if d, ok := c.data[id]; ok { return d.Logs.String() } return "" } func (c *ProgressCache) Delete(id int64) { c.mu.Lock() defer c.mu.Unlock() delete(c.data, id) } // FlushLogs returns only the log content appended since the last flush and // advances the flushed mark, so the caller can append (rather than rewrite the // whole buffer) to the DB. Empty string when there's nothing new. func (c *ProgressCache) FlushLogs(id int64) string { c.mu.Lock() defer c.mu.Unlock() d, ok := c.data[id] if !ok { return "" } full := d.Logs.String() if d.flushed >= len(full) { return "" } tail := full[d.flushed:] d.flushed = len(full) return tail }