1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
//! The bounded, drop-oldest progress bus.
//!
//! Why: `tga collect` / `tga classify` run for minutes against large corpora,
//! and #5197's TUI needs a live view of what they are doing. The producer is
//! the pipeline, whose throughput must not depend on whether anybody is
//! watching — so delivery is non-blocking and lossy by construction.
//! What: [`ProgressBus`], a cheap clonable handle over a shared ring buffer.
//! An inactive bus ([`ProgressBus::disabled`], also `Default`) drops every
//! emit on the floor, which is what every existing CLI path passes so its
//! behavior — including the current `indicatif` bars — is unchanged.
//! Test: `super::tests` covers the no-subscriber path, the drop-oldest
//! overflow policy, the dropped counter, and drain ordering.
use VecDeque;
use ;
use ;
use ProgressEvent;
/// Default ring-buffer capacity for [`ProgressBus::bounded`] callers that have
/// no reason to pick their own.
///
/// Why: 1024 events is several seconds of headroom at the emit rates the
/// pipelines actually produce (per repo, per batch — not per commit), so a TUI
/// redrawing at 10 Hz never loses anything in practice.
/// What: `1024`.
/// Test: `super::tests::default_capacity_is_used_by_new`.
pub const DEFAULT_CAPACITY: usize = 1024;
/// Shared state behind an active bus. Never exposed.
/// A non-blocking, bounded, drop-oldest channel for [`ProgressEvent`]s.
///
/// Why: progress is advisory. A pipeline must never stall, fail, or change its
/// results because a consumer is slow, absent, or has gone away — so the bus
/// never applies backpressure and never returns an error the producer has to
/// handle. The cost of that choice is that events can be lost; the bus counts
/// them ([`ProgressBus::dropped`]) so a consumer can say so rather than
/// silently rendering a gap.
///
/// What: a clonable handle. [`ProgressBus::disabled`] (and `Default`) yields
/// an *inactive* bus whose [`ProgressBus::emit`] is a no-op and whose
/// [`ProgressBus::is_active`] is `false` — this is what every non-TUI call
/// site passes. [`ProgressBus::bounded`] yields an active bus holding a ring
/// buffer of `capacity` events; when it is full, the OLDEST event is evicted
/// to make room for the newest. Drop-oldest (rather than drop-newest) is
/// deliberate: for a live display the most recent state is the useful one, and
/// a stale head is exactly what a viewer does not need.
///
/// Consumers call [`ProgressBus::drain`] on their own cadence — typically once
/// per render tick. Producers and consumers hold the internal lock only for
/// the length of a push or a `VecDeque` swap, so neither ever waits on the
/// other's real work.
///
/// Test: `super::tests::disabled_bus_swallows_every_emit`,
/// `overflow_drops_oldest_and_counts`, `drain_returns_fifo_and_empties`.