media

package
v0.8.3 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Aug 16, 2026 License: Apache-2.0 Imports: 11 Imported by: 0

Documentation

Overview

Package media 提供直播/视频服务的编排薄机制:多路流管理(Hub)、子进程监督 (Supervisor,如 ffmpeg 转码进程,含崩溃重启)与 OTel 运维指标。

边界(与框架"薄机制"一致):只做"管理与编排"这层通用苦活——注册表、生命周期、 路由、进程重启、指标埋点。**不做 policy**:每路流怎么建(hls.Stream 配置)、跑什么 命令(ffmpeg 参数)、转码档位、导出到哪(Prometheus/OTLP),全由调用方决定。

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Hub

type Hub[S Stream] struct {
	// contains filtered or unexported fields
}

Hub 管理多路直播流:streamKey → Session。解决"多路并发、重复推流、按 key 分发"。 按流类型参数化(Hub[*hls.Stream] 或 Hub[*hlsmux.Bridge] 等)。并发安全。 零值不可用,用 NewHub 构造。

func NewHub

func NewHub[S Stream](newStream func(key string) S, opts ...HubOption) *Hub[S]

NewHub 创建多路流管理器。newStream 是每路流的构造(policy:窗口、分片时长、存储/LL-HLS 等都在这里定);它同时决定了流类型 S。newStream 为 nil 会 panic。

func (*Hub[S]) Acquire

func (h *Hub[S]) Acquire(key string) (*Session[S], bool)

Acquire 为 key 创建一路流并注册。若该 key 已在推流,返回 (nil, false) 表示拒绝(防抢流); 成功返回 (session, true)。通常在 rtmp 的 PublishFunc 里调:拿到 session 就用它的 Stream 收流/分发;返回 false 时让 PublishFunc 返回 nil 拒绝这次推流。

func (*Hub[S]) Count

func (h *Hub[S]) Count() int

Count 返回当前在线流数。

func (*Hub[S]) Lookup

func (h *Hub[S]) Lookup(key string) (*Session[S], bool)

Lookup 查找一路流。

func (*Hub[S]) Metrics

func (h *Hub[S]) Metrics() *Metrics

Metrics 返回指标记录器(供调用方上报 ingest bytes / segment 等)。

func (*Hub[S]) Release

func (h *Hub[S]) Release(key string)

Release 结束一路流:取消 Session.Context(停掉绑定的后台任务)、收尾 Stream(Finish)、 从注册表移除。幂等。

func (*Hub[S]) ServeHTTP

func (h *Hub[S]) ServeHTTP(w http.ResponseWriter, r *http.Request)

ServeHTTP 路由 /{key}/… 到对应流的 Stream(去掉 key 前缀);无此流返回 404。

type HubOption

type HubOption func(*hubConfig)

HubOption 配置 Hub(与流类型无关的项)。

func WithBaseContext

func WithBaseContext(ctx context.Context) HubOption

WithBaseContext 设置各 Session Context 的父 context(默认 context.Background())。

func WithMetrics

func WithMetrics(m *Metrics) HubOption

WithMetrics 注入指标记录器(默认 NewMetrics(),基于 OTel 全局 Meter)。

type Metrics

type Metrics struct {
	// contains filtered or unexported fields
}

Metrics 是一组直播运维指标,基于 OTel 全局 MeterProvider(由 pkg/service/telemetry 配置)。未配置 telemetry 时全部为 no-op,零开销。按 stream 标签维度上报。

func NewMetrics

func NewMetrics() *Metrics

NewMetrics 从全局 Meter 创建指标(instrument 创建失败时对应项为 nil,记录时跳过)。

func (*Metrics) IngestBytes

func (m *Metrics) IngestBytes(ctx context.Context, key string, n int64)

IngestBytes 上报某路流采集入流量(由调用方在收到音视频时累加)。

func (*Metrics) Restart

func (m *Metrics) Restart(ctx context.Context, key string)

Restart 上报某路流的转码进程重启(Supervisor 内部调用)。

func (*Metrics) Segment

func (m *Metrics) Segment(ctx context.Context, key string)

Segment 上报某路流产出一个分片(由调用方在 Append/切片时调用)。

type Session

type Session[S Stream] struct {
	Key    string
	Stream S
	// contains filtered or unexported fields
}

Session 是一路直播流的运行时:一个 Stream + 一个随 Release 取消的 Context。 外部(如 ffmpeg Supervisor)把自己绑定到 Context() 上,即可在该路流结束时自动收尾。

func (*Session[S]) Context

func (s *Session[S]) Context() context.Context

Context 在该路流被 Release(结束)时取消。把 Supervisor 等后台任务绑到它上面即可随流停机。

type Stream

type Stream interface {
	http.Handler
	// Finish 在该路流被 Release 时调用,做收尾(如封为点播 / 关闭 muxer)。应幂等。
	Finish()
}

Stream 是 Hub 按 key 管理的一路流需要具备的能力:能分发播放请求(http.Handler), 且能在流结束时收尾(Finish)。pkg/hls.Stream(自研 origin)与 pkg/media/hlsmux.Bridge (gohlslib 后端)都满足它——Hub 因此与具体 HLS 实现解耦,不再绑死某一个。

type Supervisor

type Supervisor struct {
	// contains filtered or unexported fields
}

Supervisor 监督一个子进程(典型:每路流的 ffmpeg 转码进程):非正常退出时按退避 重启,ctx 取消时优雅停止(SIGTERM → 宽限 → Kill)。

命令由 newCmd 工厂提供——每次(重)启动都新建一个 *exec.Cmd,ffmpeg 参数、输入输出 管道等 policy 全在工厂里,Supervisor 只管"跑起来 + 崩了重启 + 让停就停"。

典型用法:把它绑到 Hub 的 Session.Context(),该路流结束时自动停机。

sup := media.NewSupervisor(func() *exec.Cmd { return exec.Command("ffmpeg", args...) },
    media.WithSupervisorMetrics(hub.Metrics(), key))
go sup.Run(sess.Context())

func NewSupervisor

func NewSupervisor(newCmd func() *exec.Cmd, opts ...SupervisorOption) *Supervisor

NewSupervisor 创建进程监督器。

func (*Supervisor) Run

func (s *Supervisor) Run(ctx context.Context) error

Run 启动并监督进程,阻塞直到 ctx 取消。进程自行退出(崩溃/结束)则按退避重启; ctx 取消时优雅停止当前进程并返回 nil。

type SupervisorOption

type SupervisorOption func(*Supervisor)

SupervisorOption 配置 Supervisor。

func WithRestartPolicy

func WithRestartPolicy(p *backoff.Policy) SupervisorOption

WithRestartPolicy 设置重启退避策略(默认 base 500ms、max 30s、指数)。

func WithStopGrace

func WithStopGrace(d time.Duration) SupervisorOption

WithStopGrace 设置优雅停止的宽限时长(SIGTERM 后等这么久仍未退出则 Kill,默认 5s)。

func WithSupervisorMetrics

func WithSupervisorMetrics(m *Metrics, name string) SupervisorOption

WithSupervisorMetrics 让重启计入 OTel 指标(name 作为 stream 标签)。

Directories

Path Synopsis
Package hlsmux 把 RTMP 采集到的 FLV(H.264 + AAC)喂给 bluenviron/gohlslib,产出 生产级 HLS——支持 MPEG-TS、fMP4 和 **LL-HLS(低延迟)** 三种 variant。
Package hlsmux 把 RTMP 采集到的 FLV(H.264 + AAC)喂给 bluenviron/gohlslib,产出 生产级 HLS——支持 MPEG-TS、fMP4 和 **LL-HLS(低延迟)** 三种 variant。
Package rtmp 提供一个 RTMP 采集(ingest)服务端,薄封装 github.com/yutopp/go-rtmp。
Package rtmp 提供一个 RTMP 采集(ingest)服务端,薄封装 github.com/yutopp/go-rtmp。
Package webrtc 提供 WebRTC 的 WHIP(采集)/WHEP(分发)薄机制,基于纯 Go 的 pion/webrtc(零 cgo)。
Package webrtc 提供 WebRTC 的 WHIP(采集)/WHEP(分发)薄机制,基于纯 Go 的 pion/webrtc(零 cgo)。
sfu
Package sfu 在 pkg/media/webrtc 之上提供一个「会议室」SFU 原语:多人实时音视频, 每个参会者推自己的轨道、订阅其他所有人的轨道,服务端做选择性转发(Selective Forwarding Unit,不混流不转码)。
Package sfu 在 pkg/media/webrtc 之上提供一个「会议室」SFU 原语:多人实时音视频, 每个参会者推自己的轨道、订阅其他所有人的轨道,服务端做选择性转发(Selective Forwarding Unit,不混流不转码)。

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL