rpicamera: fix restarting camera when component crashes (#3997)
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user