Repository navigation
Expand file tree
/
Copy pathstats_test.go
More file actions
314 lines (288 loc) · 9.65 KB
/
Copy pathstats_test.go
File metadata and controls
314 lines (288 loc) · 9.65 KB
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
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
package main
import (
"bytes"
"os"
"path/filepath"
"strings"
"sync"
"testing"
"time"
"github.com/floatdrop/moq-go/pkg/moqt/msf"
)
// at builds an arrival at a fixed base time plus offset milliseconds, so
// the ordering assertions do not depend on a clock.
func at(group, object uint64, ms int) arrival {
base := time.Unix(1_700_000_000, 0)
return arrival{
Group: group,
Object: object,
Bytes: 10,
At: base.Add(time.Duration(ms) * time.Millisecond),
}
}
func TestRecorderCountsObjectsThatArriveAfterALaterOne(t *testing.T) {
rec := &recorder{}
// Group 1 overtakes the tail of Group 0 — which is what two subgroup
// streams read on their own goroutines can do, and what a player has
// to survive.
for _, a := range []arrival{
at(0, 0, 0), at(0, 1, 10), at(1, 0, 20), at(0, 2, 30), at(1, 1, 40),
} {
rec.add(a)
}
if rec.outOfOrder != 1 {
t.Errorf("outOfOrder = %d, want 1 (object 0/2 arrived after 1/0)", rec.outOfOrder)
}
// sorted() must undo the interleaving: it is what the reassembled file
// is written from, so a wrong order here is a corrupt output file.
sorted := rec.sorted()
want := []arrival{at(0, 0, 0), at(0, 1, 10), at(0, 2, 30), at(1, 0, 20), at(1, 1, 40)}
for i := range want {
if sorted[i].Group != want[i].Group || sorted[i].Object != want[i].Object {
t.Fatalf("sorted[%d] = %d/%d, want %d/%d",
i, sorted[i].Group, sorted[i].Object, want[i].Group, want[i].Object)
}
}
}
func TestRecorderDoesNotCountInOrderArrivalsAsReordered(t *testing.T) {
rec := &recorder{}
for g := range uint64(3) {
for o := range uint64(4) {
rec.add(at(g, o, int(g*4+o)))
}
}
if rec.outOfOrder != 0 {
t.Errorf("outOfOrder = %d, want 0", rec.outOfOrder)
}
if rec.bytes != 12*10 {
t.Errorf("bytes = %d, want %d", rec.bytes, 12*10)
}
}
func TestGaps(t *testing.T) {
tests := []struct {
name string
arrivals []arrival
wantObjects uint64
wantGroups uint64
}{{
name: "nothing missing",
arrivals: []arrival{at(0, 0, 0), at(0, 1, 1), at(1, 0, 2), at(1, 1, 3)},
}, {
name: "an object lost mid-group",
arrivals: []arrival{at(0, 0, 0), at(0, 2, 1)},
wantObjects: 1,
}, {
name: "a whole group lost",
arrivals: []arrival{at(0, 0, 0), at(2, 0, 1)},
wantGroups: 1,
}, {
name: "a group joined after its start",
// Object 0 and 1 of group 1 never arrived, which is loss even
// though nothing between two received objects is missing.
arrivals: []arrival{at(0, 0, 0), at(1, 2, 1)},
wantObjects: 2,
}, {
name: "joined mid-broadcast",
// The subscriber's first object is 7/3. Everything before it was
// never sent to this subscriber, and counting it as loss would
// make every late join look like a delivery failure.
arrivals: []arrival{at(7, 3, 0), at(7, 4, 1), at(8, 0, 2)},
}}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
objects, groups := gaps(tc.arrivals)
if objects != tc.wantObjects {
t.Errorf("missing objects = %d, want %d", objects, tc.wantObjects)
}
if groups != tc.wantGroups {
t.Errorf("missing groups = %d, want %d", groups, tc.wantGroups)
}
})
}
}
func TestWriteMediaConcatenatesInSendOrder(t *testing.T) {
path := filepath.Join(t.TempDir(), "out.mp4")
init := []byte("INIT")
// Deliberately out of arrival order: writeMedia is handed the sorted
// view, so the file must come out in send order regardless.
sorted := []arrival{
{Group: 0, Object: 0, Payload: []byte("a")},
{Group: 0, Object: 1, Payload: []byte("b")},
{Group: 1, Object: 0, Payload: []byte("c")},
}
digest, err := writeMedia(path, init, sorted)
if err != nil {
t.Fatalf("writeMedia: %v", err)
}
got, err := os.ReadFile(path)
if err != nil {
t.Fatalf("read back: %v", err)
}
if !bytes.Equal(got, []byte("INITabc")) {
t.Errorf("file = %q, want %q", got, "INITabc")
}
// The digest is what the source is compared against, so it must cover
// the header as well as the payloads.
if digest != sha256Of("INITabc") {
t.Errorf("digest = %s, want the digest of the whole file", digest)
}
}
func TestWriteMediaSkipsWhenNoPathIsGiven(t *testing.T) {
digest, err := writeMedia("", []byte("INIT"), []arrival{{Payload: []byte("a")}})
if err != nil {
t.Fatalf("writeMedia: %v", err)
}
if digest != "" {
t.Errorf("digest = %q, want empty: with no file there is nothing to hash", digest)
}
}
func TestReportComparesAgainstTheSourceDigest(t *testing.T) {
rec := &recorder{}
rec.add(at(0, 0, 0))
rec.add(at(0, 1, 10))
var buf bytes.Buffer
rec.report(&buf, rec.sorted(), broadcast{Digest: "abc", Objects: 2, Bytes: 20}, "abc")
if !strings.Contains(buf.String(), "MATCHES the source") {
t.Errorf("matching digests not reported as a match:\n%s", buf.String())
}
buf.Reset()
rec.report(&buf, rec.sorted(), broadcast{Digest: "abc", Objects: 2, Bytes: 20}, "def")
if !strings.Contains(buf.String(), "DIFFERS from the source") {
t.Errorf("differing digests not reported as a mismatch:\n%s", buf.String())
}
}
func TestReportSaysSoWhenNothingArrived(t *testing.T) {
var buf bytes.Buffer
(&recorder{}).report(&buf, nil, broadcast{}, "")
if !strings.Contains(buf.String(), "no objects received") {
t.Errorf("empty run not reported:\n%s", buf.String())
}
}
func TestPercentile(t *testing.T) {
sorted := []time.Duration{1, 2, 3, 4, 5, 6, 7, 8, 9, 10}
for _, tc := range []struct {
p float64
want time.Duration
}{{0, 1}, {0.5, 5}, {0.9, 9}, {1, 10}} {
if got := percentile(sorted, tc.p); got != tc.want {
t.Errorf("percentile(%.2f) = %v, want %v", tc.p, got, tc.want)
}
}
}
// TestRecorderSpanSurvivesOutOfOrderTimestamps covers what concurrent
// readers make reachable: one Group's reader can stamp an arrival and then
// lose the race for the lock to a Group that landed later, so add sees
// timestamps out of order. Assigning first/last blindly would leave last
// before first, and the report's span and mean would come out negative.
func TestRecorderSpanSurvivesOutOfOrderTimestamps(t *testing.T) {
rec := &recorder{}
rec.add(at(0, 0, 50))
rec.add(at(1, 0, 10)) // stamped earlier, recorded later
rec.add(at(0, 1, 90))
rec.add(at(1, 1, 30))
if rec.last.Before(rec.first) {
t.Fatalf("last %v is before first %v: the span is negative", rec.last, rec.first)
}
if got := rec.last.Sub(rec.first); got != 80*time.Millisecond {
t.Errorf("span = %v, want 80ms (the extremes, not the last two seen)", got)
}
}
// TestDrainedReportsWhetherReadersFinished pins the bound that stops a
// still-arriving Group from being reported as loss — and the escape hatch
// that stops a reader abandoned mid-Group from holding the report forever.
func TestDrainedReportsWhetherReadersFinished(t *testing.T) {
var done sync.WaitGroup
done.Go(func() { time.Sleep(10 * time.Millisecond) })
if !drained(&done, time.Second) {
t.Error("drained reported a timeout for a reader that finished in time")
}
var stuck sync.WaitGroup
stuck.Add(1) // never released, as an abandoned stream's reader would be
if drained(&stuck, 20*time.Millisecond) {
t.Error("drained reported success for a reader that never finished")
}
}
// TestFinishRunWaitsForInFlightGroups is the regression test for a report
// that invented loss.
//
// Groups are read concurrently and the terminator catalog is a stream of
// its own, so the end-of-run signal can arrive while Objects are still
// coming in. A run that counted straight from that signal printed the
// undelivered tail as missing and wrote a truncated file — a delivery
// failure produced by the tool rather than found by it.
func TestFinishRunWaitsForInFlightGroups(t *testing.T) {
rec := &recorder{}
rec.add(at(0, 0, 0))
// A reader still draining its Group when the run is told to end.
var readers sync.WaitGroup
readers.Go(func() {
time.Sleep(30 * time.Millisecond)
rec.add(at(0, 1, 30))
})
var buf bytes.Buffer
sorted, err := finishRun(
t.Context(),
rec,
&readers,
"",
broadcast{},
mediaWriterFor(msf.PackagingCMAF, nil),
nil,
&buf,
)
if err != nil {
t.Fatalf("finishRun: %v", err)
}
if len(sorted) != 2 {
t.Fatalf("reported %d objects, want 2: the late one was counted as lost", len(sorted))
}
if strings.Contains(buf.String(), "1 missing") {
t.Errorf("report claims loss for an object that was still arriving:\n%s", buf.String())
}
}
// TestFinishRunStillReportsWhenAReaderNeverFinishes covers the other half:
// the bound exists so an abandoned stream cannot hold the report forever.
func TestFinishRunStillReportsWhenAReaderNeverFinishes(t *testing.T) {
rec := &recorder{}
rec.add(at(0, 0, 0))
var readers sync.WaitGroup
readers.Add(1) // never released, as a reader on an abandoned stream
defer readers.Done()
done := make(chan struct{})
var buf bytes.Buffer
go func() {
defer close(done)
if _, err := finishRun(
t.Context(),
rec,
&readers,
"",
broadcast{},
mediaWriterFor(msf.PackagingCMAF, nil),
nil,
&buf,
); err != nil {
t.Errorf("finishRun: %v", err)
}
}()
select {
case <-done:
if !strings.Contains(buf.String(), "delivery report") {
t.Errorf("no report was written:\n%s", buf.String())
}
case <-time.After(drainWait + 5*time.Second):
t.Fatal("finishRun never returned: the drain is unbounded")
}
}
// TestMediaStdoutKeepsTheReportOffThePipe covers -out -, which exists so
// the media can be piped into a player. Everything else the run says has to
// move aside, or the first thing down the pipe is a delivery report and no
// decoder will make sense of the stream.
func TestMediaStdoutKeepsTheReportOffThePipe(t *testing.T) {
if got := reportWriter(mediaStdout); got != os.Stderr {
t.Error("with the media on stdout the report must go to stderr")
}
if got := reportWriter("/tmp/out.mp4"); got != os.Stdout {
t.Error("with the media in a file the report belongs on stdout")
}
}