Commit 021745da authored by Joel Collins's avatar Joel Collins
Browse files

Remove separate camera stream thread

parent 94bc49cd
Loading
Loading
Loading
Loading
+3 −1
Original line number Diff line number Diff line
@@ -7,8 +7,10 @@ import time

# Look for debug flag
if "-d" in sys.argv or "--debug" in sys.argv:
    debug_app = True
    log_level = logging.DEBUG
else:
    debug_app = False
    log_level = logging.INFO


@@ -208,4 +210,4 @@ if __name__ == "__main__":

    logging.info("Starting OpenFlexure Microscope Server...")
    server = Server(app)
    server.run(host="0.0.0.0", port=5000, debug=False, zeroconf=True)
    server.run(host="0.0.0.0", port=5000, debug=debug_app, zeroconf=True)
+2 −6
Original line number Diff line number Diff line
@@ -7,7 +7,7 @@ def gen(camera):
    """Video streaming generator function."""
    while True:
        # the obtained frame is a jpeg
        frame = camera.get_frame()
        frame = camera.stream.getframe()

        yield (b"--frame\r\n" b"Content-Type: image/jpeg\r\n\r\n" + frame + b"\r\n")

@@ -33,8 +33,6 @@ class MjpegStream(PropertyView):
        will block all proceeding requests.
        """
        microscope = find_component("org.openflexure.microscope")
        # Restart stream worker thread
        microscope.camera.start_worker()

        return Response(
            gen(microscope.camera), mimetype="multipart/x-mixed-replace; boundary=frame"
@@ -53,7 +51,5 @@ class SnapshotStream(PropertyView):
        Single snapshot from the camera stream
        """
        microscope = find_component("org.openflexure.microscope")
        # Restart stream worker thread
        microscope.camera.start_worker()

        return Response(microscope.camera.get_frame(), mimetype="image/jpeg")
        return Response(microscope.camera.stream.getframe(), mimetype="image/jpeg")
+2 −100
Original line number Diff line number Diff line
# -*- coding: utf-8 -*-
import io
import logging
import threading
import time
from abc import ABCMeta, abstractmethod
from collections import namedtuple
@@ -36,7 +35,8 @@ class FrameStream(io.BytesIO):
        self.tracking = False

        # Event to track if a new frame is available since the last getvalue() call
        self.new_frame = threading.Event()
        # We use a ClientEvent so that each thread can call getvalue() independantly
        self.new_frame = ClientEvent()

    def start_tracking(self):
        if not self.tracking:
@@ -90,18 +90,12 @@ class BaseCamera(metaclass=ABCMeta):
    """

    def __init__(self):
        self.thread = None
        self.camera = None

        self.lock = StrictLock(name="Camera", timeout=None)

        self.stream = FrameStream()

        self.frame = None
        self.last_access = 0
        self.event = ClientEvent()
        self.stop = False  # Used to indicate that the stream loop should break

        self.stream_active = False
        self.record_active = False
        self.preview_active = False
@@ -144,96 +138,4 @@ class BaseCamera(metaclass=ABCMeta):

    def close(self):
        """Close the BaseCamera and all attached StreamObjects."""
        logging.info("Closing %s", (self))
        # Stop worker thread
        self.stop_worker()
        logging.info("Closed %s", (self))

    # START AND STOP WORKER THREAD

    def start_worker(self, timeout: int = 5) -> bool:
        """Start the background camera thread if it isn't running yet."""
        timeout_time = time.time() + timeout

        self.last_access = time.time()
        self.stop = False

        if not self.stream_active:
            # Start background frame thread
            self.thread = threading.Thread(target=self._thread)
            self.thread.daemon = True
            self.thread.start()

            # wait until frames are available
            logging.info("Waiting for frames")
            while self.get_frame() is None:
                if time.time() > timeout_time:
                    raise TimeoutError("Timeout waiting for frames.")
                else:
                    time.sleep(0.1)
        return True

    def stop_worker(self, timeout: int = 5) -> bool:
        """Flag worker thread for stop. Waits for thread close or timeout."""
        logging.debug("Stopping worker thread")
        timeout_time = time.time() + timeout

        if self.stream_active:
            self.stop = True
            self.thread.join()  # Wait for stream thread to exit
            logging.debug("Waiting for stream thread to exit.")

        while self.stream_active:
            if time.time() > timeout_time:
                logging.debug("Timeout waiting for worker thread close.")
                raise TimeoutError("Timeout waiting for worker thread close.")
            else:
                time.sleep(0.1)
        return True

    # HANDLE STREAM FRAMES

    def get_frame(self):
        """Return the current camera frame."""
        self.last_access = time.time()

        # wait for a signal from the camera thread
        self.event.wait()
        self.event.clear()

        return self.frame

    @abstractmethod
    def frames(self):
        """Create generator that returns frames from the camera."""

    # WORKER THREAD

    def _thread(self):
        """Camera background thread."""
        # Set the camera object's frame iterator
        frames_iterator = self.frames()
        logging.debug("Entering worker thread.")

        self.stream_active = True

        for frame in frames_iterator:
            # Store most recent frame
            self.frame = frame
            # Signal to clients that a new frame is available
            # We use this event because each client could be
            # reading frames slower than we're acquiring them.
            self.event.set()  # send signal to clients

            try:
                if self.stop is True:
                    logging.debug("Worker thread flagged for stop.")
                    frames_iterator.close()
                    break

            except AttributeError:
                pass

        logging.debug("BaseCamera worker thread exiting...")
        # Set stream_activate state
        self.stream_active = False
+59 −28
Original line number Diff line number Diff line
@@ -3,6 +3,7 @@
from __future__ import division

import logging
import threading
import time
from datetime import datetime

@@ -35,7 +36,11 @@ class MissingCamera(BaseCamera):
        self.generate_new_dummy_image()

        # Start streaming
        self.stop = False  # Used to indicate that the stream loop should break
        self.start_worker()
        # Wait until frames are available
        logging.info("Waiting for frames")
        self.stream.new_frame.wait()

    def generate_new_dummy_image(self):
        # Create a dummy image to serve in the stream
@@ -55,6 +60,60 @@ class MissingCamera(BaseCamera):

        image.save(self.stream, format="JPEG")

    def start_worker(self, **_) -> bool:
        """Start the background camera thread if it isn't running yet."""
        self.stop = False

        if not self.stream_active:
            # Start background frame thread
            self.thread = threading.Thread(target=self._thread)
            self.thread.daemon = True
            self.thread.start()
        return True

    def stop_worker(self, timeout: int = 5) -> bool:
        """Flag worker thread for stop. Waits for thread close or timeout."""
        logging.debug("Stopping worker thread")
        timeout_time = time.time() + timeout

        if self.stream_active:
            self.stop = True
            self.thread.join()  # Wait for stream thread to exit
            logging.debug("Waiting for stream thread to exit.")

        while self.stream_active:
            if time.time() > timeout_time:
                logging.debug("Timeout waiting for worker thread close.")
                raise TimeoutError("Timeout waiting for worker thread close.")
            else:
                time.sleep(0.1)
        return True

    def _thread(self):
        """Camera background thread."""
        # Set the camera object's frame iterator
        logging.debug("Entering worker thread.")

        self.stream_active = True

        while True:
            # Only serve frames at 1fps
            time.sleep(1)
            # Generate new dummy image
            self.generate_new_dummy_image()

            try:
                if self.stop is True:
                    logging.debug("Worker thread flagged for stop.")
                    break

            except AttributeError:
                pass

        logging.debug("BaseCamera worker thread exiting...")
        # Set stream_activate state
        self.stream_active = False

    @property
    def configuration(self):
        """The current camera configuration."""
@@ -65,9 +124,6 @@ class MissingCamera(BaseCamera):
        """The current read-only camera state."""
        return {}

    def initialisation(self):
        """Run any initialisation code when the frame iterator starts."""

    def close(self):
        """Close the Raspberry Pi PiCameraStreamer."""
        # Run BaseCamera close method
@@ -181,28 +237,3 @@ class MissingCamera(BaseCamera):
                output.close()
            else:
                output.flush()

    # HANDLE STREAM FRAMES

    def frames(self):
        """
        Create generator that returns frames from the camera.
        """
        # Run this initialisation method
        self.initialisation()

        # Update state
        logging.debug("STREAM ACTIVE")

        # While the iterator is not closed
        try:
            while True:
                time.sleep(1)  # Only serve frames at 1fps
                # Generate new dummy image
                self.generate_new_dummy_image()
                # Wait for the next frame and then yield it
                yield self.stream.getframe()

        # When GeneratorExit or StopIteration raised, run cleanup code
        finally:
            logging.debug("FRAME ITERATOR END")
+67 −93
Original line number Diff line number Diff line
@@ -107,8 +107,14 @@ class PiCameraStreamer(BaseCamera):
            "picamera_lst.npy"
        )  #: str: Path of .npy lens shading table file

        # Start the stream worker on init
        self.start_worker()
        # Run this initialisation method
        self._wait_for_camera()
        # Start stream recording (and set resolution)
        self.start_stream_recording()
        # Wait until frames are available
        logging.info("Waiting for frames")
        self.stream.new_frame.wait()
        logging.info("Frames incoming!")

    @property
    def configuration(self):
@@ -120,17 +126,25 @@ class PiCameraStreamer(BaseCamera):
        """The current read-only camera state."""
        return {}

    def initialisation(self):
        """Run any initialisation code when the frame iterator starts."""

    def close(self):
        """Close the Raspberry Pi PiCameraStreamer."""
        # Stop stream recording
        self.stop_stream_recording()
        # Run BaseCamera close method
        BaseCamera.close(self)
        super().close()
        # Detach Pi camera
        if self.camera:
            self.camera.close()

    def _wait_for_camera(self, timeout=5):
        """Wait for camera object, with 5 second timeout."""
        timeout_time = time.time() + timeout
        while not self.camera:
            if time.time() > timeout_time:
                raise TimeoutError("Timeout waiting for camera")
            else:
                pass

    # HANDLE SETTINGS
    def read_settings(self) -> dict:
        """
@@ -392,36 +406,6 @@ class PiCameraStreamer(BaseCamera):
            # Update state
            self.record_active = False

    def stop_stream_recording(self, splitter_port: int = 1, **kwargs) -> None:
        """
        Sets the camera resolution to the still-image resolution, and stops recording if the stream is active.

        Args:
            splitter_port (int): Splitter port to stop recording on
        """
        for k in kwargs.keys():
            logging.warning(
                "Warning, kwarg %s is invalid for stop_stream_recording.", k
            )
        with self.lock:
            # Stop the camera video recording on port 1
            try:
                self.camera.stop_recording(splitter_port=splitter_port)
            except picamerax.exc.PiCameraNotRecording:
                logging.info("Not recording on splitter_port %s", (splitter_port))
            else:
                logging.info(
                    "Stopped MJPEG stream on port %s. Switching to %s.",
                    splitter_port,
                    self.image_resolution,
                )

            # Increase the resolution for taking an image
            time.sleep(
                0.2
            )  # Sprinkled a sleep to prevent camera getting confused by rapid commands
            self.camera.resolution = self.image_resolution

    def start_stream_recording(self, splitter_port: int = 1, **kwargs) -> None:
        """
        Sets the camera resolution to the video/stream resolution, and starts recording if the stream should be active.
@@ -447,7 +431,6 @@ class PiCameraStreamer(BaseCamera):
                time.sleep(0.2)

            # If the stream should be active
            if self.stream_active:
            try:
                # Start recording on stream port
                self.camera.start_recording(
@@ -460,16 +443,46 @@ class PiCameraStreamer(BaseCamera):
                    splitter_port=splitter_port,
                )
            except picamerax.exc.PiCameraAlreadyRecording:
                    logging.info(
                        "Error while starting preview: Recording already running."
                    )
                logging.info("Error while starting preview: Recording already running.")
            else:
                self.stream_active = True
                logging.debug(
                    "Started MJPEG stream at %s on port %s",
                    self.stream_resolution,
                    splitter_port,
                )

    def stop_stream_recording(self, splitter_port: int = 1, **kwargs) -> None:
        """
        Sets the camera resolution to the still-image resolution, and stops recording if the stream is active.

        Args:
            splitter_port (int): Splitter port to stop recording on
        """
        for k in kwargs.keys():
            logging.warning(
                "Warning, kwarg %s is invalid for stop_stream_recording.", k
            )
        with self.lock:
            # Stop the camera video recording on port 1
            try:
                self.camera.stop_recording(splitter_port=splitter_port)
            except picamerax.exc.PiCameraNotRecording:
                logging.info("Not recording on splitter_port %s", (splitter_port))
            else:
                self.stream_active = False
                logging.info(
                    "Stopped MJPEG stream on port %s. Switching to %s.",
                    splitter_port,
                    self.image_resolution,
                )

            # Increase the resolution for taking an image
            time.sleep(
                0.2
            )  # Sprinkled a sleep to prevent camera getting confused by rapid commands
            self.camera.resolution = self.image_resolution

    def capture(
        self,
        output,
@@ -534,42 +547,3 @@ class PiCameraStreamer(BaseCamera):
                logging.info("Capturing to %s", (output))
                self.camera.capture(output, format="rgb", use_video_port=use_video_port)
                return output.array

    # HANDLE STREAM FRAMES

    def wait_for_camera(self, timeout=5):
        """Wait for camera object, with 5 second timeout."""
        timeout_time = time.time() + timeout
        while not self.camera:
            if time.time() > timeout_time:
                raise TimeoutError("Timeout waiting for camera")
            else:
                pass

    def frames(self):
        """
        Create generator that returns frames from the camera.

        Records video from port 1 to a byte stream,
        and iterates sequential frames.
        """
        # Run this initialisation method
        self.initialisation()
        self.wait_for_camera()

        # Start stream recording (and set resolution)
        self.start_stream_recording()

        logging.debug("STREAM ACTIVE")

        # While the iterator is not closed
        try:
            while True:
                # Wait for the next frame and then yield it
                yield self.stream.getframe()
        # When GeneratorExit or StopIteration raised, run cleanup code
        finally:
            # Stop stream recording (and set resolution)
            self.stop_stream_recording()

            logging.debug("FRAME ITERATOR END")