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 discards anything a previous run of the same id left behind.
25func (c *ProgressCache) Start(id int64) {
26 c.mu.Lock()
27 defer c.mu.Unlock()
28 c.data[id] = &liveDownload{}
29}
30
31func (c *ProgressCache) AppendLog(id int64, line string) {
32 c.mu.Lock()
33 defer c.mu.Unlock()
34 if d, ok := c.data[id]; ok {
35 d.logs.WriteString(line)
36 d.logs.WriteByte('\n')
37 }
38}
39
40// Snapshot reads under the lock so it never races AppendLog: strings.Builder is
41// not safe for concurrent read/write. Empty string if the download isn't live.
42func (c *ProgressCache) Snapshot(id int64) string {
43 c.mu.RLock()
44 defer c.mu.RUnlock()
45 if d, ok := c.data[id]; ok {
46 return d.logs.String()
47 }
48 return ""
49}
50
51func (c *ProgressCache) Delete(id int64) {
52 c.mu.Lock()
53 defer c.mu.Unlock()
54 delete(c.data, id)
55}
56
57// FlushLogs returns only what was appended since the last call and advances the
58// mark, so the caller appends to the DB instead of rewriting the whole buffer.
59func (c *ProgressCache) FlushLogs(id int64) string {
60 c.mu.Lock()
61 defer c.mu.Unlock()
62 d, ok := c.data[id]
63 if !ok {
64 return ""
65 }
66 full := d.logs.String()
67 if d.flushed >= len(full) {
68 return ""
69 }
70 tail := full[d.flushed:]
71 d.flushed = len(full)
72 return tail
73}
74