This commit is contained in:
dexter
2022-02-24 23:01:18 +08:00
parent f7e4c060ab
commit 2c03cc7263
3 changed files with 60 additions and 102 deletions

60
main.go
View File

@@ -5,7 +5,7 @@ import (
"encoding/binary" "encoding/binary"
"net" "net"
"net/http" "net/http"
"regexp" "strings"
"time" "time"
. "github.com/Monibuca/engine/v4" . "github.com/Monibuca/engine/v4"
@@ -24,10 +24,8 @@ type HDLConfig struct {
config.Pull config.Pull
} }
var streamPathReg = regexp.MustCompile(`/(hdl/)?((.+)(\.flv)|(.+))`)
func (c *HDLConfig) OnEvent(event any) { func (c *HDLConfig) OnEvent(event any) {
switch event.(type) { switch v := event.(type) {
case FirstConfig: case FirstConfig:
if c.ListenAddr != "" || c.ListenAddrTLS != "" { if c.ListenAddr != "" || c.ListenAddrTLS != "" {
plugin.Info(Green("HDL Server Start").String(), zap.String("ListenAddr", c.ListenAddr), zap.String("ListenAddrTLS", c.ListenAddrTLS)) plugin.Info(Green("HDL Server Start").String(), zap.String("ListenAddr", c.ListenAddr), zap.String("ListenAddrTLS", c.ListenAddrTLS))
@@ -35,33 +33,37 @@ func (c *HDLConfig) OnEvent(event any) {
} else { } else {
plugin.Info(Green("HDL start reuse engine port").String()) plugin.Info(Green("HDL start reuse engine port").String())
} }
case PullerPromise:
puller := v.Value
client := &HDLPuller{
Puller: puller,
}
err := client.connect()
if err == nil {
if err = plugin.Publish(puller.StreamPath, client); err == nil {
v.Resolve(util.Null)
break
}
}
client.Error(puller.RemoteURL, zap.Error(err))
v.Reject(err)
} }
} }
func (c *HDLConfig) API_Pull(rw http.ResponseWriter, r *http.Request) { func (c *HDLConfig) API_Pull(rw http.ResponseWriter, r *http.Request) {
var puller Puller err := plugin.Pull(r.URL.Query().Get("streamPath"), r.URL.Query().Get("target"), r.URL.Query().Has("save"))
puller.StreamPath = r.URL.Query().Get("streamPath") if err != nil {
puller.RemoteURL = r.URL.Query().Get("target") http.Error(rw, err.Error(), http.StatusBadRequest)
puller.Config = &c.Pull
c.PullStream(puller)
if r.URL.Query().Get("save") != "" {
c.AddPull(puller.StreamPath, puller.RemoteURL)
plugin.Modified["pull"] = c.Pull
if err := plugin.Save(); err != nil {
plugin.Error("save faild", zap.Error(err))
}
} }
} }
func (*HDLConfig) API_List(rw http.ResponseWriter, r *http.Request) { func (*HDLConfig) API_List(rw http.ResponseWriter, r *http.Request) {
util.ReturnJson(FilterStreams[*HDLPuller], time.Second, rw, r) util.ReturnJson(FilterStreams[*HDLPuller], time.Second, rw, r)
} }
var Config = new(HDLConfig)
// 确保HDLConfig实现了PullPlugin接口 // 确保HDLConfig实现了PullPlugin接口
var plugin = InstallPlugin(Config) var plugin = InstallPlugin(new(HDLConfig))
type HDLSubscriber struct { type HDLSubscriber struct {
Subscriber Subscriber
@@ -78,26 +80,14 @@ func (sub *HDLSubscriber) OnEvent(event any) {
} }
func (*HDLConfig) ServeHTTP(w http.ResponseWriter, r *http.Request) { func (*HDLConfig) ServeHTTP(w http.ResponseWriter, r *http.Request) {
parts := streamPathReg.FindStringSubmatch(r.RequestURI) streamPath := strings.TrimPrefix(r.URL.Path, "/hls")
if len(parts) == 0 {
w.WriteHeader(404)
return
}
stringPath := parts[3]
if stringPath == "" {
stringPath = parts[5]
}
w.Header().Set("Transfer-Encoding", "chunked") w.Header().Set("Transfer-Encoding", "chunked")
w.Header().Set("Content-Type", "video/x-flv") w.Header().Set("Content-Type", "video/x-flv")
sub := &HDLSubscriber{} sub := &HDLSubscriber{}
sub.ID = r.RemoteAddr sub.ID = r.RemoteAddr
sub.OnEvent(r.Context()) sub.OnEvent(r.Context()) //注入父级Context
if err := plugin.Subscribe(stringPath, sub); err == nil { sub.OnEvent(w) //注入Writer
if sub.Stream.Publisher == nil { if err := plugin.Subscribe(streamPath, sub); err == nil {
w.WriteHeader(404)
return
}
sub.Writer = w
at, vt := sub.AudioTrack, sub.VideoTrack at, vt := sub.AudioTrack, sub.VideoTrack
hasVideo := at != nil hasVideo := at != nil
hasAudio := vt != nil hasAudio := vt != nil

View File

@@ -1,23 +0,0 @@
package hdl
import (
"testing"
)
func TestHDLHandler(t *testing.T) {
tests := map[string]string{
"/hdl/abc.flv": "abc", "/hdl/abc": "abc", "/abc": "abc", "/abc.flv": "abc",
}
for name, result := range tests {
t.Run(name, func(t *testing.T) {
parts := streamPathReg.FindStringSubmatch(name)
stringPath := parts[3]
if stringPath == "" {
stringPath = parts[5]
}
if stringPath != result {
t.Fail()
}
})
}
}

79
pull.go
View File

@@ -9,15 +9,45 @@ import (
. "github.com/Monibuca/engine/v4" . "github.com/Monibuca/engine/v4"
"github.com/Monibuca/engine/v4/codec" "github.com/Monibuca/engine/v4/codec"
"github.com/Monibuca/engine/v4/log"
"github.com/Monibuca/engine/v4/util" "github.com/Monibuca/engine/v4/util"
"go.uber.org/zap" "go.uber.org/zap"
) )
func (puller *HDLPuller) pull() { func (puller *HDLPuller) connect() (err error) {
puller.ReConnectCount++ puller.ReConnectCount++
log.Info("connect", zap.String("remoteURL", puller.RemoteURL))
if strings.HasPrefix(puller.RemoteURL, "http") {
var res *http.Response
if res, err = http.Get(puller.RemoteURL); err == nil {
puller.OnEvent(res.Body)
}
} else {
var res *os.File
if res, err = os.Open(puller.RemoteURL); err == nil {
puller.OnEvent(res)
}
}
if err != nil {
log.Error("connect", zap.Error(err))
}
return
}
func (puller *HDLPuller) pull() {
var err error
defer func() {
puller.Closer.Close()
if !puller.Stream.IsClosed() {
if err = puller.connect(); err == nil {
go puller.pull()
}
} else {
puller.Info("stop", zap.String("remoteURL", puller.RemoteURL))
}
}()
head := util.Buffer(make([]byte, len(codec.FLVHeader))) head := util.Buffer(make([]byte, len(codec.FLVHeader)))
reader := bufio.NewReader(puller) reader := bufio.NewReader(puller)
_, err := io.ReadFull(reader, head) _, err = io.ReadFull(reader, head)
if err != nil { if err != nil {
return return
} }
@@ -58,49 +88,10 @@ type HDLPuller struct {
} }
func (puller *HDLPuller) OnEvent(event any) { func (puller *HDLPuller) OnEvent(event any) {
switch v := event.(type) { switch event.(type) {
case PullEvent: case SEpublish:
if v > 0 { go puller.pull() //阻塞拉流
go func(count PullEvent) {
puller.pull() //阻塞拉流
// 如果流没有被关闭,则重连,重拉
if !puller.Stream.IsClosed() {
puller.OnEvent(count)
}
}(v + 1)
} else {
// TODO: 发布失败重新发布
if plugin.Publish(puller.StreamPath, puller) == nil {
if strings.HasPrefix(puller.RemoteURL, "http") {
if res, err := http.Get(puller.RemoteURL); err == nil {
puller.Reader = res.Body
puller.Closer = res.Body
} else {
puller.Error(puller.RemoteURL, zap.Error(err))
return
}
} else {
if res, err := os.Open(puller.RemoteURL); err == nil {
puller.Reader = res
puller.Closer = res
} else {
puller.Error(puller.RemoteURL, zap.Error(err))
return
}
}
// 注入context
puller.OnEvent(Engine)
puller.OnEvent(PullEvent(1))
}
}
default: default:
puller.Publisher.OnEvent(event) puller.Publisher.OnEvent(event)
} }
} }
func (config *HDLConfig) PullStream(puller Puller) {
client := &HDLPuller{
Puller: puller,
}
client.OnEvent(PullEvent(0))
}