返回 DeepSeek-Reasonix
budget.go
根目录 / internal / historywork / budget.go
1 // Package historywork bounds reconstructable history work independently of
2 // persistence and runtime ownership. Cancelling a reader never cancels a save.
3 package historywork
4
5 import (
6 "context"
7 "errors"
8 "io"
9 "sync"
10 "sync/atomic"
11 "time"
12 )
13
14 const (
15 BatchEntries = 128
16 BatchBytes = 4 << 20
17 ReadChunk = 64 << 10
18 BytesPerSecond = 8 << 20
19 SliceDuration = 50 * time.Millisecond
20 PauseDuration = 100 * time.Millisecond
21 )
22
23 // Coordinator is shared by the catalog and historical source discovery.
24 // Background permits are held for one bounded slice, never a directory scan.
25 type Coordinator struct {
26 once sync.Once
27 background chan struct{}
28 foreground chan struct{}
29 mu sync.Mutex
30 readers int
31 changed chan struct{}
32 backgroundReady time.Time
33 readBytes atomic.Int64
34 readCalls atomic.Uint64
35 canceledReads atomic.Uint64
36 slices atomic.Uint64
37 lastSliceNanos atomic.Int64
38 }
39
40 func (c *Coordinator) init() {
41 c.once.Do(func() {
42 c.background = make(chan struct{}, 1)
43 c.foreground = make(chan struct{}, 1)
44 c.changed = make(chan struct{})
45 })
46 }
47
48 func (c *Coordinator) Foreground(ctx context.Context) (func(), error) {
49 c.init()
50 if err := ctx.Err(); err != nil {
51 return nil, err
52 }
53 select {
54 case c.foreground <- struct{}{}:
55 case <-ctx.Done():
56 return nil, ctx.Err()
57 }
58 if err := ctx.Err(); err != nil {
59 <-c.foreground
60 return nil, err
61 }
62 c.mu.Lock()
63 c.readers++
64 c.mu.Unlock()
65 var once sync.Once
66 return func() {
67 once.Do(func() {
68 c.mu.Lock()
69 c.readers--
70 close(c.changed)
71 c.changed = make(chan struct{})
72 c.mu.Unlock()
73 <-c.foreground
74 })
75 }, nil
76 }
77
78 func (c *Coordinator) Background(ctx context.Context) (func(), error) {
79 release, err := c.BackgroundSlice(ctx, false)
80 if err != nil {
81 return nil, err
82 }
83 return func() { release(0) }, nil
84 }
85
86 func (c *Coordinator) ForegroundActive() bool {
87 c.init()
88 c.mu.Lock()
89 defer c.mu.Unlock()
90 return c.readers > 0
91 }
92
93 // BackgroundSlice applies one shared rate budget across all discovery owners.
94 // P1 metadata may proceed during a foreground preparation; P2 waits for it.
95 func (c *Coordinator) BackgroundSlice(ctx context.Context, foregroundAllowed bool) (func(int64), error) {
96 return c.backgroundSlice(ctx, foregroundAllowed, false)
97 }
98
99 var ErrForegroundActive = errors.New("foreground history preparation active")
100
101 // A single scheduling loop must be able to reconsider P1 while P2 is paused.
102 // Yielding preserves the caller's iterator; this is not a failed slice.
103 func (c *Coordinator) BackgroundSliceYielding(ctx context.Context, foregroundAllowed bool) (func(int64), error) {
104 return c.backgroundSlice(ctx, foregroundAllowed, true)
105 }
106
107 func (c *Coordinator) backgroundSlice(ctx context.Context, foregroundAllowed, yield bool) (func(int64), error) {
108 c.init()
109 for {
110 if err := ctx.Err(); err != nil {
111 return nil, err
112 }
113 c.mu.Lock()
114 busy, changed, ready := c.readers > 0 && !foregroundAllowed, c.changed, c.backgroundReady
115 c.mu.Unlock()
116 if busy {
117 if yield {
118 return nil, ErrForegroundActive
119 }
120 select {
121 case <-changed:
122 continue
123 case <-ctx.Done():
124 return nil, ctx.Err()
125 }
126 }
127 if delay := time.Until(ready); delay > 0 {
128 timer := time.NewTimer(delay)
129 select {
130 case <-ctx.Done():
131 timer.Stop()
132 return nil, ctx.Err()
133 case <-timer.C:
134 }
135 }
136 select {
137 case c.background <- struct{}{}:
138 case <-ctx.Done():
139 return nil, ctx.Err()
140 }
141 c.mu.Lock()
142 foregroundBusy := c.readers > 0 && !foregroundAllowed
143 busy = foregroundBusy || time.Now().Before(c.backgroundReady)
144 c.mu.Unlock()
145 if busy {
146 <-c.background
147 if foregroundBusy && yield {
148 return nil, ErrForegroundActive
149 }
150 continue
151 }
152 var once sync.Once
153 started := time.Now()
154 return func(bytes int64) {
155 once.Do(func() {
156 c.slices.Add(1)
157 c.lastSliceNanos.Store(time.Since(started).Nanoseconds())
158 c.mu.Lock()
159 c.backgroundReady = time.Now().Add(max(PauseDuration, time.Duration(max(0, bytes))*time.Second/BytesPerSecond))
160 c.mu.Unlock()
161 <-c.background
162 })
163 }, nil
164 }
165 }
166
167 // Pause includes both the per-slice yield and the byte-rate allowance. Waiting
168 // is cancellable; tests may inject a clock through a caller's slice runner.
169 func Pause(ctx context.Context, bytes int64) error {
170 delay := max(PauseDuration, time.Duration(bytes)*time.Second/BytesPerSecond)
171 timer := time.NewTimer(delay)
172 defer timer.Stop()
173 select {
174 case <-timer.C:
175 return ctx.Err()
176 case <-ctx.Done():
177 return ctx.Err()
178 }
179 }
180
181 // Reader checks cancellation before each bounded source read, including reads
182 // performed inside a JSON decoder rather than only between decoded records.
183 type Reader struct {
184 Context context.Context
185 Source io.Reader
186 Bytes int64
187 }
188
189 func (r *Reader) Read(p []byte) (int, error) {
190 meter, _ := r.Context.Value(meterKey{}).(*Coordinator)
191 if err := r.Context.Err(); err != nil {
192 if meter != nil {
193 meter.canceledReads.Add(1)
194 }
195 return 0, err
196 }
197 if len(p) > ReadChunk {
198 p = p[:ReadChunk]
199 }
200 n, err := r.Source.Read(p)
201 r.Bytes += int64(n)
202 if meter != nil {
203 meter.readBytes.Add(int64(n))
204 meter.readCalls.Add(1)
205 }
206 return n, err
207 }
208
209 type meterKey struct{}
210
211 func (c *Coordinator) Context(ctx context.Context) context.Context {
212 return context.WithValue(ctx, meterKey{}, c)
213 }
214
215 // Counters contain no source identities or transcript text. Bytes counts only
216 // reads routed through Reader, rather than pretending a stat size was read.
217 type Diagnostics struct {
218 ForegroundActive int `json:"foregroundActive"`
219 BackgroundActive int `json:"backgroundActive"`
220 BackgroundSlices uint64 `json:"backgroundSlices"`
221 InstrumentedReadBytes int64 `json:"instrumentedReadBytes"`
222 InstrumentedReadCalls uint64 `json:"instrumentedReadCalls"`
223 CanceledReadCheckpoints uint64 `json:"canceledReadCheckpoints"`
224 LastBackgroundSliceMS float64 `json:"lastBackgroundSliceMs"`
225 }
226
227 func (c *Coordinator) Diagnostics() Diagnostics {
228 c.init()
229 c.mu.Lock()
230 defer c.mu.Unlock()
231 return Diagnostics{ForegroundActive: c.readers, BackgroundActive: len(c.background), BackgroundSlices: c.slices.Load(), InstrumentedReadBytes: c.readBytes.Load(), InstrumentedReadCalls: c.readCalls.Load(), CanceledReadCheckpoints: c.canceledReads.Load(), LastBackgroundSliceMS: float64(c.lastSliceNanos.Load()) / float64(time.Millisecond)}
232 }
233
233 lines GO