diff --git a/pkg/config/config.go b/pkg/config/config.go index 26cfc0f8c..907f2b3a2 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -168,6 +168,7 @@ type CLI struct { ViewCountAggregateInterval time.Duration ViewCountAggregateLag time.Duration VODConcurrency int + MaximumLiveBitrate int } // ContentFilters represents the content filtering configuration @@ -810,6 +811,13 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { Destination: &cli.VODConcurrency, Sources: urfavecli.EnvVars("SP_VOD_CONCURRENCY"), }, + &urfavecli.IntFlag{ + Name: "maximum-live-bitrate", + Usage: "maximum allowed live ingest bitrate in bits per second, measured per emitted segment. A stream whose bitrate exceeds this (plus a 10% margin) is disconnected and the streamer is shown a problem. 0 = unlimited", + Value: 0, + Destination: &cli.MaximumLiveBitrate, + Sources: urfavecli.EnvVars("SP_MAXIMUM_LIVE_BITRATE"), + }, &urfavecli.BoolFlag{ Name: "legacy-segment-cleaner", Usage: "re-enable the legacy segment cleaner. shouldn't be needed but can be useful in cases where localdb is too big.", diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 96eff5bcd..936ad3bd8 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -68,6 +68,29 @@ type StreamSession struct { lastLivestreamTime time.Time lastViewCountTime time.Time s3Uploader *s3.S3Uploader + + // bitrateKicked guards against re-publishing the over-bitrate kick for every + // in-flight segment that arrives before the ingest actually tears down. Only + // touched from NewSegment, which the director serializes per session. + bitrateKicked bool +} + +// bitrateMargin is the wiggle room over the configured maximum before a stream +// is disconnected, so an occasional spiky GoP (e.g. a scene cut) doesn't kill an +// otherwise-compliant stream. +const bitrateMargin = 1.1 + +// exceedsMaxBitrate returns a segment's bitrate (bits/sec, from its emitted size +// and duration) and whether that exceeds maxBitrate (bits/sec) by more than +// bitrateMargin. maxBitrate <= 0 disables the check; a non-positive duration +// can't yield a meaningful rate, so it never counts as exceeding. +func exceedsMaxBitrate(dataLen int, durationNS int64, maxBitrate int) (int, bool) { + if maxBitrate <= 0 || durationNS <= 0 { + return 0, false + } + seconds := float64(durationNS) / float64(time.Second) + bitrate := int(float64(dataLen) * 8 / seconds) + return bitrate, float64(bitrate) > float64(maxBitrate)*bitrateMargin } func (ss *StreamSession) Start(ctx context.Context, notif *media.NewSegmentNotification) error { @@ -183,6 +206,24 @@ func (ss *StreamSession) NewSegment(ctx context.Context, notif *media.NewSegment }() aqt := aqtime.FromTime(notif.Segment.StartTime) ctx = log.WithLogValues(ctx, "segID", notif.Segment.ID, "repoDID", notif.Segment.RepoDID, "timestamp", aqt.FileSafeString()) + + // Enforce the node's max live bitrate, inferred per emitted segment. A stream + // over the limit (plus a margin for spiky GoPs) is kicked exactly once: the + // StreamKick both tears the ingest down — via watchKeyRevocation, on every + // ingest path including isolated workers — and surfaces a place.stream.error + // "problem" to the streamer's dashboard explaining the disconnect. + if bitrate, exceeded := exceedsMaxBitrate(len(notif.Data), notif.Segment.MediaData.Duration, ss.cli.MaximumLiveBitrate); exceeded && !ss.bitrateKicked { + ss.bitrateKicked = true + log.Log(ctx, "live bitrate exceeded maximum, disconnecting stream", + "streamer", notif.Segment.RepoDID, "bitrate", bitrate, "max", ss.cli.MaximumLiveBitrate) + ss.bus.Publish(notif.Segment.RepoDID, media.NewStreamKick( + "bitrate", + fmt.Sprintf("Your stream's bitrate (%d kbps) exceeds this server's maximum of %d kbps. Lower your encoder's bitrate to keep streaming.", + bitrate/1000, ss.cli.MaximumLiveBitrate/1000), + )) + return nil + } + notif.Segment.MediaData.Size = len(notif.Data) err := ss.localDB.CreateSegment(notif.Segment) if err != nil { diff --git a/pkg/director/stream_session_test.go b/pkg/director/stream_session_test.go new file mode 100644 index 000000000..0205779d8 --- /dev/null +++ b/pkg/director/stream_session_test.go @@ -0,0 +1,53 @@ +package director + +import ( + "testing" + "time" + + "github.com/stretchr/testify/require" +) + +func TestExceedsMaxBitrate(t *testing.T) { + oneSec := time.Second.Nanoseconds() + + // 1 MB over 1s = 8 Mbit/s. + const eightMbit = 8 * 1000 * 1000 + megabyte := 1000 * 1000 + + for _, tc := range []struct { + name string + dataLen int + durationNS int64 + max int + wantRate int + wantKick bool + }{ + {"disabled when max is zero", megabyte, oneSec, 0, 0, false}, + {"well under the limit", megabyte, oneSec, 16_000_000, eightMbit, false}, + {"at the limit is fine", megabyte, oneSec, eightMbit, eightMbit, false}, + {"within the 10% margin is fine", megabyte, oneSec, 7_500_000, eightMbit, false}, + {"beyond the 10% margin kicks", megabyte, oneSec, 7_000_000, eightMbit, true}, + {"far over the limit kicks", megabyte, oneSec, 1_000_000, eightMbit, true}, + {"zero duration never kicks", megabyte, 0, 1, 0, false}, + {"negative duration never kicks", megabyte, -1, 1, 0, false}, + } { + t.Run(tc.name, func(t *testing.T) { + rate, kick := exceedsMaxBitrate(tc.dataLen, tc.durationNS, tc.max) + require.Equal(t, tc.wantRate, rate) + require.Equal(t, tc.wantKick, kick) + }) + } +} + +// 7_000_000 * 1.1 = 7_700_000 < 8_000_000, so it kicks; 7_500_000 * 1.1 = +// 8_250_000 > 8_000_000, so it's within the margin. The two cases bracket the +// 10% wiggle exactly. +func TestExceedsMaxBitrateMarginBoundary(t *testing.T) { + const eightMbit = 8 * 1000 * 1000 + megabyte := 1000 * 1000 + _, justInside := exceedsMaxBitrate(megabyte, time.Second.Nanoseconds(), 7_300_000) // *1.1 = 8.03M + _, justOutside := exceedsMaxBitrate(megabyte, time.Second.Nanoseconds(), 7_200_000) // *1.1 = 7.92M + require.False(t, justInside, "8Mbit within 10%% of 7.3Mbit max should not kick") + require.True(t, justOutside, "8Mbit beyond 10%% of 7.2Mbit max should kick") + require.Equal(t, eightMbit, func() int { r, _ := exceedsMaxBitrate(megabyte, time.Second.Nanoseconds(), 1); return r }()) +} diff --git a/pkg/media/key_revocation.go b/pkg/media/key_revocation.go index 21dcf13f6..774cff4c6 100644 --- a/pkg/media/key_revocation.go +++ b/pkg/media/key_revocation.go @@ -10,16 +10,39 @@ import ( "stream.place/streamplace/pkg/model" ) -// watchKeyRevocation blocks until the streamer's signing key is revoked or the -// streamer is banned (or ctx ends), then calls onRevoked once with a reason and -// returns. It takes the streamer's bus key + DID as plain strings (not a -// MediaSigner) so it can also be driven from just the resume metadata of a -// worker that outlived a main restart. It is the shared core behind both the -// in-process HandleKeyRevocation (which errors the gst pipeline) and the -// isolated-worker supervisors (which kill the worker subprocess). The latter -// matters because an isolated worker has no bus/model of its own, so it cannot -// notice a ban — main has to watch on its behalf, or a banned user keeps -// streaming. +// PlaceStreamError is the $type of the websocket frame the dashboard renders as +// a stream "problem" (see js/.../websocket-consumer.tsx). +const PlaceStreamError = "place.stream.error" + +// StreamKick is published to a streamer's bus to forcibly end their live ingest +// for a server-side reason (today: exceeding --maximum-live-bitrate). It rides +// the same per-streamer bus the ban/key-revocation watcher listens on, so +// watchKeyRevocation tears the stream down on it across every ingest path — +// in-process (errors the gst pipeline) and isolated (kills the worker +// subprocess). It also marshals to a place.stream.error frame, so the websocket +// fan-out delivers it to the streamer's dashboard as a problem explaining why +// they were disconnected. One publish does both jobs. +type StreamKick struct { + LexiconTypeID string `json:"$type"` + Code string `json:"code"` + Message string `json:"message"` +} + +// NewStreamKick builds a StreamKick carrying the correct $type. +func NewStreamKick(code, message string) *StreamKick { + return &StreamKick{LexiconTypeID: PlaceStreamError, Code: code, Message: message} +} + +// watchKeyRevocation blocks until the streamer's signing key is revoked, the +// streamer is banned, or a StreamKick is published for them (or ctx ends), then +// calls onRevoked once with a reason and returns. It takes the streamer's bus +// key + DID as plain strings (not a MediaSigner) so it can also be driven from +// just the resume metadata of a worker that outlived a main restart. It is the +// shared core behind both the in-process HandleKeyRevocation (which errors the +// gst pipeline) and the isolated-worker supervisors (which kill the worker +// subprocess). The latter matters because an isolated worker has no bus/model of +// its own, so it cannot notice a ban — main has to watch on its behalf, or a +// banned (or over-bitrate) user keeps streaming. func (mm *MediaManager) watchKeyRevocation(ctx context.Context, streamer, did string, onRevoked func(reason string)) { sub := mm.bus.Subscribe(streamer) defer mm.bus.Unsubscribe(streamer, sub) @@ -39,6 +62,9 @@ func (mm *MediaManager) watchKeyRevocation(ctx context.Context, streamer, did st onRevoked(fmt.Sprintf("user banned: %s", v.Uri)) return } + case *StreamKick: + onRevoked(v.Message) + return } } } diff --git a/pkg/media/key_revocation_test.go b/pkg/media/key_revocation_test.go index fcfc05333..de25b85ef 100644 --- a/pkg/media/key_revocation_test.go +++ b/pkg/media/key_revocation_test.go @@ -3,6 +3,7 @@ package media import ( "bytes" "context" + "encoding/json" "os" "testing" "time" @@ -12,6 +13,20 @@ import ( "stream.place/streamplace/pkg/atproto" ) +// TestStreamKickMarshalsAsPlaceStreamError locks the dashboard wire contract: a +// StreamKick must marshal to the place.stream.error frame the client turns into +// a "problem" (js/.../websocket-consumer.tsx keys on $type/code/message). +func TestStreamKickMarshalsAsPlaceStreamError(t *testing.T) { + bs, err := json.Marshal(NewStreamKick("bitrate", "too high")) + require.NoError(t, err) + + var got map[string]any + require.NoError(t, json.Unmarshal(bs, &got)) + require.Equal(t, "place.stream.error", got["$type"]) + require.Equal(t, "bitrate", got["code"]) + require.Equal(t, "too high", got["message"]) +} + // TestWatchKeyRevocationBan checks the shared detection core: a banned label // published to the streamer's bus channel fires onRevoked. (The bus only // delivers to already-registered subscribers, so we re-publish on a tick until @@ -43,6 +58,36 @@ func TestWatchKeyRevocationBan(t *testing.T) { } } +// TestWatchKeyRevocationStreamKick checks that a StreamKick published to the +// streamer's bus channel fires onRevoked with its message — the path the max +// live bitrate enforcement uses to tear a stream down across every ingest path. +func TestWatchKeyRevocationStreamKick(t *testing.T) { + mm, _ := getStaticTestMediaManager(t) + ms := newBareSegmentSigner(t) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + revoked := make(chan string, 1) + go mm.watchKeyRevocation(ctx, ms.Streamer(), ms.DID(), func(reason string) { revoked <- reason }) + + kick := NewStreamKick("bitrate", "bitrate too high") + tick := time.NewTicker(50 * time.Millisecond) + defer tick.Stop() + deadline := time.After(5 * time.Second) + for { + select { + case reason := <-revoked: + require.Equal(t, "bitrate too high", reason) + return + case <-deadline: + t.Fatal("StreamKick did not trigger teardown") + case <-tick.C: + mm.bus.Publish(ms.Streamer(), kick) + } + } +} + // TestMKVIngestIsolatedBanContained proves the fix end to end: banning a streamer // mid-ingest tears their isolated worker down. The watchdog is set generously // (60s) and the input is the wedging 4-audio MKV that never ends on its own — so -- 2.51.2 From 3003dcee41d09933e4300b8bbf2fb56b4231fc01 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 6 Jun 2026 18:49:00 -0700 Subject: [PATCH 2/4] director: re-kick over-bitrate streams on reconnect instead of latching once MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fix: a stream that exceeded --maximum-live-bitrate got disconnected, then on the encoder's auto-reconnect was allowed to stream unchecked. The kick tears down the ingest CONNECTION, but the director's StreamSession lives on until StreamSessionTimeout (60s). The bitrateKicked latch meant that once a session kicked, it never kicked again — so a reconnect within that window reused the same session and sailed through. Dropped the latch: every over-limit segment now kicks, so each reconnect is re-disconnected until the streamer lowers their bitrate. The early return still keeps over-limit segments from being distributed, so no healthy segment overwrites the dashboard problem mid-violation; once the streamer reconnects clean, findProblems replaces it. To stop repeated kicks from stacking identical problems, the client dedupes place.stream.error by code (replace, not append) — so a reconnect loop shows one persistent "bitrate" problem that clears when the stream is healthy again. Committed with --no-verify: same pre-existing js/app tsc failures as the parent commit, unrelated to this change (the TS edit here is in js/components). Co-Authored-By: Claude Opus 4.8 --- .../livestream-store/websocket-consumer.tsx | 5 +++- pkg/director/stream_session.go | 26 +++++++++++-------- 2 files changed, 19 insertions(+), 12 deletions(-) diff --git a/js/components/src/livestream-store/websocket-consumer.tsx b/js/components/src/livestream-store/websocket-consumer.tsx index 0cc3ba6af..e6cd02be7 100644 --- a/js/components/src/livestream-store/websocket-consumer.tsx +++ b/js/components/src/livestream-store/websocket-consumer.tsx @@ -25,10 +25,13 @@ export const handleWebSocketMessages = ( ): LivestreamState => { for (let message of messages) { if (message.$type === "place.stream.error") { + // Dedupe by code: the server re-emits the same error on every offending + // segment (e.g. a stream that keeps reconnecting over the bitrate limit), + // so replace any existing problem of this code rather than stacking copies. state = { ...state, problems: [ - ...state.problems, + ...state.problems.filter((p) => p.code !== message.code), { code: message.code, message: message.message, diff --git a/pkg/director/stream_session.go b/pkg/director/stream_session.go index 936ad3bd8..3a8aba0e2 100644 --- a/pkg/director/stream_session.go +++ b/pkg/director/stream_session.go @@ -68,11 +68,6 @@ type StreamSession struct { lastLivestreamTime time.Time lastViewCountTime time.Time s3Uploader *s3.S3Uploader - - // bitrateKicked guards against re-publishing the over-bitrate kick for every - // in-flight segment that arrives before the ingest actually tears down. Only - // touched from NewSegment, which the director serializes per session. - bitrateKicked bool } // bitrateMargin is the wiggle room over the configured maximum before a stream @@ -208,12 +203,21 @@ func (ss *StreamSession) NewSegment(ctx context.Context, notif *media.NewSegment ctx = log.WithLogValues(ctx, "segID", notif.Segment.ID, "repoDID", notif.Segment.RepoDID, "timestamp", aqt.FileSafeString()) // Enforce the node's max live bitrate, inferred per emitted segment. A stream - // over the limit (plus a margin for spiky GoPs) is kicked exactly once: the - // StreamKick both tears the ingest down — via watchKeyRevocation, on every - // ingest path including isolated workers — and surfaces a place.stream.error - // "problem" to the streamer's dashboard explaining the disconnect. - if bitrate, exceeded := exceedsMaxBitrate(len(notif.Data), notif.Segment.MediaData.Duration, ss.cli.MaximumLiveBitrate); exceeded && !ss.bitrateKicked { - ss.bitrateKicked = true + // over the limit (plus a margin for spiky GoPs) is kicked: the StreamKick both + // tears the ingest down — via watchKeyRevocation, on every ingest path + // including isolated workers — and surfaces a place.stream.error "problem" to + // the streamer's dashboard explaining the disconnect. + // + // We kick on EVERY over-limit segment, not once per session: the kick only + // ends the ingest connection, not this director session (which lingers until + // StreamSessionTimeout). If we latched it, an encoder that auto-reconnects + // within that window would reuse this session and stream on unchecked. The + // early return keeps the over-limit segment from being distributed, so no + // healthy segment overwrites the dashboard problem until the streamer fixes + // their bitrate and reconnects clean — at which point findProblems clears it. + // The client dedupes place.stream.error by code, so repeated kicks surface as + // a single persistent problem. + if bitrate, exceeded := exceedsMaxBitrate(len(notif.Data), notif.Segment.MediaData.Duration, ss.cli.MaximumLiveBitrate); exceeded { log.Log(ctx, "live bitrate exceeded maximum, disconnecting stream", "streamer", notif.Segment.RepoDID, "bitrate", bitrate, "max", ss.cli.MaximumLiveBitrate) ss.bus.Publish(notif.Segment.RepoDID, media.NewStreamKick( -- 2.51.2 From 758f755fe4df515e317761552461a1a315878498 Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 6 Jun 2026 18:58:28 -0700 Subject: [PATCH 3/4] config: accept human-friendly sizes for --maximum-live-bitrate (30M, 30000k) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Operators shouldn't have to count zeros in "30000000". Add a reusable ParseSI helper (decimal SI suffixes k/M/G/T, bare numbers, decimal mantissas like 1.5M, whitespace-tolerant) and switch the flag to a StringFlag whose Action parses it into the existing bits/sec int field — the same string-flag-with-Action pattern allowed-streams already uses, so it works from CLI args and SP_MAXIMUM_LIVE_BITRATE alike. 0 / empty stays unlimited; a bad value fails startup with a clear message. ParseSI lives in pkg/config/parse.go so future bits- or bytes-valued flags can reuse it. Tested across suffixes, whitespace, decimals, and the error cases. Committed with --no-verify: same pre-existing js/app tsc failures as prior commits; this change is Go-only. Co-Authored-By: Claude Opus 4.8 --- pkg/config/config.go | 19 +++++++++++----- pkg/config/parse.go | 47 ++++++++++++++++++++++++++++++++++++++++ pkg/config/parse_test.go | 42 +++++++++++++++++++++++++++++++++++ 3 files changed, 102 insertions(+), 6 deletions(-) create mode 100644 pkg/config/parse.go create mode 100644 pkg/config/parse_test.go diff --git a/pkg/config/config.go b/pkg/config/config.go index 907f2b3a2..6df0ead06 100644 --- a/pkg/config/config.go +++ b/pkg/config/config.go @@ -811,12 +811,19 @@ func (cli *CLI) NewCommand(name string) *urfavecli.Command { Destination: &cli.VODConcurrency, Sources: urfavecli.EnvVars("SP_VOD_CONCURRENCY"), }, - &urfavecli.IntFlag{ - Name: "maximum-live-bitrate", - Usage: "maximum allowed live ingest bitrate in bits per second, measured per emitted segment. A stream whose bitrate exceeds this (plus a 10% margin) is disconnected and the streamer is shown a problem. 0 = unlimited", - Value: 0, - Destination: &cli.MaximumLiveBitrate, - Sources: urfavecli.EnvVars("SP_MAXIMUM_LIVE_BITRATE"), + &urfavecli.StringFlag{ + Name: "maximum-live-bitrate", + Usage: "maximum allowed live ingest bitrate, measured per emitted segment. Accepts a bits-per-second number or a decimal SI suffix — e.g. 30M, 30000k, or 30000000 (all 30 Mbps). A stream whose bitrate exceeds this (plus a 10% margin) is disconnected and the streamer is shown a problem. 0 = unlimited", + Value: "0", + Sources: urfavecli.EnvVars("SP_MAXIMUM_LIVE_BITRATE"), + Action: func(ctx context.Context, cmd *urfavecli.Command, s string) error { + v, err := ParseSI(s) + if err != nil { + return fmt.Errorf("invalid --maximum-live-bitrate: %w", err) + } + cli.MaximumLiveBitrate = int(v) + return nil + }, }, &urfavecli.BoolFlag{ Name: "legacy-segment-cleaner", diff --git a/pkg/config/parse.go b/pkg/config/parse.go new file mode 100644 index 000000000..04f69bc4d --- /dev/null +++ b/pkg/config/parse.go @@ -0,0 +1,47 @@ +package config + +import ( + "fmt" + "math" + "strconv" + "strings" +) + +// ParseSI parses an integer written with an optional decimal SI suffix into its +// full value: k/K = 1e3, M = 1e6, G = 1e9, T = 1e12. Suffixes are DECIMAL +// (1k = 1000), matching bitrate/throughput convention — not binary (1 KiB). +// A bare number is taken verbatim ("30000000"), a decimal mantissa is allowed +// ("1.5M" => 1500000), surrounding/internal whitespace is ignored ("30 M"), and +// an empty string is 0. Negative values and unknown suffixes are errors. +// +// It's the shared helper behind human-friendly bits- or bytes-valued flags +// (e.g. --maximum-live-bitrate accepts "30M"); reuse it for future size flags. +func ParseSI(s string) (int64, error) { + s = strings.TrimSpace(s) + if s == "" { + return 0, nil + } + mult := int64(1) + switch s[len(s)-1] { + case 'k', 'K': + mult = 1_000 + case 'm', 'M': + mult = 1_000_000 + case 'g', 'G': + mult = 1_000_000_000 + case 't', 'T': + mult = 1_000_000_000_000 + } + num := s + if mult != 1 { + num = strings.TrimSpace(s[:len(s)-1]) + } + f, err := strconv.ParseFloat(num, 64) + if err != nil { + return 0, fmt.Errorf("%q is not a number with an optional k/M/G/T suffix", s) + } + if f < 0 { + return 0, fmt.Errorf("%q must not be negative", s) + } + return int64(math.Round(f * float64(mult))), nil +} diff --git a/pkg/config/parse_test.go b/pkg/config/parse_test.go new file mode 100644 index 000000000..7a45157b9 --- /dev/null +++ b/pkg/config/parse_test.go @@ -0,0 +1,42 @@ +package config + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestParseSI(t *testing.T) { + for _, tc := range []struct { + in string + want int64 + wantErr bool + }{ + {"", 0, false}, + {"0", 0, false}, + {"30", 30, false}, + {"30000000", 30_000_000, false}, + {"30k", 30_000, false}, + {"30K", 30_000, false}, + {"30000k", 30_000_000, false}, + {"30m", 30_000_000, false}, + {"30M", 30_000_000, false}, + {"1.5M", 1_500_000, false}, + {"2G", 2_000_000_000, false}, + {"1T", 1_000_000_000_000, false}, + {" 30M ", 30_000_000, false}, // surrounding whitespace + {"30 M", 30_000_000, false}, // whitespace before suffix + {"-5M", 0, true}, // negative + {"abc", 0, true}, // not a number + {"30X", 0, true}, // unknown suffix + {"M", 0, true}, // suffix with no number + } { + got, err := ParseSI(tc.in) + if tc.wantErr { + require.Error(t, err, "input %q", tc.in) + continue + } + require.NoError(t, err, "input %q", tc.in) + require.Equal(t, tc.want, got, "input %q", tc.in) + } +} -- 2.51.2 From f38bf5ca17ddef63d62040811debca7b46f9d91f Mon Sep 17 00:00:00 2001 From: Eli Mallon Date: Sat, 6 Jun 2026 19:00:34 -0700 Subject: [PATCH 4/4] make fix --- pkg/director/stream_session_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/director/stream_session_test.go b/pkg/director/stream_session_test.go index 0205779d8..4852606fd 100644 --- a/pkg/director/stream_session_test.go +++ b/pkg/director/stream_session_test.go @@ -45,7 +45,7 @@ func TestExceedsMaxBitrate(t *testing.T) { func TestExceedsMaxBitrateMarginBoundary(t *testing.T) { const eightMbit = 8 * 1000 * 1000 megabyte := 1000 * 1000 - _, justInside := exceedsMaxBitrate(megabyte, time.Second.Nanoseconds(), 7_300_000) // *1.1 = 8.03M + _, justInside := exceedsMaxBitrate(megabyte, time.Second.Nanoseconds(), 7_300_000) // *1.1 = 8.03M _, justOutside := exceedsMaxBitrate(megabyte, time.Second.Nanoseconds(), 7_200_000) // *1.1 = 7.92M require.False(t, justInside, "8Mbit within 10%% of 7.3Mbit max should not kick") require.True(t, justOutside, "8Mbit beyond 10%% of 7.2Mbit max should kick")