diff --git a/internal/staticsources/rpicamera/camera.go b/internal/staticsources/rpicamera/camera.go index c06d55a0..9db131f4 100644 --- a/internal/staticsources/rpicamera/camera.go +++ b/internal/staticsources/rpicamera/camera.go @@ -17,12 +17,13 @@ type camera struct { Params params OnData func(time.Duration, [][]byte) - cmd *exec.Cmd - pipeOut *pipe - pipeIn *pipe + cmd *exec.Cmd + pipeOut *pipe + pipeIn *pipe + finalErr error - waitDone chan error - readerDone chan error + terminate chan struct{} + done chan struct{} } func (c *camera) initialize() error { @@ -64,78 +65,92 @@ func (c *camera) initialize() error { return err } + c.terminate = make(chan struct{}) + c.done = make(chan struct{}) + + go c.run() + c.pipeOut.write(append([]byte{'c'}, c.Params.serialize()...)) - c.waitDone = make(chan error) - go func() { - c.waitDone <- c.cmd.Wait() - }() - - c.readerDone = make(chan error) - go func() { - c.readerDone <- c.readReady() - }() - - select { - case err := <-c.waitDone: - c.pipeOut.close() - c.pipeIn.close() - <-c.readerDone - freeComponent() - return fmt.Errorf("process exited unexpectedly: %v", err) - - case err := <-c.readerDone: - if err != nil { - c.pipeOut.write([]byte{'e'}) - <-c.waitDone - c.pipeOut.close() - c.pipeIn.close() - freeComponent() - return err - } - } - - c.readerDone = make(chan error) - go func() { - c.readerDone <- c.readData() - close(c.readerDone) - }() - return nil } func (c *camera) close() { - c.pipeOut.write([]byte{'e'}) - <-c.waitDone - c.pipeOut.close() - c.pipeIn.close() - <-c.readerDone + close(c.terminate) + <-c.done freeComponent() } -func (c *camera) reloadParams(params params) { - c.pipeOut.write(append([]byte{'c'}, params.serialize()...)) +func (c *camera) run() { + defer close(c.done) + c.finalErr = c.runInner() } -func (c *camera) readReady() error { - buf, err := c.pipeIn.read() - if err != nil { - return err - } +func (c *camera) runInner() error { + cmdDone := make(chan error) + go func() { + cmdDone <- c.cmd.Wait() + }() - switch buf[0] { - case 'e': - return fmt.Errorf(string(buf[1:])) + readDone := make(chan error) + go func() { + readDone <- c.runReader() + }() - case 'r': - return nil + for { + select { + case err := <-cmdDone: + c.pipeIn.close() + c.pipeOut.close() - default: - return fmt.Errorf("unexpected output from video pipe: '0x%.2x'", buf[0]) + <-readDone + + return err + + case err := <-readDone: + c.pipeIn.close() + + c.pipeOut.write([]byte{'e'}) + c.pipeOut.close() + + <-cmdDone + + return err + + case <-c.terminate: + c.pipeIn.close() + <-readDone + + c.pipeOut.write([]byte{'e'}) + c.pipeOut.close() + + <-cmdDone + + return fmt.Errorf("terminated") + } } } -func (c *camera) readData() error { +func (c *camera) runReader() error { +outer: + for { + buf, err := c.pipeIn.read() + if err != nil { + return err + } + + switch buf[0] { + case 'e': + return fmt.Errorf(string(buf[1:])) + + case 'r': + break outer + + default: + return fmt.Errorf("unexpected data from pipe: '0x%.2x'", buf[0]) + } + } + for { buf, err := c.pipeIn.read() if err != nil { @@ -159,11 +174,16 @@ func (c *camera) readData() error { c.OnData(dts, nalus) default: - return fmt.Errorf("unexpected output from pipe (%c)", buf[0]) + return fmt.Errorf("unexpected data from pipe: '0x%.2x'", buf[0]) } } } -func (c *camera) error() chan error { - return c.readerDone +func (c *camera) reloadParams(params params) { + c.pipeOut.write(append([]byte{'c'}, params.serialize()...)) +} + +func (c *camera) wait() error { + <-c.done + return c.finalErr } diff --git a/internal/staticsources/rpicamera/camera_disabled.go b/internal/staticsources/rpicamera/camera_disabled.go index 461649d6..ac84ee59 100644 --- a/internal/staticsources/rpicamera/camera_disabled.go +++ b/internal/staticsources/rpicamera/camera_disabled.go @@ -22,6 +22,6 @@ func (c *camera) close() { func (c *camera) reloadParams(_ params) { } -func (c *camera) error() chan error { +func (c *camera) wait() error { return nil } diff --git a/internal/staticsources/rpicamera/source.go b/internal/staticsources/rpicamera/source.go index 9331300f..491c444a 100644 --- a/internal/staticsources/rpicamera/source.go +++ b/internal/staticsources/rpicamera/source.go @@ -117,6 +117,12 @@ func (s *Source) Run(params defs.StaticSourceRunParams) error { }) } + defer func() { + if stream != nil { + s.Parent.SetNotReady(defs.PathSourceStaticSetNotReadyReq{}) + } + }() + cam := &camera{ Params: paramsFromConf(s.LogLevel, params.Conf), OnData: onData, @@ -127,15 +133,14 @@ func (s *Source) Run(params defs.StaticSourceRunParams) error { } defer cam.close() - defer func() { - if stream != nil { - s.Parent.SetNotReady(defs.PathSourceStaticSetNotReadyReq{}) - } + cameraErr := make(chan error) + go func() { + cameraErr <- cam.wait() }() for { select { - case err := <-cam.error(): + case err := <-cameraErr: return err case cnf := <-params.ReloadConf: