# -*- coding: utf-8 -*- import time import os import shutil import threading import datetime import logging import threading import gevent from abc import ABCMeta, abstractmethod from labthings.core.lock import StrictLock from labthings.core.event import ClientEvent class BaseCamera(metaclass=ABCMeta): """ Base implementation of StreamingCamera. """ def __init__(self): self.thread = None self.camera = None self.lock = StrictLock(timeout=1, name="Camera") self.frame = None self.last_access = 0 self.event = ClientEvent() self.stop = False # Used to indicate that the stream loop should break self.stream_timeout = 20 self.stream_timeout_enabled = False self.stream_active = False self.record_active = False @property @abstractmethod def configuration(self): """The current camera configuration.""" pass @property @abstractmethod def state(self): """The current read-only camera state.""" pass @property def settings(self): return self.read_settings() @abstractmethod def update_settings(self, config: dict): """Update settings from a config dictionary""" with self.lock: # Apply valid config params to camera object for key, value in config.items(): # For each provided setting if hasattr(self, key): # If the instance has a matching property setattr(self, key, value) # Set to the target value @abstractmethod def read_settings(self) -> dict: """Return the current settings as a dictionary""" return {} def __enter__(self): """Create camera on context enter.""" return self def __exit__(self, exc_type, exc_value, traceback): """Close camera stream on context exit.""" self.close() def close(self): """Close the BaseCamera and all attached StreamObjects.""" logging.info("Closing {}".format(self)) # Stop worker thread self.stop_worker() logging.info("Closed {}".format(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: # Spawn a greenlet to handle stream # start background frame thread self.thread = gevent.spawn(self._thread) # 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.""" pass # WORKER THREAD def _thread(self): """Camera background thread.""" self.frames_iterator = self.frames() logging.debug("Entering worker thread.") self.stream_active = True for frame in self.frames_iterator: self.frame = frame self.event.set() # send signal to clients # Handle timeout if ( self.stream_timeout_enabled and ( # If using timeout time.time() - self.last_access > self.stream_timeout ) and not self.preview_active # And GPU preview is not active ): self.frames_iterator.close() break try: if self.stop is True: logging.debug("Worker thread flagged for stop.") self.frames_iterator.close() break except AttributeError: pass logging.debug("BaseCamera worker thread exiting...") # Set stream_activate state self.stream_active = False