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), } } // Start discards anything a previous run of the same id left behind. func (c *ProgressCache) Start(id int64) { c.mu.Lock() defer c.mu.Unlock() c.data[id] = &liveDownload{} } 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 reads under the lock so it never races 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 what was appended since the last call and advances the // mark, so the caller appends to the DB instead of rewriting the whole buffer. 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 }