progress_cache.go
⎇
Raw
1package service
2
3import (
4 "strings"
5 "sync"
6)
7
8type liveDownload struct {
9 logs strings.Builder
10 flushed int // length of logs already persisted to the DB
11}
12
13type ProgressCache struct {
14 mu sync.RWMutex
15 data map[int64]*liveDownload
16}
17
18func NewProgressCache() *ProgressCache {
19 return &ProgressCache{
20 data: make(map[int64]*liveDownload),
21 }
22}
23
24// Start begins buffering live output for a download, discarding anything a
25// previous run of the same id left behind.
26func (c *ProgressCache) Start(id int64) {
27 c.mu.Lock()
28 defer c.mu.Unlock()
29 c.data[id] = &liveDownload{}
30}
31
32func (c *ProgressCache) AppendLog(id int64, line string) {
33 c.mu.Lock()
34 defer c.mu.Unlock()
35 if d, ok := c.data[id]; ok {
36 d.logs.WriteString(line)
37 d.logs.WriteByte('\n')
38 }
39}
40
41// Snapshot returns the full in-memory log buffer for a live download, read under
42// the lock so it never races the writer in AppendLog (strings.Builder is not
43// safe for concurrent read/write). Empty string if the download isn't live.
44func (c *ProgressCache) Snapshot(id int64) string {
45 c.mu.RLock()
46 defer c.mu.RUnlock()
47 if d, ok := c.data[id]; ok {
48 return d.logs.String()
49 }
50 return ""
51}
52
53func (c *ProgressCache) Delete(id int64) {
54 c.mu.Lock()
55 defer c.mu.Unlock()
56 delete(c.data, id)
57}
58
59// FlushLogs returns only the log content appended since the last flush and
60// advances the flushed mark, so the caller can append (rather than rewrite the
61// whole buffer) to the DB. Empty string when there's nothing new.
62func (c *ProgressCache) FlushLogs(id int64) string {
63 c.mu.Lock()
64 defer c.mu.Unlock()
65 d, ok := c.data[id]
66 if !ok {
67 return ""
68 }
69 full := d.logs.String()
70 if d.flushed >= len(full) {
71 return ""
72 }
73 tail := full[d.flushed:]
74 d.flushed = len(full)
75 return tail
76}
77