Files
cam2ip/handlers/stream.go

135 lines
2.4 KiB
Go

package handlers
import (
"bytes"
"log"
"sync"
"time"
"github.com/gen2brain/cam2ip/image"
)
const errorBackoff = 100 * time.Millisecond
// Stream captures frames in a single loop and broadcasts the encoded JPEG to all subscribers.
type Stream struct {
open func() (ImageReader, error)
lazy bool
delay int
quality int
reader ImageReader
mu sync.Mutex
cond *sync.Cond
subs map[chan []byte]struct{}
}
// NewStream returns a new Stream and starts its capture loop. open acquires the camera: once up front, or while clients are connected when lazy.
func NewStream(open func() (ImageReader, error), delay, quality int, lazy bool) *Stream {
s := &Stream{
open: open,
lazy: lazy,
delay: delay,
quality: quality,
subs: make(map[chan []byte]struct{}),
}
s.cond = sync.NewCond(&s.mu)
if !lazy {
s.reader, _ = open()
}
go s.capture()
return s
}
// capture reads, encodes and broadcasts frames while subscribers are connected.
func (s *Stream) capture() {
for {
s.mu.Lock()
for len(s.subs) == 0 {
if s.lazy && s.reader != nil {
s.reader.Close()
s.reader = nil
}
s.cond.Wait()
}
s.mu.Unlock()
if s.reader == nil {
reader, err := s.open()
if err != nil {
log.Printf("stream: open: %v", err)
time.Sleep(errorBackoff)
continue
}
s.reader = reader
}
img, err := s.reader.Read()
if err != nil {
log.Printf("stream: read: %v", err)
time.Sleep(errorBackoff)
continue
}
buf := new(bytes.Buffer)
if err := image.NewEncoder(buf, s.quality).Encode(img); err != nil {
log.Printf("stream: encode: %v", err)
time.Sleep(errorBackoff)
continue
}
s.broadcast(buf.Bytes())
if s.delay > 0 {
time.Sleep(time.Duration(s.delay) * time.Millisecond)
}
}
}
// broadcast sends the frame to every subscriber, dropping stale frames so slow clients never block.
func (s *Stream) broadcast(frame []byte) {
s.mu.Lock()
defer s.mu.Unlock()
for ch := range s.subs {
select {
case <-ch:
default:
}
select {
case ch <- frame:
default:
}
}
}
// subscribe registers a new subscriber and returns its frame channel.
func (s *Stream) subscribe() chan []byte {
ch := make(chan []byte, 1)
s.mu.Lock()
s.subs[ch] = struct{}{}
s.mu.Unlock()
s.cond.Signal()
return ch
}
// unsubscribe removes a subscriber.
func (s *Stream) unsubscribe(ch chan []byte) {
s.mu.Lock()
delete(s.subs, ch)
s.mu.Unlock()
}