diff --git a/pkg/cmd/streamplace.go b/pkg/cmd/streamplace.go index 53fe5d2a..8e82e4bb 100644 --- a/pkg/cmd/streamplace.go +++ b/pkg/cmd/streamplace.go @@ -498,7 +498,10 @@ func runMain(ctx context.Context, build *config.BuildFlags, platformJobs []jobFu if cli.WHIPTest != "" { group.Go(func() error { - err := WHIP(strings.Split(cli.WHIPTest, " ")) + // Parse WHIPTest string using the whip command's flag parser + whipCmd := makeWhipCommand(build) + args := strings.Split(cli.WHIPTest, " ") + err := whipCmd.Run(ctx, append([]string{"streamplace", "whip"}, args...)) log.Warn(ctx, "WHIP test complete, sleeping for 3 seconds and shutting down gstreamer") time.Sleep(time.Second * 3) // gst.Deinit() @@ -595,8 +598,28 @@ func makeWhepCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "whep", Usage: "WHEP client", + Flags: []urfavecli.Flag{ + &urfavecli.IntFlag{ + Name: "count", + Usage: "number of concurrent streams (for load testing)", + Value: 1, + }, + &urfavecli.DurationFlag{ + Name: "duration", + Usage: "stop after this long", + }, + &urfavecli.StringFlag{ + Name: "endpoint", + Usage: "endpoint to send the WHEP request to", + }, + }, Action: func(ctx context.Context, cmd *urfavecli.Command) error { - return WHEP(cmd.Args().Slice()) + return WHEP( + ctx, + cmd.Int("count"), + cmd.Duration("duration"), + cmd.String("endpoint"), + ) }, } } @@ -605,8 +628,50 @@ func makeWhipCommand(build *config.BuildFlags) *urfavecli.Command { return &urfavecli.Command{ Name: "whip", Usage: "WHIP client", + Flags: []urfavecli.Flag{ + &urfavecli.StringFlag{ + Name: "stream-key", + Usage: "stream key", + }, + &urfavecli.IntFlag{ + Name: "count", + Usage: "number of concurrent streams (for load testing)", + Value: 1, + }, + &urfavecli.IntFlag{ + Name: "viewers", + Usage: "number of viewers to simulate per stream", + }, + &urfavecli.DurationFlag{ + Name: "duration", + Usage: "duration of the stream", + }, + &urfavecli.StringFlag{ + Name: "file", + Usage: "file to stream (needs to be an MP4 containing H264 video and Opus audio)", + Required: true, + }, + &urfavecli.StringFlag{ + Name: "endpoint", + Usage: "endpoint to send the WHIP request to", + Value: "http://127.0.0.1:38080", + }, + &urfavecli.DurationFlag{ + Name: "freeze-after", + Usage: "freeze the stream after the given duration", + }, + }, Action: func(ctx context.Context, cmd *urfavecli.Command) error { - return WHIP(cmd.Args().Slice()) + return WHIP( + ctx, + cmd.String("stream-key"), + cmd.Int("count"), + cmd.Int("viewers"), + cmd.Duration("duration"), + cmd.String("file"), + cmd.String("endpoint"), + cmd.Duration("freeze-after"), + ) }, } } diff --git a/pkg/cmd/whep.go b/pkg/cmd/whep.go index bfc40a11..a30ca0b2 100644 --- a/pkg/cmd/whep.go +++ b/pkg/cmd/whep.go @@ -2,7 +2,6 @@ package cmd import ( "context" - "flag" "fmt" "io" "net/http" @@ -15,27 +14,16 @@ import ( "stream.place/streamplace/pkg/log" ) -func WHEP(args []string) error { - fs := flag.NewFlagSet("whep", flag.ExitOnError) - count := fs.Int("count", 1, "number of concurrent streams (for load testing)") - duration := fs.Duration("duration", 0, "stop after this long") - endpoint := fs.String("endpoint", "", "endpoint to send the WHEP request to") - err := fs.Parse(args) - - if err != nil { - return err - } - - ctx := context.Background() - if *duration > 0 { +func WHEP(ctx context.Context, count int, duration time.Duration, endpoint string) error { + if duration > 0 { var cancel context.CancelFunc - ctx, cancel = context.WithTimeout(ctx, *duration) + ctx, cancel = context.WithTimeout(ctx, duration) defer cancel() } w := &WHEPClient{ - Endpoint: *endpoint, - Count: *count, + Endpoint: endpoint, + Count: count, } return w.WHEP(ctx) diff --git a/pkg/cmd/whip.go b/pkg/cmd/whip.go index 9a377a48..2f67e8b6 100644 --- a/pkg/cmd/whip.go +++ b/pkg/cmd/whip.go @@ -2,7 +2,6 @@ package cmd import ( "context" - "flag" "fmt" "io" "net/http" @@ -20,38 +19,25 @@ import ( "stream.place/streamplace/pkg/media" ) -func WHIP(args []string) error { - fs := flag.NewFlagSet("whip", flag.ExitOnError) - streamKey := fs.String("stream-key", "", "stream key") - count := fs.Int("count", 1, "number of concurrent streams (for load testing)") - viewers := fs.Int("viewers", 0, "number of viewers to simulate per stream") - duration := fs.Duration("duration", 0, "duration of the stream") - file := fs.String("file", "", "file to stream (needs to be an MP4 containing H264 video and Opus audio)") - endpoint := fs.String("endpoint", "http://127.0.0.1:38080", "endpoint to send the WHIP request to") - freezeAfter := fs.Duration("freeze-after", 0, "freeze the stream after the given duration") - err := fs.Parse(args) - if *file == "" { +func WHIP(ctx context.Context, streamKey string, count int, viewers int, duration time.Duration, file string, endpoint string, freezeAfter time.Duration) error { + if file == "" { return fmt.Errorf("file is required") } - if err != nil { - return err - } gstinit.InitGST() - ctx := context.Background() - if *duration > 0 { + if duration > 0 { var cancel context.CancelFunc - ctx, cancel = context.WithTimeout(ctx, *duration) + ctx, cancel = context.WithTimeout(ctx, duration) defer cancel() } w := &WHIPClient{ - StreamKey: *streamKey, - File: *file, - Endpoint: *endpoint, - Count: *count, - FreezeAfter: *freezeAfter, - Viewers: *viewers, + StreamKey: streamKey, + File: file, + Endpoint: endpoint, + Count: count, + FreezeAfter: freezeAfter, + Viewers: viewers, } return w.WHIP(ctx)