169 lines
4.6 KiB
Go
169 lines
4.6 KiB
Go
package media
|
|
|
|
import (
|
|
"bytes"
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"net/url"
|
|
"regexp"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
var validPath = regexp.MustCompile(`^[A-Za-z0-9_-]{1,96}$`)
|
|
|
|
type Source struct {
|
|
Path, URI, Username, Password string
|
|
}
|
|
|
|
type PathStatus struct {
|
|
Exists, Ready bool
|
|
Readers int
|
|
}
|
|
|
|
type Controller interface {
|
|
Health(context.Context) error
|
|
Apply(context.Context, Source) error
|
|
Delete(context.Context, string) error
|
|
Status(context.Context, string) (PathStatus, error)
|
|
}
|
|
|
|
type HTTPController struct {
|
|
base string
|
|
client *http.Client
|
|
}
|
|
|
|
func NewHTTPController(base string) (*HTTPController, error) {
|
|
if err := validateControlAPI(base); err != nil {
|
|
return nil, err
|
|
}
|
|
transport := http.DefaultTransport.(*http.Transport).Clone()
|
|
transport.Proxy = nil
|
|
return &HTTPController{base: strings.TrimRight(base, "/"), client: &http.Client{
|
|
Timeout: 5 * time.Second, Transport: transport,
|
|
CheckRedirect: func(*http.Request, []*http.Request) error {
|
|
return errors.New("MediaMTX Control API redirect rejected")
|
|
},
|
|
}}, nil
|
|
}
|
|
|
|
func (c *HTTPController) Health(ctx context.Context) error {
|
|
return c.request(ctx, http.MethodGet, "/v3/config/global/get", nil, nil)
|
|
}
|
|
|
|
func (c *HTTPController) Apply(ctx context.Context, source Source) error {
|
|
if !validPath.MatchString(source.Path) {
|
|
return errors.New("invalid MediaMTX path")
|
|
}
|
|
u, err := url.Parse(source.URI)
|
|
if err != nil || u.User != nil || u.Host == "" || (u.Scheme != "rtsp" && u.Scheme != "rtsps") {
|
|
return errors.New("invalid credential-free RTSP source")
|
|
}
|
|
payload := map[string]any{"source": u.String(), "sourceOnDemand": true, "rtspTransport": "tcp"}
|
|
if source.Username != "" {
|
|
// MediaMTX v1.19.3 has no sourceUser/sourcePass fields. Credentials
|
|
// are assembled only for this loopback request and are never stored,
|
|
// logged, or returned by a Sense endpoint.
|
|
u.User = url.UserPassword(source.Username, source.Password)
|
|
payload["source"] = u.String()
|
|
}
|
|
data, err := json.Marshal(payload)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
configured, err := c.configured(ctx, source.Path)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
method, action := http.MethodPost, "add"
|
|
if configured {
|
|
method, action = http.MethodPatch, "patch"
|
|
}
|
|
return c.request(ctx, method, "/v3/config/paths/"+action+"/"+url.PathEscape(source.Path), bytes.NewReader(data), nil)
|
|
}
|
|
|
|
func (c *HTTPController) configured(ctx context.Context, path string) (bool, error) {
|
|
status := 0
|
|
err := c.request(ctx, http.MethodGet, "/v3/config/paths/get/"+url.PathEscape(path), nil, &status)
|
|
if status == http.StatusNotFound {
|
|
return false, nil
|
|
}
|
|
return err == nil, err
|
|
}
|
|
|
|
func (c *HTTPController) Delete(ctx context.Context, path string) error {
|
|
if !validPath.MatchString(path) {
|
|
return errors.New("invalid MediaMTX path")
|
|
}
|
|
status := 0
|
|
err := c.request(ctx, http.MethodDelete, "/v3/config/paths/delete/"+url.PathEscape(path), nil, &status)
|
|
if status == http.StatusNotFound {
|
|
return nil
|
|
}
|
|
return err
|
|
}
|
|
|
|
func (c *HTTPController) Status(ctx context.Context, path string) (PathStatus, error) {
|
|
statusCode := 0
|
|
var raw struct {
|
|
Ready bool `json:"ready"`
|
|
Readers []any `json:"readers"`
|
|
}
|
|
err := c.requestJSON(ctx, http.MethodGet, "/v3/paths/get/"+url.PathEscape(path), &statusCode, &raw)
|
|
if statusCode == http.StatusNotFound {
|
|
return PathStatus{}, nil
|
|
}
|
|
if err != nil {
|
|
return PathStatus{}, err
|
|
}
|
|
return PathStatus{Exists: true, Ready: raw.Ready, Readers: len(raw.Readers)}, nil
|
|
}
|
|
|
|
func (c *HTTPController) request(ctx context.Context, method, path string, body io.Reader, statusOut *int) error {
|
|
return c.requestJSON(ctx, method, path, body, statusOut, nil)
|
|
}
|
|
|
|
func (c *HTTPController) requestJSON(ctx context.Context, method, path string, args ...any) error {
|
|
var body io.Reader
|
|
var statusOut *int
|
|
var target any
|
|
for _, arg := range args {
|
|
switch value := arg.(type) {
|
|
case io.Reader:
|
|
body = value
|
|
case *int:
|
|
statusOut = value
|
|
default:
|
|
target = value
|
|
}
|
|
}
|
|
req, err := http.NewRequestWithContext(ctx, method, c.base+path, body)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if body != nil {
|
|
req.Header.Set("Content-Type", "application/json")
|
|
}
|
|
res, err := c.client.Do(req)
|
|
if err != nil {
|
|
return fmt.Errorf("MediaMTX Control API unavailable: %w", err)
|
|
}
|
|
defer res.Body.Close()
|
|
if statusOut != nil {
|
|
*statusOut = res.StatusCode
|
|
}
|
|
if res.StatusCode < 200 || res.StatusCode >= 300 {
|
|
_, _ = io.Copy(io.Discard, io.LimitReader(res.Body, 1<<20))
|
|
return fmt.Errorf("MediaMTX Control API returned %d", res.StatusCode)
|
|
}
|
|
if target == nil {
|
|
_, _ = io.Copy(io.Discard, io.LimitReader(res.Body, 1<<20))
|
|
return nil
|
|
}
|
|
return json.NewDecoder(io.LimitReader(res.Body, 1<<20)).Decode(target)
|
|
}
|