Repository navigation
Expand file tree
/
Copy pathmain.go
More file actions
129 lines (118 loc) · 4.72 KB
/
Copy pathmain.go
File metadata and controls
129 lines (118 loc) · 4.72 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
// Command video streams a local video file to a relay as a CMSF
// (CMAF-packaged) broadcast, and measures what comes back out of it.
//
// It exists to answer one question about a suspected delivery fault: is
// the transport at fault, or the encoder feeding it? Every byte the
// publisher sends is known in advance, so the subscriber can say exactly
// what arrived, in what order and how late, and can reassemble the
// Objects it received into a file that either matches the source byte for
// byte or does not. A clean run points at the capture and encode path; a
// dirty one points here.
//
// Objects are CMAF chunks of one frame each, and Groups start at sync
// samples (draft-ietf-moq-cmsf-01 §3.3, §3.4), so per-Object timings are
// per-frame timings. Only the file's first video track is published;
// audio, subtitle and data tracks are ignored.
package main
import (
"context"
"flag"
"fmt"
"log/slog"
"os"
"os/signal"
"syscall"
"time"
"github.com/lmittmann/tint"
dialpkg "github.com/floatdrop/moq-go/internal/dial"
"github.com/floatdrop/moq-go/pkg/moqt/msf"
"github.com/floatdrop/moq-go/pkg/moqt/session"
)
// defaultNamespace is the namespace both modes use unless -ns says
// otherwise.
const defaultNamespace = "moq-example/video"
func main() {
addr := flag.String("addr", "localhost:4433", "relay address (host:port or a moqt:// URI)")
namespace := flag.String("ns", defaultNamespace, "track namespace to publish under / subscribe to")
in := flag.String("in", "", "publish: path of the video file to stream (required)")
out := flag.String("out", "", "subscribe: path to write the received media to")
rate := flag.Float64("rate", 1, "publish: pacing multiplier; 0 sends as fast as the transport takes it")
loops := flag.Int("loop", 1, "publish: passes over the file; 0 repeats until interrupted")
gop := flag.Int("gop", 0, "publish: minimum Objects per Group; 0 starts a Group at every sync sample")
delay := flag.Duration("delay", 2*time.Second,
"publish: wait between the catalog and the first frame, so a subscriber can be in place")
wait := flag.Duration("wait", 30*time.Second,
"subscribe: how long to wait for a publisher on the namespace before giving up")
packaging := flag.String("packaging", msf.PackagingCMAF,
`subscribe: packaging of the track to read — "cmaf" for a broadcast this tool published, `+
`"legacy" for the bare-bitstream tracks moq-lite/hang serves`)
flag.Parse()
if flag.NArg() < 1 {
fmt.Fprintln(os.Stderr, "usage: video [flags] publish|subscribe")
flag.PrintDefaults()
os.Exit(1)
}
slog.SetDefault(slog.New(tint.NewTextHandler(os.Stderr, &tint.Options{
Level: slog.LevelInfo,
TimeFormat: time.TimeOnly,
})))
// Keep SIGPIPE from killing the run. Go's runtime turns a broken-pipe
// write to stdout or stderr into a fatal signal unless the program has
// asked for it — and the documented way to watch a stream is to pipe
// -out - into a player, whose quitting is a normal end. Left alone, the
// delivery report on stderr would never print, which is the one thing a
// run exists to produce. Notified and ignored, the write returns EPIPE
// and the wind-down reports it.
signal.Ignore(syscall.SIGPIPE)
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
go func() {
<-ctx.Done()
slog.Info("signal received, shutting down (Ctrl+C again to force-exit)")
stop()
}()
var err error
switch flag.Arg(0) {
case "publish":
if *in == "" {
fmt.Fprintln(os.Stderr, "publish needs -in <file.mp4>")
os.Exit(1)
}
err = publish(ctx, *addr, publishOptions{
Namespace: *namespace,
File: *in,
Rate: *rate,
Loops: *loops,
MinGroupObjects: *gop,
Delay: *delay,
})
case "subscribe":
if *packaging != msf.PackagingCMAF && *packaging != legacyPackaging {
fmt.Fprintf(os.Stderr, "unknown -packaging %q (want %q or %q)\n",
*packaging, msf.PackagingCMAF, legacyPackaging)
os.Exit(1)
}
err = subscribe(ctx, *addr, subscribeOptions{
Namespace: *namespace,
Out: *out,
Wait: *wait,
Packaging: *packaging,
})
default:
fmt.Fprintf(os.Stderr, "unknown mode %q (want publish or subscribe)\n", flag.Arg(0))
os.Exit(1)
}
if err != nil {
slog.Error("fatal", tint.Err(err))
os.Exit(1)
}
}
// dial establishes a QUIC connection and completes the MOQT client
// handshake. addr may be a bare host:port or a moqt:// URI whose authority
// and path/query ride the AUTHORITY / PATH Setup Options (see
// internal/dial).
func dial(ctx context.Context, addr string) (*session.Session, error) {
return dialpkg.QUIC(ctx, addr, dialpkg.Options{
Implementation: "video/0.1",
InsecureSkipVerify: true, // dev-only debug client; certs not verified by design
})
}