Something went wrong. Try again.
This repository has no description
Something went wrong. Try again.
3.4 kB · 113 lines
Go
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114package director
import ( "context" "sync"
"github.com/streamplace/oatproxy/pkg/oatproxy" "golang.org/x/sync/errgroup" "stream.place/streamplace/pkg/atproto" "stream.place/streamplace/pkg/bus" "stream.place/streamplace/pkg/config" "stream.place/streamplace/pkg/localdb" "stream.place/streamplace/pkg/log" "stream.place/streamplace/pkg/media" "stream.place/streamplace/pkg/model" "stream.place/streamplace/pkg/replication" "stream.place/streamplace/pkg/statedb")
// director is responsible for managing the lifecycle of a stream, making business// logic decisions about when to do things like// - size of the in-memory segment cache// - transcoding// - thumbnail generation
type Director struct { mm *media.MediaManager mod model.Model cli *config.CLI bus *bus.Bus streamSessions map[string]*StreamSession streamSessionsMu sync.Mutex op *oatproxy.OATProxy statefulDB *statedb.StatefulDB replicator replication.Replicator localDB localdb.LocalDB atsync *atproto.ATProtoSynchronizer}
func NewDirector(mm *media.MediaManager, mod model.Model, cli *config.CLI, bus *bus.Bus, op *oatproxy.OATProxy, statefulDB *statedb.StatefulDB, replicator replication.Replicator, ldb localdb.LocalDB, atsync *atproto.ATProtoSynchronizer) *Director { return &Director{ mm: mm, mod: mod, cli: cli, bus: bus, streamSessions: make(map[string]*StreamSession), streamSessionsMu: sync.Mutex{}, op: op, statefulDB: statefulDB, replicator: replicator, localDB: ldb, atsync: atsync, }}
func (d *Director) Start(ctx context.Context) error { newSeg := d.mm.NewSegment() ctx, cancel := context.WithCancel(ctx) defer cancel() g, ctx := errgroup.WithContext(ctx) for { select { case <-ctx.Done(): cancel() return g.Wait() case not := <-newSeg: d.streamSessionsMu.Lock() ss, ok := d.streamSessions[not.Segment.RepoDID] if !ok { ss = &StreamSession{ lp: nil, repoDID: not.Segment.RepoDID, mm: d.mm, mod: d.mod, cli: d.cli, bus: d.bus, segmentChan: make(chan struct{}), sourceLane: newLane(), renditions: newLane(), op: d.op, packets: make([]bus.PacketizedSegment, 0), started: make(chan struct{}), statefulDB: d.statefulDB, replicator: d.replicator, // Initialize notification channels (buffered size 1 for coalescing) statusUpdateChan: make(chan struct{}, 1), originUpdateChan: make(chan struct{}, 1), livestreamUpdateChan: make(chan struct{}, 1), viewCountUpdateChan: make(chan struct{}, 1), localDB: d.localDB, atsync: d.atsync, } d.streamSessions[not.Segment.RepoDID] = ss g.Go(func() error { err := ss.Start(ctx, not) if err != nil { log.Error(ctx, "could not start stream session", "error", err) } d.streamSessionsMu.Lock() delete(d.streamSessions, not.Segment.RepoDID) d.streamSessionsMu.Unlock() return nil }) } d.streamSessionsMu.Unlock()
err := ss.NewSegment(ctx, not) if err != nil { log.Error(ctx, "could not add segment to stream session", "error", err) } } }}