mirror of
https://github.com/m1k1o/neko.git
synced 2024-07-24 14:40:50 +12:00
GST pipelines refactor.
This commit is contained in:
parent
c10b2212d1
commit
16d762b6ae
@ -37,7 +37,7 @@ func (manager *CaptureManagerCtx) StopBroadcastPipeline() {
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
manager.broadcast.DestroyPipeline()
|
manager.broadcast.Stop()
|
||||||
manager.logger.Info().Msgf("Stopping broadcast pipeline...")
|
manager.logger.Info().Msgf("Stopping broadcast pipeline...")
|
||||||
manager.broadcast = nil
|
manager.broadcast = nil
|
||||||
}
|
}
|
||||||
|
@ -68,11 +68,6 @@ GstElement *gstreamer_send_create_pipeline(char *pipeline) {
|
|||||||
return gst_parse_launch(pipeline, &error);
|
return gst_parse_launch(pipeline, &error);
|
||||||
}
|
}
|
||||||
|
|
||||||
void gstreamer_send_destroy_pipeline(GstElement *pipeline) {
|
|
||||||
gst_element_set_state(pipeline, GST_STATE_NULL);
|
|
||||||
gst_object_unref(pipeline);
|
|
||||||
}
|
|
||||||
|
|
||||||
void gstreamer_send_start_pipeline(GstElement *pipeline, int pipelineId) {
|
void gstreamer_send_start_pipeline(GstElement *pipeline, int pipelineId) {
|
||||||
SampleHandlerUserData *s = calloc(1, sizeof(SampleHandlerUserData));
|
SampleHandlerUserData *s = calloc(1, sizeof(SampleHandlerUserData));
|
||||||
s->pipelineId = pipelineId;
|
s->pipelineId = pipelineId;
|
||||||
@ -95,4 +90,5 @@ void gstreamer_send_play_pipeline(GstElement *pipeline) {
|
|||||||
|
|
||||||
void gstreamer_send_stop_pipeline(GstElement *pipeline) {
|
void gstreamer_send_stop_pipeline(GstElement *pipeline) {
|
||||||
gst_element_set_state(pipeline, GST_STATE_NULL);
|
gst_element_set_state(pipeline, GST_STATE_NULL);
|
||||||
|
gst_object_unref(pipeline);
|
||||||
}
|
}
|
||||||
|
@ -210,11 +210,6 @@ func CreatePipeline(pipelineStr string, codecName string, clockRate float32) (*P
|
|||||||
return p, nil
|
return p, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Destroy GStreamer Pipeline
|
|
||||||
func (p *Pipeline) DestroyPipeline() {
|
|
||||||
C.gstreamer_send_destroy_pipeline(p.Pipeline)
|
|
||||||
}
|
|
||||||
|
|
||||||
// Start starts the GStreamer Pipeline
|
// Start starts the GStreamer Pipeline
|
||||||
func (p *Pipeline) Start() {
|
func (p *Pipeline) Start() {
|
||||||
C.gstreamer_send_start_pipeline(p.Pipeline, C.int(p.id))
|
C.gstreamer_send_start_pipeline(p.Pipeline, C.int(p.id))
|
||||||
|
@ -9,7 +9,6 @@
|
|||||||
extern void goHandlePipelineBuffer(void *buffer, int bufferLen, int samples, int pipelineId);
|
extern void goHandlePipelineBuffer(void *buffer, int bufferLen, int samples, int pipelineId);
|
||||||
|
|
||||||
GstElement *gstreamer_send_create_pipeline(char *pipeline);
|
GstElement *gstreamer_send_create_pipeline(char *pipeline);
|
||||||
void gstreamer_send_destroy_pipeline(GstElement *pipeline);
|
|
||||||
|
|
||||||
void gstreamer_send_start_pipeline(GstElement *pipeline, int pipelineId);
|
void gstreamer_send_start_pipeline(GstElement *pipeline, int pipelineId);
|
||||||
void gstreamer_send_play_pipeline(GstElement *pipeline);
|
void gstreamer_send_play_pipeline(GstElement *pipeline);
|
||||||
|
@ -18,7 +18,8 @@ type CaptureManagerCtx struct {
|
|||||||
audio *gst.Pipeline
|
audio *gst.Pipeline
|
||||||
broadcast *gst.Pipeline
|
broadcast *gst.Pipeline
|
||||||
config *config.Capture
|
config *config.Capture
|
||||||
shutdown chan bool
|
audio_stop chan bool
|
||||||
|
video_stop chan bool
|
||||||
emmiter events.EventEmmiter
|
emmiter events.EventEmmiter
|
||||||
streaming bool
|
streaming bool
|
||||||
broadcasting bool
|
broadcasting bool
|
||||||
@ -29,7 +30,8 @@ type CaptureManagerCtx struct {
|
|||||||
func New(desktop types.DesktopManager, config *config.Capture) *CaptureManagerCtx {
|
func New(desktop types.DesktopManager, config *config.Capture) *CaptureManagerCtx {
|
||||||
return &CaptureManagerCtx{
|
return &CaptureManagerCtx{
|
||||||
logger: log.With().Str("module", "capture").Logger(),
|
logger: log.With().Str("module", "capture").Logger(),
|
||||||
shutdown: make(chan bool),
|
audio_stop: make(chan bool),
|
||||||
|
video_stop: make(chan bool),
|
||||||
emmiter: events.New(),
|
emmiter: events.New(),
|
||||||
config: config,
|
config: config,
|
||||||
streaming: false,
|
streaming: false,
|
||||||
@ -48,35 +50,15 @@ func (manager *CaptureManagerCtx) Start() {
|
|||||||
manager.logger.Warn().Err(err).Msg("unable to change screen size")
|
manager.logger.Warn().Err(err).Msg("unable to change screen size")
|
||||||
}
|
}
|
||||||
|
|
||||||
manager.CreateVideoPipeline()
|
|
||||||
manager.CreateAudioPipeline()
|
|
||||||
manager.StartBroadcastPipeline()
|
manager.StartBroadcastPipeline()
|
||||||
|
|
||||||
go func() {
|
|
||||||
defer func() {
|
|
||||||
manager.logger.Info().Msg("shutdown")
|
|
||||||
}()
|
|
||||||
|
|
||||||
for {
|
|
||||||
select {
|
|
||||||
case <-manager.shutdown:
|
|
||||||
return
|
|
||||||
case sample := <-manager.video.Sample:
|
|
||||||
manager.emmiter.Emit("video", sample)
|
|
||||||
case sample := <-manager.audio.Sample:
|
|
||||||
manager.emmiter.Emit("audio", sample)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (manager *CaptureManagerCtx) Shutdown() error {
|
func (manager *CaptureManagerCtx) Shutdown() error {
|
||||||
manager.logger.Info().Msgf("capture shutting down")
|
manager.logger.Info().Msgf("capture shutting down")
|
||||||
manager.video.DestroyPipeline()
|
manager.audio_stop <- true
|
||||||
manager.audio.DestroyPipeline()
|
manager.video_stop <- true
|
||||||
manager.StopBroadcastPipeline()
|
manager.StopBroadcastPipeline()
|
||||||
|
|
||||||
manager.shutdown <- true
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -103,16 +85,16 @@ func (manager *CaptureManagerCtx) OnAudioFrame(listener func(sample types.Sample
|
|||||||
func (manager *CaptureManagerCtx) StartStream() {
|
func (manager *CaptureManagerCtx) StartStream() {
|
||||||
manager.logger.Info().Msgf("Pipelines starting...")
|
manager.logger.Info().Msgf("Pipelines starting...")
|
||||||
|
|
||||||
manager.video.Start()
|
manager.createVideoPipeline()
|
||||||
manager.audio.Start()
|
manager.createAudioPipeline()
|
||||||
manager.streaming = true
|
manager.streaming = true
|
||||||
}
|
}
|
||||||
|
|
||||||
func (manager *CaptureManagerCtx) StopStream() {
|
func (manager *CaptureManagerCtx) StopStream() {
|
||||||
manager.logger.Info().Msgf("Pipelines shutting down...")
|
manager.logger.Info().Msgf("Pipelines stopping...")
|
||||||
|
|
||||||
manager.video.Stop()
|
manager.audio_stop <- true
|
||||||
manager.audio.Stop()
|
manager.video_stop <- true
|
||||||
manager.streaming = false
|
manager.streaming = false
|
||||||
}
|
}
|
||||||
|
|
||||||
@ -120,7 +102,19 @@ func (manager *CaptureManagerCtx) Streaming() bool {
|
|||||||
return manager.streaming
|
return manager.streaming
|
||||||
}
|
}
|
||||||
|
|
||||||
func (manager *CaptureManagerCtx) CreateVideoPipeline() {
|
func (manager *CaptureManagerCtx) ChangeResolution(width int, height int, rate int) error {
|
||||||
|
manager.video_stop <- true
|
||||||
|
manager.StopBroadcastPipeline()
|
||||||
|
|
||||||
|
defer func() {
|
||||||
|
manager.createVideoPipeline()
|
||||||
|
manager.StartBroadcastPipeline()
|
||||||
|
}()
|
||||||
|
|
||||||
|
return manager.desktop.ChangeScreenSize(width, height, rate)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (manager *CaptureManagerCtx) createVideoPipeline() {
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
manager.logger.Info().
|
manager.logger.Info().
|
||||||
@ -138,9 +132,34 @@ func (manager *CaptureManagerCtx) CreateVideoPipeline() {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
manager.logger.Panic().Err(err).Msg("unable to create video pipeline")
|
manager.logger.Panic().Err(err).Msg("unable to create video pipeline")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
manager.logger.Info().
|
||||||
|
Str("pipeline", manager.video.Src).
|
||||||
|
Msgf("Starting video pipeline...")
|
||||||
|
|
||||||
|
manager.video.Start()
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
manager.logger.Debug().Msg("started emitting video data")
|
||||||
|
|
||||||
|
defer func() {
|
||||||
|
manager.logger.Debug().Msg("stopped emitting video data")
|
||||||
|
}()
|
||||||
|
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-manager.video_stop:
|
||||||
|
manager.logger.Info().Msgf("Stopping video pipeline...")
|
||||||
|
manager.video.Stop()
|
||||||
|
return
|
||||||
|
case sample := <-manager.video.Sample:
|
||||||
|
manager.emmiter.Emit("video", sample)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
func (manager *CaptureManagerCtx) CreateAudioPipeline() {
|
func (manager *CaptureManagerCtx) createAudioPipeline() {
|
||||||
var err error
|
var err error
|
||||||
|
|
||||||
manager.logger.Info().
|
manager.logger.Info().
|
||||||
@ -158,20 +177,29 @@ func (manager *CaptureManagerCtx) CreateAudioPipeline() {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
manager.logger.Panic().Err(err).Msg("unable to create audio pipeline")
|
manager.logger.Panic().Err(err).Msg("unable to create audio pipeline")
|
||||||
}
|
}
|
||||||
}
|
|
||||||
|
|
||||||
func (manager *CaptureManagerCtx) ChangeResolution(width int, height int, rate int) error {
|
manager.logger.Info().
|
||||||
manager.video.DestroyPipeline()
|
Str("pipeline", manager.audio.Src).
|
||||||
manager.StopBroadcastPipeline()
|
Msgf("Starting audio pipeline...")
|
||||||
|
|
||||||
|
manager.audio.Start()
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
manager.logger.Debug().Msg("started emitting audio data")
|
||||||
|
|
||||||
defer func() {
|
defer func() {
|
||||||
manager.CreateVideoPipeline()
|
manager.logger.Debug().Msg("stopped emitting audio data")
|
||||||
|
|
||||||
manager.video.Start()
|
|
||||||
manager.logger.Info().Msg("starting video pipeline...")
|
|
||||||
|
|
||||||
manager.StartBroadcastPipeline()
|
|
||||||
}()
|
}()
|
||||||
|
|
||||||
return manager.desktop.ChangeScreenSize(width, height, rate)
|
for {
|
||||||
|
select {
|
||||||
|
case <-manager.audio_stop:
|
||||||
|
manager.logger.Info().Msgf("Stopping audio pipeline...")
|
||||||
|
manager.audio.Stop()
|
||||||
|
return
|
||||||
|
case sample := <-manager.audio.Sample:
|
||||||
|
manager.emmiter.Emit("audio", sample)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
}
|
}
|
||||||
|
Loading…
Reference in New Issue
Block a user