152 lines
5.0 KiB
Go
152 lines
5.0 KiB
Go
package media_test
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"net"
|
|
"os"
|
|
"path/filepath"
|
|
"testing"
|
|
"time"
|
|
|
|
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media"
|
|
"git.ilapage.cn/ila/yovision/Sense/server/app/sense/media_shard"
|
|
"gorm.io/driver/sqlite"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
func freeAddress(t *testing.T) string {
|
|
t.Helper()
|
|
listener, err := net.Listen("tcp", "127.0.0.1:0")
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
address := listener.Addr().String()
|
|
if err = listener.Close(); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
return address
|
|
}
|
|
|
|
func startTestMediaMTX(t *testing.T, binary string) string {
|
|
t.Helper()
|
|
apiAddress, rtspAddress := freeAddress(t), freeAddress(t)
|
|
configPath := filepath.Join(t.TempDir(), "mediamtx.yml")
|
|
config := fmt.Sprintf("logLevel: warn\napi: true\napiAddress: %s\nrtspAddress: %s\nrtspTransports: [tcp]\nrtmp: false\nhls: false\nwebrtc: false\nsrt: false\nmoq: false\nplayback: false\npaths: {}\n", apiAddress, rtspAddress)
|
|
if err := os.WriteFile(configPath, []byte(config), 0o600); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
supervisor := media.NewSupervisor(binary, configPath)
|
|
if err := supervisor.Start(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(func() {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = supervisor.Stop(ctx)
|
|
})
|
|
controller, err := media.NewHTTPController("http://" + apiAddress)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
deadline := time.Now().Add(8 * time.Second)
|
|
for {
|
|
if err = controller.Health(context.Background()); err == nil {
|
|
break
|
|
}
|
|
if time.Now().After(deadline) {
|
|
t.Fatalf("MediaMTX did not become ready: %v", err)
|
|
}
|
|
time.Sleep(100 * time.Millisecond)
|
|
}
|
|
return "http://" + apiAddress
|
|
}
|
|
|
|
func TestTwoRealMediaMTXShardsAreIndependentlyUsable(t *testing.T) {
|
|
binary := os.Getenv("SENSE_MEDIAMTX_TEST_BINARY")
|
|
if binary == "" {
|
|
t.Skip("set SENSE_MEDIAMTX_TEST_BINARY to run the real MediaMTX integration")
|
|
}
|
|
first, second := startTestMediaMTX(t, binary), startTestMediaMTX(t, binary)
|
|
db, err := gorm.Open(sqlite.Open("file:real_media_shards?mode=memory&cache=shared"), &gorm.Config{})
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err = db.AutoMigrate(&media_shard.Shard{}, &media_shard.Assignment{}); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
service := media_shard.NewService(db, func(ctx context.Context, endpoint string) error {
|
|
controller, buildErr := media.NewHTTPController(endpoint)
|
|
if buildErr != nil {
|
|
return buildErr
|
|
}
|
|
return controller.Health(ctx)
|
|
})
|
|
specs := []media_shard.Spec{{ID: "one", Name: "测试分片一", Mode: "external", ControlAPI: first, Capacity: 2}, {ID: "two", Name: "测试分片二", Mode: "external", ControlAPI: second, Capacity: 2}}
|
|
if err = service.SyncSpecs(context.Background(), specs); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err = service.RefreshAll(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
for _, routeID := range []string{"route-a", "route-b", "route-c", "route-d"} {
|
|
if _, err = service.EnsureAssignment(context.Background(), routeID); err != nil {
|
|
t.Fatalf("assign %s: %v", routeID, err)
|
|
}
|
|
}
|
|
result, err := service.List(context.Background())
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if result.Summary.Available != 2 || result.Summary.Assigned != 4 || result.Summary.Configured != 4 {
|
|
t.Fatalf("unexpected real shard summary: %+v", result.Summary)
|
|
}
|
|
}
|
|
|
|
func TestRealMediaMTXControlLifecycle(t *testing.T) {
|
|
binary := os.Getenv("SENSE_MEDIAMTX_TEST_BINARY")
|
|
if binary == "" {
|
|
t.Skip("set SENSE_MEDIAMTX_TEST_BINARY to run the real MediaMTX integration")
|
|
}
|
|
apiAddress, rtspAddress := freeAddress(t), freeAddress(t)
|
|
configPath := filepath.Join(t.TempDir(), "mediamtx.yml")
|
|
config := fmt.Sprintf("logLevel: warn\napi: true\napiAddress: %s\nrtspAddress: %s\nrtspTransports: [tcp]\nrtmp: false\nhls: false\nwebrtc: false\nsrt: false\nmoq: false\nplayback: false\npaths: {}\n", apiAddress, rtspAddress)
|
|
if err := os.WriteFile(configPath, []byte(config), 0o600); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
controller, err := media.NewHTTPController("http://" + apiAddress)
|
|
if err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
supervisor := media.NewSupervisor(binary, configPath)
|
|
if err = supervisor.Start(context.Background()); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
t.Cleanup(func() {
|
|
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
|
defer cancel()
|
|
_ = supervisor.Stop(ctx)
|
|
})
|
|
deadline := time.Now().Add(8 * time.Second)
|
|
for {
|
|
err = controller.Health(context.Background())
|
|
if err == nil {
|
|
break
|
|
}
|
|
if time.Now().After(deadline) {
|
|
t.Fatalf("MediaMTX did not become ready: %v", err)
|
|
}
|
|
time.Sleep(100 * time.Millisecond)
|
|
}
|
|
source := media.Source{Path: "sense_integration", URI: "rtsp://127.0.0.1:65530/test", Username: "synthetic-user", Password: "synthetic-password"}
|
|
if err = controller.Apply(context.Background(), source); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
if err = controller.Apply(context.Background(), source); err != nil {
|
|
t.Fatalf("replace must be idempotent: %v", err)
|
|
}
|
|
if err = controller.Delete(context.Background(), source.Path); err != nil {
|
|
t.Fatal(err)
|
|
}
|
|
}
|