Compare commits
4 Commits
27dcbd1b40
...
aac8c664d8
| Author | SHA1 | Date | |
|---|---|---|---|
| aac8c664d8 | |||
| 05f356fa10 | |||
| 9b6607fc48 | |||
| 774361df47 |
+64
-29
@@ -75,26 +75,27 @@ func printUsage(w io.Writer) {
|
||||
}
|
||||
|
||||
func checkMXLargs(args appArgs) {
|
||||
if args.VideoFlowId == "" && args.AudioFlowId == "" {
|
||||
return
|
||||
}
|
||||
|
||||
if args.Domain == "" {
|
||||
fmt.Fprintln(
|
||||
os.Stderr,
|
||||
"You must provide a domain when a feed UUID is configured",
|
||||
)
|
||||
checkDomain := func(label, domain string) {
|
||||
if domain == "" {
|
||||
fmt.Fprintf(os.Stderr, "%s domain is required when its UUID is configured\n", label)
|
||||
printUsage(os.Stderr)
|
||||
os.Exit(2)
|
||||
}
|
||||
|
||||
fi, err := os.Stat(args.Domain)
|
||||
if err != nil || !fi.IsDir() {
|
||||
fmt.Fprintln(os.Stderr, "Invalid MXL domain:", args.Domain)
|
||||
info, err := os.Stat(domain)
|
||||
if err != nil || !info.IsDir() {
|
||||
fmt.Fprintf(os.Stderr, "Invalid %s MXL domain: %s\n", label, domain)
|
||||
fmt.Fprintln(os.Stderr, "Domain must be a directory in tmpfs")
|
||||
printUsage(os.Stderr)
|
||||
os.Exit(2)
|
||||
}
|
||||
}
|
||||
|
||||
if args.VideoFlowId != "" {
|
||||
checkDomain("video", args.VideoDomain)
|
||||
}
|
||||
if args.AudioFlowId != "" {
|
||||
checkDomain("audio", args.AudioDomain)
|
||||
}
|
||||
}
|
||||
|
||||
func main() {
|
||||
@@ -107,7 +108,15 @@ func main() {
|
||||
flagSet.SortFlags = false
|
||||
flagSet.Usage = func() { printUsage(os.Stderr) }
|
||||
flagSet.BoolVarP(&args.ShowHelp, "help", "h", false, "Show help message and exit")
|
||||
flagSet.StringVarP(&args.Domain, "domain", "d", "", "MXL domain directory")
|
||||
flagSet.StringVarP(
|
||||
&args.Domain,
|
||||
"domain",
|
||||
"d",
|
||||
"",
|
||||
"Default MXL domain for feeds without a specific domain",
|
||||
)
|
||||
flagSet.StringVar(&args.VideoDomain, "video-domain", "", "MXL domain for the video feed")
|
||||
flagSet.StringVar(&args.AudioDomain, "audio-domain", "", "MXL domain for the audio feed")
|
||||
flagSet.StringVarP(&args.VideoFlowId, "video", "v", "", "Video flow UUID")
|
||||
flagSet.StringVarP(&args.AudioFlowId, "audio", "a", "", "Audio flow UUID")
|
||||
flagSet.BoolVarP(
|
||||
@@ -152,9 +161,25 @@ func main() {
|
||||
fmt.Fprintln(os.Stderr, "invalid retry configuration:", err)
|
||||
os.Exit(2)
|
||||
}
|
||||
if args.VideoDomain == "" {
|
||||
args.VideoDomain = args.Domain
|
||||
}
|
||||
if args.AudioDomain == "" {
|
||||
args.AudioDomain = args.Domain
|
||||
}
|
||||
if !args.ListAudio && !args.ListGPU {
|
||||
checkMXLargs(args)
|
||||
}
|
||||
if args.SyncRequested &&
|
||||
args.VideoFlowId != "" &&
|
||||
args.AudioFlowId != "" &&
|
||||
args.VideoDomain != args.AudioDomain {
|
||||
fmt.Fprintln(
|
||||
os.Stderr,
|
||||
"--sync currently requires video and audio to use the same MXL domain",
|
||||
)
|
||||
os.Exit(2)
|
||||
}
|
||||
// path selection
|
||||
useLegacySync := args.SyncRequested &&
|
||||
args.VideoFlowId != "" &&
|
||||
@@ -292,7 +317,7 @@ func main() {
|
||||
switch {
|
||||
case useLegacySync:
|
||||
syncSrc, err = source.OpenSameDomainSync(
|
||||
args.Domain,
|
||||
args.VideoDomain,
|
||||
args.VideoFlowId,
|
||||
args.AudioFlowId,
|
||||
)
|
||||
@@ -378,7 +403,8 @@ func main() {
|
||||
|
||||
// GUI state (accessible from doReconnect + goroutine)
|
||||
var (
|
||||
domainStr string = args.Domain
|
||||
videoDomainStr string = args.VideoDomain
|
||||
audioDomainStr string = args.AudioDomain
|
||||
videoStr string = args.VideoFlowId
|
||||
audioStr string = args.AudioFlowId
|
||||
showStats bool = true
|
||||
@@ -552,7 +578,7 @@ func main() {
|
||||
videoConfig := playback.FeedConfig{}
|
||||
if videoStr != "" {
|
||||
videoConfig = playback.FeedConfig{
|
||||
Domain: domainStr,
|
||||
Domain: videoDomainStr,
|
||||
UUID: videoStr,
|
||||
Active: true,
|
||||
}
|
||||
@@ -561,7 +587,7 @@ func main() {
|
||||
audioConfig := playback.FeedConfig{}
|
||||
if audioStr != "" {
|
||||
audioConfig = playback.FeedConfig{
|
||||
Domain: domainStr,
|
||||
Domain: audioDomainStr,
|
||||
UUID: audioStr,
|
||||
Active: true,
|
||||
}
|
||||
@@ -572,11 +598,19 @@ func main() {
|
||||
return
|
||||
}
|
||||
// legacy
|
||||
if videoDomainStr != audioDomainStr {
|
||||
log.Printf(
|
||||
"sync reconnect rejected: video domain %q differs from audio domain %q",
|
||||
videoDomainStr,
|
||||
audioDomainStr,
|
||||
)
|
||||
return
|
||||
}
|
||||
select {
|
||||
case <-control:
|
||||
default:
|
||||
}
|
||||
control <- reconnectParams{domain: domainStr, video: videoStr, audio: audioStr}
|
||||
control <- reconnectParams{domain: videoDomainStr, video: videoStr, audio: audioStr}
|
||||
}
|
||||
|
||||
playbackDone := make(chan struct{})
|
||||
@@ -593,7 +627,7 @@ func main() {
|
||||
err := videoSlot.Run(
|
||||
ctx,
|
||||
playback.FeedConfig{
|
||||
Domain: args.Domain,
|
||||
Domain: args.VideoDomain,
|
||||
UUID: args.VideoFlowId,
|
||||
Active: args.VideoFlowId != "",
|
||||
},
|
||||
@@ -610,7 +644,7 @@ func main() {
|
||||
err := audioSlot.Run(
|
||||
ctx,
|
||||
playback.FeedConfig{
|
||||
Domain: args.Domain,
|
||||
Domain: args.AudioDomain,
|
||||
UUID: args.AudioFlowId,
|
||||
Active: args.AudioFlowId != "",
|
||||
},
|
||||
@@ -669,7 +703,7 @@ func main() {
|
||||
return
|
||||
}
|
||||
log.Printf("source: %v", err)
|
||||
params := reconnectParams{domain: domainStr, video: videoStr, audio: audioStr}
|
||||
params := reconnectParams{domain: videoDomainStr, video: videoStr, audio: audioStr}
|
||||
select {
|
||||
case <-control:
|
||||
default:
|
||||
@@ -731,7 +765,7 @@ func main() {
|
||||
log.Printf("source: %v", err)
|
||||
// Request a reconnect after the read failure.
|
||||
params := reconnectParams{
|
||||
domain: domainStr,
|
||||
domain: videoDomainStr,
|
||||
video: videoStr,
|
||||
audio: audioStr,
|
||||
}
|
||||
@@ -763,7 +797,7 @@ func main() {
|
||||
}
|
||||
log.Printf("source: %v", err)
|
||||
params := reconnectParams{
|
||||
domain: domainStr,
|
||||
domain: videoDomainStr,
|
||||
video: videoStr,
|
||||
audio: audioStr,
|
||||
}
|
||||
@@ -941,7 +975,8 @@ func main() {
|
||||
cimgui.End()
|
||||
}
|
||||
cimgui.Begin("Connection")
|
||||
cimgui.InputTextWithHint("Domain", "/dev/shm/mxl", &domainStr, 0, nil)
|
||||
cimgui.InputTextWithHint("Video domain", "/dev/shm/mxl", &videoDomainStr, 0, nil)
|
||||
cimgui.InputTextWithHint("Audio domain", "/dev/shm/mxl", &audioDomainStr, 0, nil)
|
||||
cimgui.InputTextWithHint("Video UUID", "", &videoStr, 0, nil)
|
||||
cimgui.InputTextWithHint("Audio UUID", "", &audioStr, 0, nil)
|
||||
if cimgui.Button("Connect") {
|
||||
@@ -954,7 +989,7 @@ func main() {
|
||||
videoActive = false
|
||||
enqueueVideoConfig(
|
||||
playback.FeedConfig{
|
||||
Domain: domainStr,
|
||||
Domain: videoDomainStr,
|
||||
UUID: videoStr,
|
||||
Active: false,
|
||||
})
|
||||
@@ -965,7 +1000,7 @@ func main() {
|
||||
if cimgui.Button("Resume video") {
|
||||
videoActive = true
|
||||
enqueueVideoConfig(playback.FeedConfig{
|
||||
Domain: domainStr,
|
||||
Domain: videoDomainStr,
|
||||
UUID: videoStr,
|
||||
Active: true,
|
||||
})
|
||||
@@ -1016,7 +1051,7 @@ func main() {
|
||||
if cimgui.Button("Stop audio") {
|
||||
audioActive = false
|
||||
enqueueAudioConfig(playback.FeedConfig{
|
||||
Domain: domainStr,
|
||||
Domain: audioDomainStr,
|
||||
UUID: audioStr,
|
||||
Active: false,
|
||||
})
|
||||
@@ -1027,7 +1062,7 @@ func main() {
|
||||
if cimgui.Button("Resume audio") {
|
||||
audioActive = true
|
||||
enqueueAudioConfig(playback.FeedConfig{
|
||||
Domain: domainStr,
|
||||
Domain: audioDomainStr,
|
||||
UUID: audioStr,
|
||||
Active: true,
|
||||
})
|
||||
|
||||
@@ -14,7 +14,7 @@ Size=200,200
|
||||
Collapsed=0
|
||||
|
||||
[Window][Connection]
|
||||
Pos=185,295
|
||||
Pos=177,324
|
||||
Size=640,354
|
||||
Collapsed=0
|
||||
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
package playback
|
||||
|
||||
import "context"
|
||||
|
||||
// SyncFrame contains one video frame and its corresponding audio batch.
|
||||
//
|
||||
// Both payloads may borrow reader-owned memory and are valid only until the
|
||||
// next ReadSync call or until the reader is closed.
|
||||
type SyncFrame struct {
|
||||
Video VideoFrame
|
||||
Audio AudioFrame
|
||||
}
|
||||
|
||||
// SyncReader reads synchronized audio/video pairs.
|
||||
//
|
||||
// ReadSync must not be called again until both frame payloads have been
|
||||
// consumed.
|
||||
type SyncReader interface {
|
||||
ReadSync(context.Context) (SyncFrame, error)
|
||||
Close() error
|
||||
}
|
||||
|
||||
// SyncReaderFactory opens one synchronized reader for two configured feeds.
|
||||
//
|
||||
// Native MXL implementations may require the feeds to use the same domain.
|
||||
// Future manual synchronization may support different domains behind another
|
||||
// implementation of this interface.
|
||||
type SyncReaderFactory interface {
|
||||
OpenSync(
|
||||
context.Context,
|
||||
FeedConfig,
|
||||
FeedConfig,
|
||||
) (SyncReader, error)
|
||||
}
|
||||
@@ -0,0 +1,56 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
)
|
||||
|
||||
func runSyncAttempt(
|
||||
ctx context.Context,
|
||||
factory SyncReaderFactory,
|
||||
videoSink VideoSink,
|
||||
audioSink AudioSink,
|
||||
videoConfig FeedConfig,
|
||||
audioConfig FeedConfig,
|
||||
) (resultErr error) {
|
||||
reader, err := factory.OpenSync(
|
||||
ctx,
|
||||
videoConfig,
|
||||
audioConfig,
|
||||
)
|
||||
if err != nil {
|
||||
return fmt.Errorf("open sync group: %w", err)
|
||||
}
|
||||
|
||||
defer func() {
|
||||
if closeErr := reader.Close(); closeErr != nil {
|
||||
closeErr = fmt.Errorf("close sync group: %w", closeErr)
|
||||
resultErr = errors.Join(resultErr, closeErr)
|
||||
}
|
||||
}()
|
||||
|
||||
for {
|
||||
frame, err := reader.ReadSync(ctx)
|
||||
if err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
return fmt.Errorf("read sync group: %w", err)
|
||||
}
|
||||
|
||||
if err := videoSink.ConsumeVideo(ctx, frame.Video); err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
return &videoSinkError{err: err}
|
||||
}
|
||||
|
||||
if err := audioSink.ConsumeAudio(ctx, frame.Audio); err != nil {
|
||||
if ctx.Err() != nil {
|
||||
return ctx.Err()
|
||||
}
|
||||
return &audioSinkError{err: err}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,241 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
)
|
||||
|
||||
type fakeSyncFactory struct {
|
||||
reader SyncReader
|
||||
err error
|
||||
calls int
|
||||
videoConfig FeedConfig
|
||||
audioConfig FeedConfig
|
||||
}
|
||||
|
||||
func (f *fakeSyncFactory) OpenSync(
|
||||
_ context.Context,
|
||||
videoConfig FeedConfig,
|
||||
audioConfig FeedConfig,
|
||||
) (SyncReader, error) {
|
||||
f.calls++
|
||||
f.videoConfig = videoConfig
|
||||
f.audioConfig = audioConfig
|
||||
return f.reader, f.err
|
||||
}
|
||||
|
||||
type fakeSyncReader struct {
|
||||
frames []SyncFrame
|
||||
readErr error
|
||||
closeErr error
|
||||
readCalls int
|
||||
closed bool
|
||||
read func(context.Context) (SyncFrame, error)
|
||||
}
|
||||
|
||||
func (r *fakeSyncReader) ReadSync(ctx context.Context) (SyncFrame, error) {
|
||||
r.readCalls++
|
||||
if r.read != nil {
|
||||
return r.read(ctx)
|
||||
}
|
||||
if len(r.frames) == 0 {
|
||||
return SyncFrame{}, r.readErr
|
||||
}
|
||||
frame := r.frames[0]
|
||||
r.frames = r.frames[1:]
|
||||
return frame, nil
|
||||
}
|
||||
|
||||
func (r *fakeSyncReader) Close() error {
|
||||
r.closed = true
|
||||
return r.closeErr
|
||||
}
|
||||
|
||||
type orderedVideoSink struct {
|
||||
order *[]string
|
||||
err error
|
||||
frame VideoFrame
|
||||
}
|
||||
|
||||
func (s *orderedVideoSink) ConsumeVideo(_ context.Context, frame VideoFrame) error {
|
||||
*s.order = append(*s.order, "video")
|
||||
s.frame = frame
|
||||
return s.err
|
||||
}
|
||||
|
||||
type orderedAudioSink struct {
|
||||
order *[]string
|
||||
err error
|
||||
frame AudioFrame
|
||||
}
|
||||
|
||||
func (s *orderedAudioSink) ConsumeAudio(_ context.Context, frame AudioFrame) error {
|
||||
*s.order = append(*s.order, "audio")
|
||||
s.frame = frame
|
||||
return s.err
|
||||
}
|
||||
|
||||
func TestRunSyncAttemptOpenFailure(t *testing.T) {
|
||||
openErr := errors.New("open failed")
|
||||
factory := &fakeSyncFactory{err: openErr}
|
||||
videoConfig := FeedConfig{Domain: "/mxl", UUID: "video", Active: true}
|
||||
audioConfig := FeedConfig{Domain: "/mxl", UUID: "audio", Active: true}
|
||||
|
||||
err := runSyncAttempt(
|
||||
context.Background(),
|
||||
factory,
|
||||
&fakeVideoSink{},
|
||||
&fakeAudioSink{},
|
||||
videoConfig,
|
||||
audioConfig,
|
||||
)
|
||||
|
||||
if !errors.Is(err, openErr) {
|
||||
t.Fatalf("runSyncAttempt() error = %v, want %v", err, openErr)
|
||||
}
|
||||
if factory.calls != 1 || factory.videoConfig != videoConfig || factory.audioConfig != audioConfig {
|
||||
t.Fatalf(
|
||||
"factory call = %d, video %#v, audio %#v",
|
||||
factory.calls,
|
||||
factory.videoConfig,
|
||||
factory.audioConfig,
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunSyncAttemptConsumesBorrowedPairInOrderThenReturnsReadError(t *testing.T) {
|
||||
readErr := errors.New("sync read failed")
|
||||
videoPayload := []byte{1, 2, 3, 4}
|
||||
audioSamples := [][]byte{{5, 6, 7, 8}}
|
||||
want := SyncFrame{
|
||||
Video: VideoFrame{Index: 10, Payload: videoPayload},
|
||||
Audio: AudioFrame{Index: 20, Samples: audioSamples},
|
||||
}
|
||||
reader := &fakeSyncReader{frames: []SyncFrame{want}, readErr: readErr}
|
||||
var order []string
|
||||
videoSink := &orderedVideoSink{order: &order}
|
||||
audioSink := &orderedAudioSink{order: &order}
|
||||
|
||||
err := runSyncAttempt(
|
||||
context.Background(),
|
||||
&fakeSyncFactory{reader: reader},
|
||||
videoSink,
|
||||
audioSink,
|
||||
FeedConfig{},
|
||||
FeedConfig{},
|
||||
)
|
||||
|
||||
if !errors.Is(err, readErr) {
|
||||
t.Fatalf("runSyncAttempt() error = %v, want %v", err, readErr)
|
||||
}
|
||||
if len(order) != 2 || order[0] != "video" || order[1] != "audio" {
|
||||
t.Fatalf("sink order = %v, want [video audio]", order)
|
||||
}
|
||||
if &videoSink.frame.Payload[0] != &videoPayload[0] {
|
||||
t.Fatal("video payload was copied")
|
||||
}
|
||||
if &audioSink.frame.Samples[0][0] != &audioSamples[0][0] {
|
||||
t.Fatal("audio samples were copied")
|
||||
}
|
||||
if reader.readCalls != 2 || !reader.closed {
|
||||
t.Fatalf("reader calls = %d, closed = %t; want 2, true", reader.readCalls, reader.closed)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunSyncAttemptVideoSinkFailureSkipsAudio(t *testing.T) {
|
||||
sinkErr := errors.New("video output failed")
|
||||
reader := &fakeSyncReader{frames: []SyncFrame{{}}}
|
||||
var order []string
|
||||
|
||||
err := runSyncAttempt(
|
||||
context.Background(),
|
||||
&fakeSyncFactory{reader: reader},
|
||||
&orderedVideoSink{order: &order, err: sinkErr},
|
||||
&orderedAudioSink{order: &order},
|
||||
FeedConfig{},
|
||||
FeedConfig{},
|
||||
)
|
||||
|
||||
if !errors.Is(err, sinkErr) {
|
||||
t.Fatalf("runSyncAttempt() error = %v, want %v", err, sinkErr)
|
||||
}
|
||||
var typedErr *videoSinkError
|
||||
if !errors.As(err, &typedErr) {
|
||||
t.Fatalf("runSyncAttempt() error type = %T, want *videoSinkError", err)
|
||||
}
|
||||
if len(order) != 1 || order[0] != "video" {
|
||||
t.Fatalf("sink order = %v, want [video]", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunSyncAttemptAudioSinkFailureFollowsVideo(t *testing.T) {
|
||||
sinkErr := errors.New("audio output failed")
|
||||
reader := &fakeSyncReader{frames: []SyncFrame{{}}}
|
||||
var order []string
|
||||
|
||||
err := runSyncAttempt(
|
||||
context.Background(),
|
||||
&fakeSyncFactory{reader: reader},
|
||||
&orderedVideoSink{order: &order},
|
||||
&orderedAudioSink{order: &order, err: sinkErr},
|
||||
FeedConfig{},
|
||||
FeedConfig{},
|
||||
)
|
||||
|
||||
if !errors.Is(err, sinkErr) {
|
||||
t.Fatalf("runSyncAttempt() error = %v, want %v", err, sinkErr)
|
||||
}
|
||||
var typedErr *audioSinkError
|
||||
if !errors.As(err, &typedErr) {
|
||||
t.Fatalf("runSyncAttempt() error type = %T, want *audioSinkError", err)
|
||||
}
|
||||
if len(order) != 2 || order[0] != "video" || order[1] != "audio" {
|
||||
t.Fatalf("sink order = %v, want [video audio]", order)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunSyncAttemptCancellationClosesReader(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
reader := &fakeSyncReader{
|
||||
read: func(ctx context.Context) (SyncFrame, error) {
|
||||
cancel()
|
||||
return SyncFrame{}, ctx.Err()
|
||||
},
|
||||
}
|
||||
|
||||
err := runSyncAttempt(
|
||||
ctx,
|
||||
&fakeSyncFactory{reader: reader},
|
||||
&fakeVideoSink{},
|
||||
&fakeAudioSink{},
|
||||
FeedConfig{},
|
||||
FeedConfig{},
|
||||
)
|
||||
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("runSyncAttempt() error = %v, want context.Canceled", err)
|
||||
}
|
||||
if !reader.closed {
|
||||
t.Fatal("reader was not closed")
|
||||
}
|
||||
}
|
||||
|
||||
func TestRunSyncAttemptJoinsReadAndCloseErrors(t *testing.T) {
|
||||
readErr := errors.New("read failed")
|
||||
closeErr := errors.New("close failed")
|
||||
reader := &fakeSyncReader{readErr: readErr, closeErr: closeErr}
|
||||
|
||||
err := runSyncAttempt(
|
||||
context.Background(),
|
||||
&fakeSyncFactory{reader: reader},
|
||||
&fakeVideoSink{},
|
||||
&fakeAudioSink{},
|
||||
FeedConfig{},
|
||||
FeedConfig{},
|
||||
)
|
||||
|
||||
if !errors.Is(err, readErr) || !errors.Is(err, closeErr) {
|
||||
t.Fatalf("runSyncAttempt() error = %v, want read and close errors", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,179 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"fmt"
|
||||
)
|
||||
|
||||
var (
|
||||
ErrSyncFactoryRequired = errors.New("sync reader factory is required")
|
||||
ErrSyncVideoSinkRequired = errors.New("sync video sink is required")
|
||||
ErrSyncAudioSinkRequired = errors.New("sync audio sink is required")
|
||||
ErrSyncRetryDeciderRequired = errors.New("sync retry decider is required")
|
||||
ErrSyncFeedsInactive = errors.New("both sync feeds must be active")
|
||||
)
|
||||
|
||||
type SyncWorker struct {
|
||||
factory SyncReaderFactory
|
||||
videoSink VideoSink
|
||||
audioSink AudioSink
|
||||
retry RetryPolicy
|
||||
shouldRetry retryDecider
|
||||
observer StatusObserver
|
||||
wait waitFunc
|
||||
}
|
||||
|
||||
func NewSyncWorker(
|
||||
factory SyncReaderFactory,
|
||||
videoSink VideoSink,
|
||||
audioSink AudioSink,
|
||||
retry RetryPolicy,
|
||||
shouldRetry func(error) bool,
|
||||
observer StatusObserver,
|
||||
) (*SyncWorker, error) {
|
||||
if factory == nil {
|
||||
return nil, ErrSyncFactoryRequired
|
||||
}
|
||||
if videoSink == nil {
|
||||
return nil, ErrSyncVideoSinkRequired
|
||||
}
|
||||
if audioSink == nil {
|
||||
return nil, ErrSyncAudioSinkRequired
|
||||
}
|
||||
if err := retry.Validate(); err != nil {
|
||||
return nil, fmt.Errorf("validate sync retry policy: %w", err)
|
||||
}
|
||||
if shouldRetry == nil {
|
||||
return nil, ErrSyncRetryDeciderRequired
|
||||
}
|
||||
|
||||
return &SyncWorker{
|
||||
factory: factory,
|
||||
videoSink: videoSink,
|
||||
audioSink: audioSink,
|
||||
retry: retry,
|
||||
shouldRetry: shouldRetry,
|
||||
observer: observer,
|
||||
wait: waitForRetry,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (w *SyncWorker) emit(status Status) {
|
||||
if w.observer != nil {
|
||||
w.observer(status)
|
||||
}
|
||||
}
|
||||
|
||||
func (w *SyncWorker) Run(
|
||||
ctx context.Context,
|
||||
videoConfig FeedConfig,
|
||||
audioConfig FeedConfig,
|
||||
) error {
|
||||
if err := videoConfig.Validate(); err != nil {
|
||||
return fmt.Errorf("validate sync video config: %w", err)
|
||||
}
|
||||
if err := audioConfig.Validate(); err != nil {
|
||||
return fmt.Errorf("validate sync audio config: %w", err)
|
||||
}
|
||||
if !videoConfig.Active || !audioConfig.Active {
|
||||
return ErrSyncFeedsInactive
|
||||
}
|
||||
|
||||
attemptNumber := 0
|
||||
var latestRetry retryEvent
|
||||
|
||||
attempt := func(ctx context.Context) (bool, error) {
|
||||
attemptNumber++
|
||||
|
||||
state := StateConnecting
|
||||
if attemptNumber > 1 {
|
||||
state = StateReconnecting
|
||||
}
|
||||
w.emit(Status{
|
||||
Unit: UnitSync,
|
||||
State: state,
|
||||
Attempt: attemptNumber,
|
||||
})
|
||||
|
||||
attemptAudioSink := &stabilityAudioSink{
|
||||
sink: w.audioSink,
|
||||
onStable: func() {
|
||||
w.emit(Status{
|
||||
Unit: UnitSync,
|
||||
State: StatePlaying,
|
||||
Attempt: attemptNumber,
|
||||
})
|
||||
},
|
||||
}
|
||||
err := runSyncAttempt(
|
||||
ctx,
|
||||
w.factory,
|
||||
w.videoSink,
|
||||
attemptAudioSink,
|
||||
videoConfig,
|
||||
audioConfig,
|
||||
)
|
||||
return attemptAudioSink.stable, err
|
||||
}
|
||||
|
||||
decide := func(err error) bool {
|
||||
var videoErr *videoSinkError
|
||||
if errors.As(err, &videoErr) {
|
||||
return false
|
||||
}
|
||||
var audioErr *audioSinkError
|
||||
if errors.As(err, &audioErr) {
|
||||
return false
|
||||
}
|
||||
return w.shouldRetry(err)
|
||||
}
|
||||
observeRetry := func(event retryEvent) {
|
||||
latestRetry = event
|
||||
if !event.WillRetry {
|
||||
return
|
||||
}
|
||||
w.emit(Status{
|
||||
Unit: UnitSync,
|
||||
State: StateReconnecting,
|
||||
Attempt: attemptNumber + 1,
|
||||
FailedAttempts: event.FailedAttempts,
|
||||
RetryIn: event.RetryIn,
|
||||
Err: event.Err,
|
||||
})
|
||||
}
|
||||
err := runWithRetry(
|
||||
ctx,
|
||||
w.retry,
|
||||
attempt,
|
||||
decide,
|
||||
w.wait,
|
||||
observeRetry,
|
||||
)
|
||||
if ctx.Err() != nil {
|
||||
w.emit(Status{
|
||||
Unit: UnitSync,
|
||||
State: StateStopping,
|
||||
})
|
||||
w.emit(Status{
|
||||
Unit: UnitSync,
|
||||
State: StateIdle,
|
||||
})
|
||||
return ctx.Err()
|
||||
}
|
||||
if err != nil {
|
||||
w.emit(Status{
|
||||
Unit: UnitSync,
|
||||
State: StateFailed,
|
||||
Attempt: attemptNumber,
|
||||
FailedAttempts: latestRetry.FailedAttempts,
|
||||
Err: err,
|
||||
})
|
||||
return err
|
||||
}
|
||||
w.emit(Status{
|
||||
Unit: UnitSync,
|
||||
State: StateIdle,
|
||||
})
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,261 @@
|
||||
package playback
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
type syncOpenResult struct {
|
||||
reader SyncReader
|
||||
err error
|
||||
}
|
||||
|
||||
type scriptedSyncFactory struct {
|
||||
results []syncOpenResult
|
||||
calls int
|
||||
}
|
||||
|
||||
func (f *scriptedSyncFactory) OpenSync(
|
||||
context.Context,
|
||||
FeedConfig,
|
||||
FeedConfig,
|
||||
) (SyncReader, error) {
|
||||
if f.calls >= len(f.results) {
|
||||
return nil, errors.New("unexpected sync open attempt")
|
||||
}
|
||||
result := f.results[f.calls]
|
||||
f.calls++
|
||||
return result.reader, result.err
|
||||
}
|
||||
|
||||
func activeSyncConfigs() (FeedConfig, FeedConfig) {
|
||||
return FeedConfig{Domain: "/mxl", UUID: "video", Active: true},
|
||||
FeedConfig{Domain: "/mxl", UUID: "audio", Active: true}
|
||||
}
|
||||
|
||||
func newTestSyncWorker(
|
||||
t *testing.T,
|
||||
factory SyncReaderFactory,
|
||||
videoSink VideoSink,
|
||||
audioSink AudioSink,
|
||||
maxAttempts int,
|
||||
observer StatusObserver,
|
||||
) *SyncWorker {
|
||||
t.Helper()
|
||||
worker, err := NewSyncWorker(
|
||||
factory,
|
||||
videoSink,
|
||||
audioSink,
|
||||
testRetryPolicy(maxAttempts),
|
||||
func(error) bool { return true },
|
||||
observer,
|
||||
)
|
||||
if err != nil {
|
||||
t.Fatalf("NewSyncWorker() error = %v", err)
|
||||
}
|
||||
worker.wait = func(context.Context, time.Duration) error { return nil }
|
||||
return worker
|
||||
}
|
||||
|
||||
func TestNewSyncWorkerValidatesDependencies(t *testing.T) {
|
||||
factory := &scriptedSyncFactory{}
|
||||
videoSink := &fakeVideoSink{}
|
||||
audioSink := &fakeAudioSink{}
|
||||
retry := testRetryPolicy(3)
|
||||
decide := func(error) bool { return true }
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
factory SyncReaderFactory
|
||||
videoSink VideoSink
|
||||
audioSink AudioSink
|
||||
retry RetryPolicy
|
||||
shouldRetry func(error) bool
|
||||
wantErr error
|
||||
}{
|
||||
{"missing factory", nil, videoSink, audioSink, retry, decide, ErrSyncFactoryRequired},
|
||||
{"missing video sink", factory, nil, audioSink, retry, decide, ErrSyncVideoSinkRequired},
|
||||
{"missing audio sink", factory, videoSink, nil, retry, decide, ErrSyncAudioSinkRequired},
|
||||
{"invalid retry", factory, videoSink, audioSink, RetryPolicy{}, decide, ErrInvalidRetryDelay},
|
||||
{"missing decider", factory, videoSink, audioSink, retry, nil, ErrSyncRetryDeciderRequired},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
worker, err := NewSyncWorker(
|
||||
tt.factory, tt.videoSink, tt.audioSink, tt.retry, tt.shouldRetry, nil,
|
||||
)
|
||||
if worker != nil {
|
||||
t.Fatal("NewSyncWorker() worker is not nil")
|
||||
}
|
||||
if !errors.Is(err, tt.wantErr) {
|
||||
t.Fatalf("NewSyncWorker() error = %v, want %v", err, tt.wantErr)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncWorkerRejectsInvalidOrInactiveFeeds(t *testing.T) {
|
||||
video, audio := activeSyncConfigs()
|
||||
tests := []struct {
|
||||
name string
|
||||
video FeedConfig
|
||||
audio FeedConfig
|
||||
want error
|
||||
}{
|
||||
{"invalid video", FeedConfig{Active: true}, audio, ErrActiveFeedNotConfigured},
|
||||
{"invalid audio", video, FeedConfig{Active: true}, ErrActiveFeedNotConfigured},
|
||||
{"inactive video", FeedConfig{Domain: video.Domain, UUID: video.UUID}, audio, ErrSyncFeedsInactive},
|
||||
{"inactive audio", video, FeedConfig{Domain: audio.Domain, UUID: audio.UUID}, ErrSyncFeedsInactive},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
factory := &scriptedSyncFactory{}
|
||||
worker := newTestSyncWorker(t, factory, &fakeVideoSink{}, &fakeAudioSink{}, 3, nil)
|
||||
err := worker.Run(context.Background(), tt.video, tt.audio)
|
||||
if !errors.Is(err, tt.want) {
|
||||
t.Fatalf("Run() error = %v, want %v", err, tt.want)
|
||||
}
|
||||
if factory.calls != 0 {
|
||||
t.Fatalf("factory calls = %d, want 0", factory.calls)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncWorkerExhaustsOpenRetries(t *testing.T) {
|
||||
openErr := errors.New("sync producer unavailable")
|
||||
factory := &scriptedSyncFactory{results: []syncOpenResult{{err: openErr}, {err: openErr}}}
|
||||
var statuses []Status
|
||||
worker := newTestSyncWorker(t, factory, &fakeVideoSink{}, &fakeAudioSink{}, 2,
|
||||
func(status Status) { statuses = append(statuses, status) })
|
||||
video, audio := activeSyncConfigs()
|
||||
|
||||
err := worker.Run(context.Background(), video, audio)
|
||||
if !errors.Is(err, openErr) {
|
||||
t.Fatalf("Run() error = %v, want %v", err, openErr)
|
||||
}
|
||||
if factory.calls != 2 {
|
||||
t.Fatalf("factory calls = %d, want 2", factory.calls)
|
||||
}
|
||||
want := []State{StateConnecting, StateReconnecting, StateReconnecting, StateFailed}
|
||||
if len(statuses) != len(want) {
|
||||
t.Fatalf("statuses = %+v, want %d entries", statuses, len(want))
|
||||
}
|
||||
for i, state := range want {
|
||||
if statuses[i].Unit != UnitSync || statuses[i].State != state {
|
||||
t.Errorf("status %d = %+v, want unit=%v state=%v", i, statuses[i], UnitSync, state)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncWorkerStablePairResetsRetryCounter(t *testing.T) {
|
||||
readErr := errors.New("sync disconnected")
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
lastReader := &fakeSyncReader{read: func(ctx context.Context) (SyncFrame, error) {
|
||||
cancel()
|
||||
return SyncFrame{}, ctx.Err()
|
||||
}}
|
||||
factory := &scriptedSyncFactory{results: []syncOpenResult{
|
||||
{reader: &fakeSyncReader{frames: []SyncFrame{{}}, readErr: readErr}},
|
||||
{reader: &fakeSyncReader{frames: []SyncFrame{{}}, readErr: readErr}},
|
||||
{reader: lastReader},
|
||||
}}
|
||||
var statuses []Status
|
||||
worker := newTestSyncWorker(t, factory, &fakeVideoSink{}, &fakeAudioSink{}, 2,
|
||||
func(status Status) { statuses = append(statuses, status) })
|
||||
video, audio := activeSyncConfigs()
|
||||
|
||||
err := worker.Run(ctx, video, audio)
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("Run() error = %v, want context.Canceled", err)
|
||||
}
|
||||
if factory.calls != 3 {
|
||||
t.Fatalf("factory calls = %d, want 3", factory.calls)
|
||||
}
|
||||
var retryFailures []int
|
||||
for _, status := range statuses {
|
||||
if status.State == StateReconnecting && status.RetryIn > 0 {
|
||||
retryFailures = append(retryFailures, status.FailedAttempts)
|
||||
}
|
||||
}
|
||||
if len(retryFailures) != 2 || retryFailures[0] != 1 || retryFailures[1] != 1 {
|
||||
t.Fatalf("retry failure counts = %v, want [1 1]", retryFailures)
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncWorkerDoesNotRetrySinkFailures(t *testing.T) {
|
||||
sinkErr := errors.New("output failed")
|
||||
tests := []struct {
|
||||
name string
|
||||
videoSink VideoSink
|
||||
audioSink AudioSink
|
||||
}{
|
||||
{"video", &fakeVideoSink{err: sinkErr}, &fakeAudioSink{}},
|
||||
{"audio", &fakeVideoSink{}, &fakeAudioSink{err: sinkErr}},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
reader := &fakeSyncReader{frames: []SyncFrame{{}}}
|
||||
factory := &scriptedSyncFactory{results: []syncOpenResult{{reader: reader}}}
|
||||
deciderCalls := 0
|
||||
worker, err := NewSyncWorker(factory, tt.videoSink, tt.audioSink, testRetryPolicy(0),
|
||||
func(error) bool { deciderCalls++; return true }, nil)
|
||||
if err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
worker.wait = func(context.Context, time.Duration) error { return nil }
|
||||
video, audio := activeSyncConfigs()
|
||||
err = worker.Run(context.Background(), video, audio)
|
||||
if !errors.Is(err, sinkErr) {
|
||||
t.Fatalf("Run() error = %v, want %v", err, sinkErr)
|
||||
}
|
||||
if factory.calls != 1 || deciderCalls != 0 || !reader.closed {
|
||||
t.Fatalf("calls=%d deciderCalls=%d closed=%t", factory.calls, deciderCalls, reader.closed)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestSyncWorkerEmitsPlayingThenStopsOnCancellation(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
reader := &fakeSyncReader{
|
||||
frames: []SyncFrame{{}},
|
||||
read: func(ctx context.Context) (SyncFrame, error) {
|
||||
cancel()
|
||||
return SyncFrame{}, ctx.Err()
|
||||
},
|
||||
}
|
||||
// Preserve the first frame before switching to the cancellation callback.
|
||||
readCalls := 0
|
||||
reader.read = func(ctx context.Context) (SyncFrame, error) {
|
||||
readCalls++
|
||||
if readCalls == 1 {
|
||||
return SyncFrame{}, nil
|
||||
}
|
||||
cancel()
|
||||
return SyncFrame{}, ctx.Err()
|
||||
}
|
||||
factory := &scriptedSyncFactory{results: []syncOpenResult{{reader: reader}}}
|
||||
var statuses []Status
|
||||
worker := newTestSyncWorker(t, factory, &fakeVideoSink{}, &fakeAudioSink{}, 3,
|
||||
func(status Status) { statuses = append(statuses, status) })
|
||||
video, audio := activeSyncConfigs()
|
||||
|
||||
err := worker.Run(ctx, video, audio)
|
||||
if !errors.Is(err, context.Canceled) {
|
||||
t.Fatalf("Run() error = %v, want context.Canceled", err)
|
||||
}
|
||||
want := []State{StateConnecting, StatePlaying, StateStopping, StateIdle}
|
||||
if len(statuses) != len(want) {
|
||||
t.Fatalf("statuses = %+v, want %v", statuses, want)
|
||||
}
|
||||
for i, state := range want {
|
||||
if statuses[i].Unit != UnitSync || statuses[i].State != state {
|
||||
t.Errorf("status %d = %+v, want unit=%v state=%v", i, statuses[i], UnitSync, state)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user