video working
This commit is contained in:
commit
b7abd95296
51 files changed
+68060
No files matched your search
@@ -0,0 +1,35 @@
|
||||
package common
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"net"
|
||||
|
||||
"github.com/alexflint/go-arg"
|
||||
)
|
||||
|
||||
type Config struct {
|
||||
Logging ConfigLogging `arg:"-"`
|
||||
Http ConfigHTTP `arg:"-"`
|
||||
}
|
||||
|
||||
type ConfigLogging struct {
|
||||
Level string `arg:"--log-level,env:LOG_LEVEL" default:"debug"`
|
||||
}
|
||||
|
||||
type ConfigHTTP struct {
|
||||
Host string `arg:"--http-host,env:HTTP_HOST" default:"0.0.0.0"`
|
||||
Port uint `arg:"--http-port,env:HTTP_PORT" default:"8080"`
|
||||
PathGetLiveness string `arg:"env:HTTP_PATH_LIVENESS" default:"/healthz"`
|
||||
PathGetReadiness string `arg:"env:HTTP_PATH_READINESS" default:"/readyz"`
|
||||
StaticPath *string `arg:"--http-static,env:HTTP_STATIC"`
|
||||
}
|
||||
|
||||
func (c *ConfigHTTP) Address() string {
|
||||
return net.JoinHostPort(c.Host, fmt.Sprint(c.Port))
|
||||
}
|
||||
|
||||
func (c *Config) MustParse() {
|
||||
arg.MustParse(&c.Logging)
|
||||
arg.MustParse(&c.Http)
|
||||
arg.MustParse(c)
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
package common
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
)
|
||||
|
||||
func NewLogger(cfg *ConfigLogging) (*zap.Logger, error) {
|
||||
var c zap.Config
|
||||
if cfg.Level == "debug" {
|
||||
c = zap.NewDevelopmentConfig()
|
||||
} else {
|
||||
c = zap.NewProductionConfig()
|
||||
}
|
||||
|
||||
c.EncoderConfig.TimeKey = "time"
|
||||
c.EncoderConfig.EncodeTime = zapcore.TimeEncoderOfLayout(time.RFC3339)
|
||||
|
||||
lg, err := c.Build()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return lg, nil
|
||||
}
|
||||
@@ -0,0 +1,193 @@
|
||||
package server
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"mime"
|
||||
"net/http"
|
||||
"os"
|
||||
"path"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gitea.heinrich.blue/PHI/webrtc-gst/pkg/common"
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/pion/webrtc/v3"
|
||||
"go.uber.org/zap"
|
||||
"nhooyr.io/websocket"
|
||||
"nhooyr.io/websocket/wsjson"
|
||||
)
|
||||
|
||||
type SignalingHandle struct {
|
||||
Id string
|
||||
Recv chan webrtc.SessionDescription
|
||||
Trcv chan webrtc.SessionDescription
|
||||
}
|
||||
|
||||
type HttpServer struct {
|
||||
http.Server
|
||||
lg *zap.Logger
|
||||
Hndl chan *SignalingHandle
|
||||
}
|
||||
|
||||
func NewSignalingHandle(id string) SignalingHandle {
|
||||
return SignalingHandle{
|
||||
Id: id,
|
||||
Recv: make(chan webrtc.SessionDescription, 10),
|
||||
Trcv: make(chan webrtc.SessionDescription, 10),
|
||||
}
|
||||
}
|
||||
|
||||
func NewHttpServer(lg *zap.Logger, cfg *common.Config) *HttpServer {
|
||||
h := HttpServer{
|
||||
Hndl: make(chan *SignalingHandle, 10),
|
||||
lg: lg,
|
||||
}
|
||||
|
||||
engine := gin.Default()
|
||||
engine.GET("/signaling/:id", h.signalingHandler)
|
||||
|
||||
if cfg.Http.StaticPath != nil {
|
||||
engine.NoRoute(func(c *gin.Context) {
|
||||
t := path.Join(*cfg.Http.StaticPath, c.Request.URL.Path)
|
||||
|
||||
if i, err := os.Stat(t); err == nil && !i.IsDir() {
|
||||
h.compress(c, t)
|
||||
return
|
||||
}
|
||||
|
||||
h.compress(c, path.Join(*cfg.Http.StaticPath, "index.html"))
|
||||
})
|
||||
}
|
||||
|
||||
// set out handler
|
||||
h.Handler = engine
|
||||
|
||||
return &h
|
||||
}
|
||||
|
||||
func (h *HttpServer) compress(c *gin.Context, t string) {
|
||||
encoding := c.Request.Header.Get("Accept-Encoding")
|
||||
|
||||
switch true {
|
||||
case strings.Contains(encoding, "br"):
|
||||
lt := t + ".br"
|
||||
if i, err := os.Stat(lt); err == nil {
|
||||
fmt.Println(i.Name())
|
||||
c.Header("Content-Type", mime.TypeByExtension(filepath.Ext(t)))
|
||||
c.Header("Content-Encoding", "br")
|
||||
c.Header("Vary", "Accept-Encoding")
|
||||
c.File(lt)
|
||||
return
|
||||
}
|
||||
case strings.Contains(encoding, "gzip"):
|
||||
lt := t + ".br"
|
||||
if i, err := os.Stat(lt); err == nil {
|
||||
fmt.Println(i.Name())
|
||||
c.Header("Content-Type", mime.TypeByExtension(filepath.Ext(t)))
|
||||
c.Header("Content-Encoding", "gzip")
|
||||
c.File(lt)
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
fmt.Println(t)
|
||||
c.File(t)
|
||||
}
|
||||
|
||||
func (h *HttpServer) ListenAndServe(ctx context.Context, addr string) {
|
||||
go func() {
|
||||
|
||||
// set the configured address
|
||||
h.Addr = addr
|
||||
|
||||
// and listen
|
||||
if err := h.Server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||
h.lg.Fatal("listen failed", zap.Error(err))
|
||||
}
|
||||
|
||||
}()
|
||||
}
|
||||
|
||||
func (h *HttpServer) TearDown() error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
if err := h.Server.Shutdown(ctx); err != nil {
|
||||
return fmt.Errorf("server forced to shutdown: %s", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *HttpServer) signalingHandler(c *gin.Context) {
|
||||
id := c.Param("id")
|
||||
|
||||
conn, err := websocket.Accept(c.Writer, c.Request, nil)
|
||||
if err != nil {
|
||||
h.lg.Error("failed to upgrade websocket", zap.Error(err))
|
||||
c.Status(http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
// tear down when going down
|
||||
defer conn.Close(websocket.StatusInternalError, "the sky is falling")
|
||||
|
||||
ctx, cancel := context.WithCancel(c.Request.Context())
|
||||
defer cancel()
|
||||
|
||||
hndl := NewSignalingHandle(id)
|
||||
|
||||
// push hndl outside
|
||||
h.Hndl <- &hndl
|
||||
|
||||
// loop read
|
||||
go func() {
|
||||
for {
|
||||
// wait and read/parse message
|
||||
var sdp webrtc.SessionDescription
|
||||
err := wsjson.Read(ctx, conn, &sdp)
|
||||
if err != nil {
|
||||
status := websocket.CloseStatus(err)
|
||||
if status == websocket.StatusNormalClosure || status == websocket.StatusGoingAway {
|
||||
|
||||
// close channels to signal closure
|
||||
close(hndl.Recv)
|
||||
close(hndl.Trcv)
|
||||
|
||||
h.lg.Info("socket closed")
|
||||
break
|
||||
}
|
||||
|
||||
h.lg.Error("failed to decode message", zap.Error(err))
|
||||
break
|
||||
}
|
||||
|
||||
h.lg.Info("received message", zap.String("type", sdp.Type.String()))
|
||||
|
||||
// forward in channel
|
||||
hndl.Recv <- sdp
|
||||
}
|
||||
|
||||
cancel()
|
||||
}()
|
||||
|
||||
// loop write
|
||||
for sdp := range hndl.Trcv {
|
||||
err := wsjson.Write(ctx, conn, sdp)
|
||||
if err != nil {
|
||||
status := websocket.CloseStatus(err)
|
||||
if status == websocket.StatusNormalClosure || status == websocket.StatusGoingAway {
|
||||
h.lg.Info("socket closed")
|
||||
break
|
||||
}
|
||||
|
||||
h.lg.Error("failed to encode message", zap.Error(err))
|
||||
break
|
||||
}
|
||||
|
||||
h.lg.Info("tranceived message", zap.String("type", sdp.Type.String()))
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,119 @@
|
||||
package webrtc
|
||||
|
||||
import (
|
||||
"github.com/pion/webrtc/v3/pkg/media"
|
||||
"github.com/tinyzimmer/go-gst/gst"
|
||||
"github.com/tinyzimmer/go-gst/gst/app"
|
||||
)
|
||||
|
||||
func init() {
|
||||
gst.Init(nil)
|
||||
}
|
||||
|
||||
func setCallback(sink *app.Sink, ch chan media.Sample) {
|
||||
sink.SetCallbacks(&app.SinkCallbacks{
|
||||
NewSampleFunc: func(sink *app.Sink) gst.FlowReturn {
|
||||
// Pull the sample that triggered this callback
|
||||
sample := sink.PullSample()
|
||||
if sample == nil {
|
||||
return gst.FlowEOS
|
||||
}
|
||||
|
||||
// Retrieve the buffer from the sample
|
||||
buffer := sample.GetBuffer()
|
||||
if buffer == nil {
|
||||
return gst.FlowError
|
||||
}
|
||||
|
||||
// At this point, buffer is only a reference to an existing memory region somewhere.
|
||||
// When we want to access its content, we have to map it while requesting the required
|
||||
// mode of access (read, read/write).
|
||||
data := buffer.Map(gst.MapRead).AsUint8Slice()
|
||||
defer buffer.Unmap()
|
||||
|
||||
ch <- media.Sample{Data: data, Duration: buffer.Duration()}
|
||||
|
||||
return gst.FlowOK
|
||||
},
|
||||
})
|
||||
}
|
||||
|
||||
func CreateVideoPipeline(s string) (*gst.Pipeline, <-chan media.Sample, error) {
|
||||
// Create a pipeline
|
||||
pipeline, err := gst.NewPipeline("")
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// Create the src
|
||||
src, err := gst.NewElement(s)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// Create the enc
|
||||
enc, err := gst.NewElement("vp8enc")
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
enc.SetProperty("error-resilient", "partitions")
|
||||
enc.SetProperty("keyframe-max-dist", 10)
|
||||
enc.SetProperty("auto-alt-ref", true)
|
||||
enc.SetProperty("cpu-used", 5)
|
||||
enc.SetProperty("deadline", 1)
|
||||
|
||||
// Create the sink
|
||||
sink, err := app.NewAppSink()
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
ch := make(chan media.Sample, 100)
|
||||
setCallback(sink, ch)
|
||||
|
||||
// Add the elements to the pipeline
|
||||
pipeline.AddMany(src, enc, sink.Element)
|
||||
|
||||
// link the elements
|
||||
gst.ElementLinkMany(src, enc, sink.Element)
|
||||
|
||||
return pipeline, ch, nil
|
||||
}
|
||||
|
||||
func CreateAudioPipeline(s string) (*gst.Pipeline, <-chan media.Sample, error) {
|
||||
// Create a pipeline
|
||||
pipeline, err := gst.NewPipeline("")
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// Create the src
|
||||
src, err := gst.NewElement(s)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// Create the enc
|
||||
enc, err := gst.NewElement("opusenc")
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
// Create the sink
|
||||
sink, err := app.NewAppSink()
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
|
||||
ch := make(chan media.Sample, 100)
|
||||
setCallback(sink, ch)
|
||||
|
||||
// Add the elements to the pipeline
|
||||
pipeline.AddMany(src, enc, sink.Element)
|
||||
|
||||
// link the elements
|
||||
gst.ElementLinkMany(src, enc, sink.Element)
|
||||
|
||||
return pipeline, ch, nil
|
||||
}
|
||||
@@ -0,0 +1,278 @@
|
||||
package webrtc
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
|
||||
"gitea.heinrich.blue/PHI/webrtc-gst/pkg/server"
|
||||
"github.com/pion/webrtc/v3"
|
||||
"github.com/pion/webrtc/v3/pkg/media"
|
||||
"github.com/tinyzimmer/go-gst/gst"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
const STUN_SERVER = "stun:stun.l.google.com:19302"
|
||||
|
||||
type WebrtcHandler struct {
|
||||
lg *zap.Logger
|
||||
mu *sync.Mutex
|
||||
audioPipeline *gst.Pipeline
|
||||
videoPipeline *gst.Pipeline
|
||||
peerHandles map[string]*PeerHandle
|
||||
}
|
||||
|
||||
type PeerHandle struct {
|
||||
audioTrack *webrtc.TrackLocalStaticSample
|
||||
videoTrack *webrtc.TrackLocalStaticSample
|
||||
}
|
||||
|
||||
func NewWebrtcHandler(ctx context.Context, lg *zap.Logger, ch <-chan *server.SignalingHandle) error {
|
||||
wh := WebrtcHandler{
|
||||
lg: lg,
|
||||
mu: &sync.Mutex{},
|
||||
peerHandles: make(map[string]*PeerHandle, 0),
|
||||
}
|
||||
|
||||
err := wh.handleAudioSamples(ctx, "audiotestsrc")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
err = wh.handleVideoSamples(ctx, "videotestsrc")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
go func() {
|
||||
for {
|
||||
select {
|
||||
case sh := <-ch:
|
||||
wh.lg.Info("new connection", zap.String("id", sh.Id))
|
||||
|
||||
// create new peer handle
|
||||
wh.createPeerHandle(ctx, sh)
|
||||
|
||||
// make sure the pipelines are running
|
||||
wh.startPipelines()
|
||||
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (wh *WebrtcHandler) startPipelines() {
|
||||
|
||||
if wh.audioPipeline.GetState() != gst.StatePlaying {
|
||||
err := wh.audioPipeline.SetState(gst.StatePlaying)
|
||||
if err != nil {
|
||||
wh.lg.Fatal("failed to start audio pipeline", zap.Error(err))
|
||||
}
|
||||
wh.lg.Info("started audio pipeline")
|
||||
}
|
||||
|
||||
if wh.videoPipeline.GetState() != gst.StatePlaying {
|
||||
err := wh.videoPipeline.SetState(gst.StatePlaying)
|
||||
if err != nil {
|
||||
wh.lg.Fatal("failed to start video pipeline", zap.Error(err))
|
||||
}
|
||||
wh.lg.Info("started video pipeline")
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func (wh *WebrtcHandler) stopPipelines() {
|
||||
var err error
|
||||
|
||||
err = wh.audioPipeline.SetState(gst.StatePaused)
|
||||
if err != nil {
|
||||
wh.lg.Fatal("failed to pause audio pipeline", zap.Error(err))
|
||||
}
|
||||
|
||||
err = wh.videoPipeline.SetState(gst.StatePaused)
|
||||
if err != nil {
|
||||
wh.lg.Fatal("failed to pause video pipeline", zap.Error(err))
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
func (wh *WebrtcHandler) handleAudioSamples(ctx context.Context, src string) error {
|
||||
var err error
|
||||
var audioCh <-chan media.Sample
|
||||
wh.audioPipeline, audioCh, err = CreateAudioPipeline(src)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
go func() {
|
||||
wh.lg.Info("wait for audio sample")
|
||||
for {
|
||||
select {
|
||||
case data := <-audioCh:
|
||||
for id, ph := range wh.peerHandles {
|
||||
err := ph.audioTrack.WriteSample(data)
|
||||
if err != nil {
|
||||
wh.lg.Error("failed to write audio sample", zap.String("id", id), zap.Error(err))
|
||||
}
|
||||
}
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (wh *WebrtcHandler) handleVideoSamples(ctx context.Context, src string) error {
|
||||
var err error
|
||||
var videoCh <-chan media.Sample
|
||||
wh.videoPipeline, videoCh, err = CreateVideoPipeline(src)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
go func() {
|
||||
wh.lg.Info("wait for video sample")
|
||||
for {
|
||||
select {
|
||||
case data := <-videoCh:
|
||||
for id, ph := range wh.peerHandles {
|
||||
err := ph.videoTrack.WriteSample(data)
|
||||
if err != nil {
|
||||
wh.lg.Error("failed to write video sample", zap.String("id", id), zap.Error(err))
|
||||
}
|
||||
}
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (wh *WebrtcHandler) createPeerHandle(ctx context.Context, sh *server.SignalingHandle) error {
|
||||
wh.mu.Lock()
|
||||
defer wh.mu.Unlock()
|
||||
|
||||
// Prepare the configuration
|
||||
config := webrtc.Configuration{
|
||||
ICEServers: []webrtc.ICEServer{{URLs: []string{STUN_SERVER}}},
|
||||
}
|
||||
|
||||
// Create a new RTCPeerConnection
|
||||
peerConnection, err := webrtc.NewPeerConnection(config)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Set the handler for ICE connection state
|
||||
// This will notify you when the peer has connected/disconnected
|
||||
peerConnection.OnICEConnectionStateChange(func(connectionState webrtc.ICEConnectionState) {
|
||||
wh.lg.Info("connection-state has changed", zap.String("state", connectionState.String()))
|
||||
|
||||
if connectionState == webrtc.ICEConnectionStateDisconnected {
|
||||
peerConnection.Close()
|
||||
|
||||
// remove this handle
|
||||
wh.mu.Lock()
|
||||
delete(wh.peerHandles, sh.Id)
|
||||
wh.mu.Unlock()
|
||||
}
|
||||
})
|
||||
|
||||
// Create a audio track
|
||||
audioTrack, err := webrtc.NewTrackLocalStaticSample(webrtc.RTPCodecCapability{MimeType: "audio/opus"}, "audio", "pion1")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = peerConnection.AddTrack(audioTrack)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Create a video track
|
||||
videoTrack, err := webrtc.NewTrackLocalStaticSample(webrtc.RTPCodecCapability{MimeType: "video/vp8"}, "video", "pion2")
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
_, err = peerConnection.AddTrack(videoTrack)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
onOfferReceived := func(offer webrtc.SessionDescription) error {
|
||||
|
||||
// Set the remote SessionDescription
|
||||
err = peerConnection.SetRemoteDescription(offer)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Create an answer
|
||||
answer, err := peerConnection.CreateAnswer(nil)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Create channel that is blocked until ICE Gathering is complete
|
||||
gatherComplete := webrtc.GatheringCompletePromise(peerConnection)
|
||||
|
||||
// Sets the LocalDescription, and starts our UDP listeners
|
||||
err = peerConnection.SetLocalDescription(answer)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
<-gatherComplete
|
||||
|
||||
// Send the answer
|
||||
sh.Trcv <- *peerConnection.LocalDescription()
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
hndl := PeerHandle{
|
||||
audioTrack: audioTrack,
|
||||
videoTrack: videoTrack,
|
||||
}
|
||||
|
||||
go func() {
|
||||
for {
|
||||
select {
|
||||
case offer, ok := <-sh.Recv:
|
||||
// return when channel is closed
|
||||
if !ok {
|
||||
wh.mu.Lock()
|
||||
defer wh.mu.Unlock()
|
||||
|
||||
// remove this handle
|
||||
delete(wh.peerHandles, sh.Id)
|
||||
|
||||
// when there is no left over pause the pipelines
|
||||
if len(wh.peerHandles) == 0 {
|
||||
wh.stopPipelines()
|
||||
}
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
err := onOfferReceived(offer)
|
||||
if err != nil {
|
||||
wh.lg.Error("failed to handle offer", zap.Error(err))
|
||||
}
|
||||
case <-ctx.Done():
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
// add handle to list
|
||||
wh.peerHandles[sh.Id] = &hndl
|
||||
|
||||
return nil
|
||||
}
|
||||
Reference in new issue
Block a user