Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,10 @@ jobs:
uses: golangci/golangci-lint-action@v9
with:
version: v2.13.2
only-new-issues: true
# The GitHub PR diff API rejects reviews larger than 20,000 lines.
# Keep the new-code gate independent of that API limit.
args: >-
--new-from-rev=${{ github.event_name == 'pull_request' && format('origin/{0}', github.base_ref) || 'HEAD^' }}

test:
name: Test
Expand Down
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ configs/liveforge.local.yaml
# Recordings
/data/
/recordings/
/module/record/recordings/

# Node (Playwright test tooling)
node_modules/
Expand Down
43 changes: 31 additions & 12 deletions README.md

Large diffs are not rendered by default.

43 changes: 31 additions & 12 deletions README.zh-CN.md

Large diffs are not rendered by default.

78 changes: 58 additions & 20 deletions agent-manifest.json

Large diffs are not rendered by default.

61 changes: 40 additions & 21 deletions config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,10 @@ package config

import "time"

// DefaultRuntimeSourceMaxBytes is the default complete-document limit for
// file and network-backed runtime configuration sources.
const DefaultRuntimeSourceMaxBytes int64 = 4 << 20

// Config is the root configuration for the streaming server.
type Config struct {
Server ServerConfig `yaml:"server"`
Expand Down Expand Up @@ -41,7 +45,8 @@ type RuntimeConfig struct {
}

type RuntimeFileSourceConfig struct {
Path string `yaml:"path"`
Path string `yaml:"path"`
MaxBytes int64 `yaml:"max_bytes"`
}

type RuntimeHTTPSourceConfig struct {
Expand All @@ -66,6 +71,7 @@ type RuntimeRedisSourceConfig struct {
Hash string `yaml:"hash"`
VersionKey string `yaml:"version_key"`
TLS bool `yaml:"tls"`
MaxBytes int64 `yaml:"max_bytes"`
}

// AudioCodecConfig controls audio transcoding between protocols.
Expand Down Expand Up @@ -103,9 +109,10 @@ type LimitsConfig struct {

// RateLimitConfig holds per-IP HTTP rate limiting settings.
type RateLimitConfig struct {
Enabled bool `yaml:"enabled"`
Rate float64 `yaml:"rate"` // requests per second per IP
Burst int `yaml:"burst"` // max burst size per IP
Enabled bool `yaml:"enabled"`
Rate float64 `yaml:"rate"` // requests per second per IP
Burst int `yaml:"burst"` // max burst size per IP
TrustedProxies []string `yaml:"trusted_proxies"`
}

// RTMPConfig holds RTMP module settings.
Expand Down Expand Up @@ -231,18 +238,20 @@ type SIPAuth struct {

// SIPGatewayConfig holds SIP-to-stream gateway settings.
type SIPGatewayConfig struct {
Enabled bool `yaml:"enabled"`
StreamPrefix string `yaml:"stream_prefix"` // stream key prefix (default "sip")
RTPPortRange []int `yaml:"rtp_port_range"` // [min, max] for RTP port allocation
Codecs []string `yaml:"codecs"` // preferred codecs (default: opus, PCMA, PCMU)
MaxCalls int `yaml:"max_calls"` // max concurrent calls (default 100)
Enabled bool `yaml:"enabled"`
StreamPrefix string `yaml:"stream_prefix"` // stream key prefix (default "sip")
RTPPortRange []int `yaml:"rtp_port_range"` // [min, max] for RTP port allocation
Codecs []string `yaml:"codecs"` // preferred codecs (default: opus, PCMA, PCMU)
MaxCalls int `yaml:"max_calls"` // max concurrent calls (default 100)
MaxLabSessions int `yaml:"max_lab_sessions"` // max active local protocol labs (default 16)
}

// GB28181Config holds GB28181 module settings.
type GB28181Config struct {
Enabled bool `yaml:"enabled"`
StreamPrefix string `yaml:"stream_prefix"`
RTPPortRange []int `yaml:"rtp_port_range"`
MaxLabSessions int `yaml:"max_lab_sessions"` // max active local protocol labs (default 16)
SSRC SSRCConfig `yaml:"ssrc"`
Keepalive KeepaliveConfig `yaml:"keepalive"`
AutoInvite bool `yaml:"auto_invite"`
Expand Down Expand Up @@ -283,15 +292,22 @@ type SlowConsumerConfig struct {

// StreamConfig holds stream-level settings.
type StreamConfig struct {
GOPCache bool `yaml:"gop_cache"`
GOPCacheNum int `yaml:"gop_cache_num"`
RingBufferSize int `yaml:"ring_buffer_size"`
IdleTimeout time.Duration `yaml:"idle_timeout"`
NoPublisherTimeout time.Duration `yaml:"no_publisher_timeout"`
SlowConsumer SlowConsumerConfig `yaml:"slow_consumer"`
Simulcast SimulcastConfig `yaml:"simulcast"`
Feedback FeedbackConfig `yaml:"feedback"`
}
GOPCache bool `yaml:"gop_cache"`
GOPCacheNum int `yaml:"gop_cache_num"`
GOPCacheMaxFrames int `yaml:"gop_cache_max_frames"`
GOPCacheMaxDuration time.Duration `yaml:"gop_cache_max_duration"`
GOPCacheMaxBytes int64 `yaml:"gop_cache_max_bytes"`
RingBufferSize int `yaml:"ring_buffer_size"`
IdleTimeout time.Duration `yaml:"idle_timeout"`
NoPublisherTimeout time.Duration `yaml:"no_publisher_timeout"`
SlowConsumer SlowConsumerConfig `yaml:"slow_consumer"`
Simulcast SimulcastConfig `yaml:"simulcast"`
Feedback FeedbackConfig `yaml:"feedback"`
}

// DefaultGOPCacheMaxFrames is the defensive cardinality bound used when a
// stream is constructed directly without passing through config validation.
const DefaultGOPCacheMaxFrames = 300

// SimulcastConfig holds simulcast layer settings.
type SimulcastConfig struct {
Expand Down Expand Up @@ -499,9 +515,12 @@ type DVRConfig struct {

// MetricsConfig holds Prometheus metrics settings.
type MetricsConfig struct {
Enabled bool `yaml:"enabled"`
Listen string `yaml:"listen"`
Path string `yaml:"path"`
Enabled bool `yaml:"enabled"`
Listen string `yaml:"listen"`
Path string `yaml:"path"`
StreamDetail bool `yaml:"stream_detail"`
StreamDetailLimit int `yaml:"stream_detail_limit"`
StreamDetailAllowlist []string `yaml:"stream_detail_allowlist"`
}

// APIConfig holds the management API settings.
Expand Down
124 changes: 119 additions & 5 deletions config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"fmt"
"os"
"path/filepath"
"strings"
"testing"
"time"
)
Expand Down Expand Up @@ -72,6 +73,109 @@ func TestLoadConfigDefaults(t *testing.T) {
if cfg.Stream.RingBufferSize != 1024 {
t.Errorf("expected default ring_buffer_size 1024, got %d", cfg.Stream.RingBufferSize)
}
if cfg.Stream.GOPCacheMaxFrames <= 0 || cfg.Stream.GOPCacheMaxDuration <= 0 || cfg.Stream.GOPCacheMaxBytes <= 0 {
t.Fatalf("GOP cache bounds must have positive defaults: frames=%d duration=%s bytes=%d", cfg.Stream.GOPCacheMaxFrames, cfg.Stream.GOPCacheMaxDuration, cfg.Stream.GOPCacheMaxBytes)
}
if cfg.SIP.Gateway.MaxLabSessions != 16 || cfg.GB28181.MaxLabSessions != 16 {
t.Fatalf("expected default protocol lab session ceilings of 16, got SIP=%d GB=%d", cfg.SIP.Gateway.MaxLabSessions, cfg.GB28181.MaxLabSessions)
}
}

func TestLoadConfigGOPCacheBounds(t *testing.T) {
path := filepath.Join(t.TempDir(), "config.yaml")
doc := "stream:\n gop_cache_max_frames: 7\n gop_cache_max_duration: 2s\n gop_cache_max_bytes: 8192\n"
if err := os.WriteFile(path, []byte(doc), 0o600); err != nil {
t.Fatal(err)
}
cfg, err := Load(path)
if err != nil {
t.Fatal(err)
}
if cfg.Stream.GOPCacheMaxFrames != 7 || cfg.Stream.GOPCacheMaxDuration != 2*time.Second || cfg.Stream.GOPCacheMaxBytes != 8192 {
t.Fatalf("loaded GOP bounds = frames=%d duration=%s bytes=%d", cfg.Stream.GOPCacheMaxFrames, cfg.Stream.GOPCacheMaxDuration, cfg.Stream.GOPCacheMaxBytes)
}
}

func TestValidateRejectsNonPositiveRingBufferSize(t *testing.T) {
cfg := Defaults()
cfg.Stream.RingBufferSize = 0
if err := Validate(cfg); err == nil || !strings.Contains(err.Error(), "stream.ring_buffer_size") {
t.Fatalf("Validate() error = %v, want ring buffer size rejection", err)
}
}

func TestValidateMetricsStreamDetailLimitAllowsZeroAndRejectsNegative(t *testing.T) {
zero := Defaults()
zero.Metrics.StreamDetailLimit = 0
if err := Validate(zero); err != nil {
t.Fatalf("Validate() rejected metrics.stream_detail_limit=0: %v", err)
}

negative := Defaults()
negative.Metrics.StreamDetailLimit = -1
if err := Validate(negative); err == nil || !strings.Contains(err.Error(), "metrics.stream_detail_limit must not be negative") {
t.Fatalf("Validate() error = %v, want negative metrics stream detail limit rejection", err)
}
}

func TestValidateRejectsInvalidTrustedProxy(t *testing.T) {
cfg := Defaults()
cfg.Limits.RateLimit.TrustedProxies = []string{"127.0.0.1", "not-a-network"}
if err := Validate(cfg); err == nil || !strings.Contains(err.Error(), "limits.rate_limit.trusted_proxies[1]") {
t.Fatalf("Validate() error = %v, want invalid trusted proxy rejection", err)
}
}

func TestValidateRejectsUnboundedEnabledGOPCache(t *testing.T) {
cfg := Defaults()
cfg.Stream.GOPCache = true
cfg.Stream.GOPCacheMaxFrames = 0
cfg.Stream.GOPCacheMaxBytes = 0
if err := Validate(cfg); err == nil || !strings.Contains(err.Error(), "gop_cache_max_frames or stream.gop_cache_max_bytes") {
t.Fatalf("Validate() error = %v, want hard GOP bound rejection", err)
}
}

func TestValidateRecordFormatAndMaxSize(t *testing.T) {
for _, format := range []string{"flv", "fmp4", "mp4", "ts", "hls", " HLS "} {
cfg := Defaults()
cfg.Record.Format = format
if err := Validate(cfg); err != nil {
t.Errorf("Validate() rejected record format %q: %v", format, err)
}
}

for _, maxSize := range []string{"", "0", "0MB", "512KB", "1GB", " 256mb "} {
cfg := Defaults()
cfg.Record.Segment.MaxSize = maxSize
if err := Validate(cfg); err != nil {
t.Errorf("Validate() rejected record.segment.max_size %q: %v", maxSize, err)
}
}

for _, test := range []struct {
name string
field string
value string
}{
{name: "format", field: "record.format", value: "webm"},
{name: "fractional size", field: "record.segment.max_size", value: "1.5MB"},
{name: "negative size", field: "record.segment.max_size", value: "-1MB"},
{name: "unknown suffix", field: "record.segment.max_size", value: "1TB"},
{name: "overflow size", field: "record.segment.max_size", value: "9223372036854775808B"},
} {
t.Run(test.name, func(t *testing.T) {
cfg := Defaults()
if test.field == "record.format" {
cfg.Record.Format = test.value
} else {
cfg.Record.Segment.MaxSize = test.value
}
if err := Validate(cfg); err == nil || !strings.Contains(err.Error(), test.field) {
t.Fatalf("Validate() error = %v, want %s rejection", err, test.field)
}
})
}
}

func TestLoadConfigEnvExpansion(t *testing.T) {
Expand Down Expand Up @@ -229,7 +333,9 @@ http_stream:
container: "ts"
`
tmpFile := filepath.Join(t.TempDir(), "test.yaml")
os.WriteFile(tmpFile, []byte(yaml), 0644)
if err := os.WriteFile(tmpFile, []byte(yaml), 0600); err != nil {
t.Fatal(err)
}
cfg, err := Load(tmpFile)
if err != nil {
t.Fatalf("Load: %v", err)
Expand Down Expand Up @@ -257,7 +363,9 @@ http_stream:
listen: ":8080"
`
tmpFile := filepath.Join(t.TempDir(), "test.yaml")
os.WriteFile(tmpFile, []byte(yaml), 0644)
if err := os.WriteFile(tmpFile, []byte(yaml), 0600); err != nil {
t.Fatal(err)
}
cfg, err := Load(tmpFile)
if err != nil {
t.Fatalf("Load: %v", err)
Expand Down Expand Up @@ -303,7 +411,9 @@ http_stream:
container: "mpegts"
`
tmpFile := filepath.Join(t.TempDir(), "test.yaml")
os.WriteFile(tmpFile, []byte(yaml), 0644)
if err := os.WriteFile(tmpFile, []byte(yaml), 0600); err != nil {
t.Fatal(err)
}
cfg, err := Load(tmpFile)
if err != nil {
t.Fatalf("Load: %v", err)
Expand All @@ -320,7 +430,9 @@ http_stream:
container: "mpeg-ts"
`
tmpFile := filepath.Join(t.TempDir(), "test.yaml")
os.WriteFile(tmpFile, []byte(yaml), 0644)
if err := os.WriteFile(tmpFile, []byte(yaml), 0600); err != nil {
t.Fatal(err)
}
cfg, err := Load(tmpFile)
if err != nil {
t.Fatalf("Load: %v", err)
Expand Down Expand Up @@ -401,7 +513,9 @@ func TestLoadConfigInvalidPath(t *testing.T) {

func TestLoadConfigInvalidYAML(t *testing.T) {
tmpFile := filepath.Join(t.TempDir(), "bad.yaml")
os.WriteFile(tmpFile, []byte("{{invalid yaml"), 0644)
if err := os.WriteFile(tmpFile, []byte("{{invalid yaml"), 0600); err != nil {
t.Fatal(err)
}
_, err := Load(tmpFile)
if err == nil {
t.Error("expected error for invalid YAML")
Expand Down
20 changes: 15 additions & 5 deletions config/loader.go
Original file line number Diff line number Diff line change
Expand Up @@ -203,11 +203,15 @@ func defaults() *Config {
SIP: SIPConfig{
Listen: ":5060",
Transport: []string{"udp", "tcp"},
Gateway: SIPGatewayConfig{MaxLabSessions: 16},
},
Stream: StreamConfig{
GOPCache: true,
GOPCacheNum: 1,
RingBufferSize: 1024,
GOPCache: true,
GOPCacheNum: 1,
GOPCacheMaxFrames: DefaultGOPCacheMaxFrames,
GOPCacheMaxDuration: 10 * time.Second,
GOPCacheMaxBytes: 32 * 1024 * 1024,
RingBufferSize: 1024,
SlowConsumer: SlowConsumerConfig{
Enabled: true,
LagWarnRatio: 0.5,
Expand Down Expand Up @@ -236,13 +240,18 @@ func defaults() *Config {
Audit: AuditConfig{MaxEntries: 1000},
},
Metrics: MetricsConfig{
Listen: ":9090",
Path: "/metrics",
Listen: ":9090",
Path: "/metrics",
StreamDetailLimit: 100,
},
Runtime: RuntimeConfig{
Source: "file",
PollInterval: 30 * time.Second,
LoadTimeout: 10 * time.Second,
File: RuntimeFileSourceConfig{MaxBytes: DefaultRuntimeSourceMaxBytes},
HTTP: RuntimeHTTPSourceConfig{MaxBytes: DefaultRuntimeSourceMaxBytes},
Consul: RuntimeConsulSourceConfig{MaxBytes: DefaultRuntimeSourceMaxBytes},
Redis: RuntimeRedisSourceConfig{MaxBytes: DefaultRuntimeSourceMaxBytes},
},
Record: RecordConfig{Format: "fmp4"},
DVR: DVRConfig{
Expand All @@ -252,6 +261,7 @@ func defaults() *Config {
SegmentDuration: 6 * time.Second,
CleanupInterval: 30 * time.Second,
},
GB28181: GB28181Config{MaxLabSessions: 16},
Cluster: ClusterConfig{
SRT: ClusterSRTConfig{
Latency: 120 * time.Millisecond,
Expand Down
40 changes: 40 additions & 0 deletions config/runtime/error_redaction.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
package runtime

import (
"net/url"
"regexp"
"strings"
)

var errorURLPattern = regexp.MustCompile(`(?i)(?:https?|rediss?|consul)://[^\s"'<>]+`)

// RedactError removes URL credentials, query values, fragments, and line
// breaks before an error crosses a logging or management API boundary.
func RedactError(err error) string {
if err == nil {
return ""
}
message := strings.NewReplacer("\r", " ", "\n", " ").Replace(err.Error())
return errorURLPattern.ReplaceAllStringFunc(message, redactErrorURL)
}

func redactErrorURL(raw string) string {
trailing := ""
for len(raw) > 0 && strings.ContainsRune(".,;:)", rune(raw[len(raw)-1])) {
trailing = string(raw[len(raw)-1]) + trailing
raw = raw[:len(raw)-1]
}
parsed, err := url.Parse(raw)
if err != nil || parsed.Scheme == "" || parsed.Host == "" {
return "REDACTED_URL" + trailing
}
parsed.User = nil
if parsed.Path != "" && parsed.Path != "/" {
parsed.Path = "/REDACTED"
parsed.RawPath = ""
}
parsed.RawQuery = "__liveforge_redacted__=1"
parsed.ForceQuery = false
parsed.Fragment = ""
return parsed.String() + trailing
}
Loading