diff --git a/openflexure_microscope/__init__.py b/openflexure_microscope/__init__.py index 0b72e60d..c23ac0c0 100644 --- a/openflexure_microscope/__init__.py +++ b/openflexure_microscope/__init__.py @@ -1,4 +1,4 @@ -__all__ = ['Microscope', 'config'] +__all__ = ['Microscope', 'config', 'task', 'lock'] __version__ = "0.1.0" from .microscope import Microscope diff --git a/openflexure_microscope/api/app.py b/openflexure_microscope/api/app.py index dc655e4f..8d4ca70c 100644 --- a/openflexure_microscope/api/app.py +++ b/openflexure_microscope/api/app.py @@ -23,8 +23,9 @@ from flask_cors import CORS from openflexure_microscope.api.utilities import list_routes from openflexure_microscope import Microscope, config +from openflexure_microscope.lock import LockError from openflexure_microscope.camera.pi import StreamingCamera -from openflexure_stage import OpenFlexureStage +from openflexure_microscope.stage.openflexure import Stage import atexit import logging, sys @@ -66,6 +67,16 @@ def _handle_http_exception(e): for code in default_exceptions: app.errorhandler(code)(_handle_http_exception) +def _handle_lock_exception(e): + return make_response( + jsonify({ + 'status_code': 423, + 'error': e.code, + 'details': e.message + }), + 423) + +app.errorhandler(LockError)(_handle_lock_exception) # After app starts, but before first request, attach hardware to global microscope @app.before_first_request @@ -77,7 +88,7 @@ def attach_microscope(): logging.debug("Creating camera object...") api_camera = StreamingCamera(config=openflexurerc) logging.debug("Creating stage object...") - api_stage = OpenFlexureStage("/dev/ttyUSB0") + api_stage = Stage("/dev/ttyUSB0") logging.debug("Attaching devices to microscope...") api_microscope.attach( diff --git a/openflexure_microscope/camera/base.py b/openflexure_microscope/camera/base.py index cb8bf69b..07e5ee5e 100644 --- a/openflexure_microscope/camera/base.py +++ b/openflexure_microscope/camera/base.py @@ -19,6 +19,7 @@ except ImportError: from .capture import CaptureObject, capture_from_dict, BASE_CAPTURE_PATH from openflexure_microscope.config import USER_CONFIG_DIR from openflexure_microscope.utilities import entry_by_id +from openflexure_microscope.lock import StrictLock def last_entry(object_list: list): @@ -29,7 +30,6 @@ def last_entry(object_list: list): return None - def generate_basename(): """Return a default filename based on the capture datetime""" return datetime.datetime.now().strftime("%Y-%m-%d_%H-%M-%S") @@ -87,6 +87,8 @@ class BaseCamera(object): self.thread = None #: Background thread reading frames from camera self.camera = None #: Camera object + self.lock = StrictLock(timeout=1) #: Strict lock controlling thread access to camera hardware + self.frame = None #: bytes: Current frame is stored here by background thread self.last_access = 0 #: time: Time of last client access to the camera self.event = CameraEvent() diff --git a/openflexure_microscope/camera/pi.py b/openflexure_microscope/camera/pi.py index 7b4a52ee..5458856e 100644 --- a/openflexure_microscope/camera/pi.py +++ b/openflexure_microscope/camera/pi.py @@ -201,70 +201,73 @@ class StreamingCamera(BaseCamera): paused_stream = False - # Apply valid config params to Picamera object - if not self.state['record_active']: # If not recording a video + with self.lock: - logging.debug("Applying config:") - logging.debug(config) + # Apply valid config params to Picamera object + if not self.state['record_active']: # If not recording a video - # Pause stream while changing settings - if self.state['stream_active']: # If stream is active - logging.info("Pausing stream to update config.") - self.stop_stream_recording() # Pause stream - paused_stream = True # Remember to unpause stream when done + logging.debug("Applying config:") + logging.debug(config) - # PiCamera parameters (applied directly to PiCamera object) - if 'picamera_params' in config: - for key in PICAMERA_KEYS: - if key in config['picamera_params']: - logging.debug("Setting parameter {}: {}".format(key, config['picamera_params'][key])) - setattr(self.camera, key, config['picamera_params'][key]) + # Pause stream while changing settings + if self.state['stream_active']: # If stream is active + logging.info("Pausing stream to update config.") + self.stop_stream_recording() # Pause stream + paused_stream = True # Remember to unpause stream when done - # StreamingCamera parameters (applied via StreamingCamera config) - for key in CONFIG_KEYS: - if key in config: - logging.debug("Setting parameter {}: {}".format(key, config[key])) - self._config[key] = config[key] + # PiCamera parameters (applied directly to PiCamera object) + if 'picamera_params' in config: + for key in PICAMERA_KEYS: + if key in config['picamera_params']: + logging.debug("Setting parameter {}: {}".format(key, config['picamera_params'][key])) + setattr(self.camera, key, config['picamera_params'][key]) - # Richard's library to set analog and digital gains - if 'analog_gain' in config: - set_analog_gain(self.camera, config['analog_gain']) - if 'digital_gain' in config: - set_digital_gain(self.camera, config['digital_gain']) + # StreamingCamera parameters (applied via StreamingCamera config) + for key in CONFIG_KEYS: + if key in config: + logging.debug("Setting parameter {}: {}".format(key, config[key])) + self._config[key] = config[key] - # Handle lens shading, if a new one is given - if 'shading_table_path' in config: - self.apply_shading_table() + # Richard's library to set analog and digital gains + if 'analog_gain' in config: + set_analog_gain(self.camera, config['analog_gain']) + if 'digital_gain' in config: + set_digital_gain(self.camera, config['digital_gain']) - # If stream was paused to update config, unpause - if paused_stream: - logging.info("Resuming stream.") - self.start_stream_recording() + # Handle lens shading, if a new one is given + if 'shading_table_path' in config: + self.apply_shading_table() - else: - raise Exception( - "Cannot update camera config while recording is active.") + # If stream was paused to update config, unpause + if paused_stream: + logging.info("Resuming stream.") + self.start_stream_recording() + + else: + raise Exception( + "Cannot update camera config while recording is active.") def set_zoom(self, zoom_value: float=1.) -> None: """ Change the camera zoom, handling recentering and scaling. """ - self.state['zoom_value'] = float(zoom_value) - if self.state['zoom_value'] < 1: - self.state['zoom_value'] = 1 - # Richard's code for zooming ! - fov = self.camera.zoom - centre = np.array([fov[0] + fov[2]/2.0, fov[1] + fov[3]/2.0]) - size = 1.0/self.state['zoom_value'] - # If the new zoom value would be invalid, move the centre to - # keep it within the camera's sensor (this is only relevant - # when zooming out, if the FoV is not centred on (0.5, 0.5) - for i in range(2): - if np.abs(centre[i] - 0.5) + size/2 > 0.5: - centre[i] = 0.5 + (1.0 - size)/2 * np.sign(centre[i]-0.5) - logging.info("setting zoom, centre {}, size {}".format(centre, size)) - new_fov = (centre[0] - size/2, centre[1] - size/2, size, size) - self.camera.zoom = new_fov + with self.lock: + self.state['zoom_value'] = float(zoom_value) + if self.state['zoom_value'] < 1: + self.state['zoom_value'] = 1 + # Richard's code for zooming ! + fov = self.camera.zoom + centre = np.array([fov[0] + fov[2]/2.0, fov[1] + fov[3]/2.0]) + size = 1.0/self.state['zoom_value'] + # If the new zoom value would be invalid, move the centre to + # keep it within the camera's sensor (this is only relevant + # when zooming out, if the FoV is not centred on (0.5, 0.5) + for i in range(2): + if np.abs(centre[i] - 0.5) + size/2 > 0.5: + centre[i] = 0.5 + (1.0 - size)/2 * np.sign(centre[i]-0.5) + logging.info("setting zoom, centre {}, size {}".format(centre, size)) + new_fov = (centre[0] - size/2, centre[1] - size/2, size, size) + self.camera.zoom = new_fov # LAUNCH ACTIONS @@ -299,48 +302,50 @@ class StreamingCamera(BaseCamera): output_object (str/BytesIO): Target object. """ - # Start recording method only if a current recording is not running - if not self.state['record_active']: + with self.lock: + # Start recording method only if a current recording is not running + if not self.state['record_active']: + + # If output is a StreamObject + if isinstance(output, CaptureObject): + # Lock the StreamObject while recording + output.lock() + # Set target to capture stream + output_stream = output.stream + else: + output_stream = output + + # Start the camera video recording on port 2 + logging.info("Recording to {}".format(output)) + + self.camera.start_recording( + output_stream, + format=fmt, + splitter_port=2, + resize=self.config['video_resolution'], + quality=quality) + + # Update state dictionary + self.state['record_active'] = True + + return output - # If output is a StreamObject - if isinstance(output, CaptureObject): - # Lock the StreamObject while recording - output.lock() - # Set target to capture stream - output_stream = output.stream else: - output_stream = output - - # Start the camera video recording on port 2 - logging.info("Recording to {}".format(output)) - - self.camera.start_recording( - output_stream, - format=fmt, - splitter_port=2, - resize=self.config['video_resolution'], - quality=quality) - - # Update state dictionary - self.state['record_active'] = True - - return output - - else: - logging.error( - "Cannot start a new recording\ - until the current recording has stopped.") - return None + logging.error( + "Cannot start a new recording\ + until the current recording has stopped.") + return None def stop_recording(self) -> bool: """Stop the last started video recording on splitter port 2.""" - # Stop the camera video recording on port 2 - logging.info("Stopping recording") - self.camera.stop_recording(splitter_port=2) - logging.info("Recording stopped") + with self.lock: + # Stop the camera video recording on port 2 + logging.info("Stopping recording") + self.camera.stop_recording(splitter_port=2) + logging.info("Recording stopped") - # Update state dictionary - self.state['record_active'] = False + # Update state dictionary + self.state['record_active'] = False def stop_stream_recording( self, @@ -378,7 +383,6 @@ class StreamingCamera(BaseCamera): splitter_port (int): Splitter port to start recording on resolution ((int, int)): Resolution to set the camera to, before starting recording. Defaults to `self.config['video_resolution']`. """ - # If no stream object exists if not hasattr(self, 'stream'): self.stream = io.BytesIO() # Create a stream object @@ -421,32 +425,33 @@ class StreamingCamera(BaseCamera): resize ((int, int)): Resize the captured image. """ - # If output is a StreamObject - if isinstance(output, CaptureObject): - # Set target to capture stream - output_stream = output.stream - else: - output_stream = output + with self.lock: + # If output is a StreamObject + if isinstance(output, CaptureObject): + # Set target to capture stream + output_stream = output.stream + else: + output_stream = output - logging.info("Capturing to {}".format(output)) + logging.info("Capturing to {}".format(output)) - # Set resolution and stop stream recording if necesarry - if not use_video_port: - self.stop_stream_recording() + # Set resolution and stop stream recording if necesarry + if not use_video_port: + self.stop_stream_recording() - self.camera.capture( - output_stream, - format=fmt, - quality=100, - resize=resize, - bayer=(not use_video_port) and bayer, - use_video_port=use_video_port) + self.camera.capture( + output_stream, + format=fmt, + quality=100, + resize=resize, + bayer=(not use_video_port) and bayer, + use_video_port=use_video_port) - # Set resolution and start stream recording if necesarry - if not use_video_port: - self.start_stream_recording() + # Set resolution and start stream recording if necesarry + if not use_video_port: + self.start_stream_recording() - return output + return output def yuv( self, @@ -459,38 +464,39 @@ class StreamingCamera(BaseCamera): use_video_port (bool): Capture from the video port used for streaming. Lower resolution, faster. resize ((int, int)): Resize the captured image. """ - if use_video_port: - resolution = self.config['video_resolution'] - else: - resolution = self.config['numpy_resolution'] + with self.lock: + if use_video_port: + resolution = self.config['video_resolution'] + else: + resolution = self.config['numpy_resolution'] - if resize: - size = resize - else: - size = resolution - - if not use_video_port: - self.stop_stream_recording(resolution=resolution) - - logging.debug("Creating PiYUVArray") - with picamera.array.PiYUVArray(self.camera, size=size) as output: - - logging.info("Capturing to {}".format(output)) - - self.camera.capture( - output, - resize=size, - format='yuv', - use_video_port=use_video_port) + if resize: + size = resize + else: + size = resolution if not use_video_port: - self.start_stream_recording() + self.stop_stream_recording(resolution=resolution) - if rgb: - logging.debug("Converting to RGB") - return output.rgb_array - else: - return output.array + logging.debug("Creating PiYUVArray") + with picamera.array.PiYUVArray(self.camera, size=size) as output: + + logging.info("Capturing to {}".format(output)) + + self.camera.capture( + output, + resize=size, + format='yuv', + use_video_port=use_video_port) + + if not use_video_port: + self.start_stream_recording() + + if rgb: + logging.debug("Converting to RGB") + return output.rgb_array + else: + return output.array def array( self, @@ -502,35 +508,35 @@ class StreamingCamera(BaseCamera): use_video_port (bool): Capture from the video port used for streaming. Lower resolution, faster. resize ((int, int)): Resize the captured image. """ + with self.lock: + if use_video_port: + resolution = self.config['video_resolution'] + else: + resolution = self.config['numpy_resolution'] - if use_video_port: - resolution = self.config['video_resolution'] - else: - resolution = self.config['numpy_resolution'] + if resize: + size = resize + else: + size = resolution - if resize: - size = resize - else: - size = resolution + # Always pause stream, to prevent resizer memory issues + self.stop_stream_recording(resolution=resolution) - # Always pause stream, to prevent resizer memory issues - self.stop_stream_recording(resolution=resolution) + logging.debug("Creating PiRGBArray") + with picamera.array.PiRGBArray(self.camera, size=size) as output: - logging.debug("Creating PiRGBArray") - with picamera.array.PiRGBArray(self.camera, size=size) as output: + logging.info("Capturing to {}".format(output)) - logging.info("Capturing to {}".format(output)) + self.camera.capture( + output, + resize=size, + format='rgb', + use_video_port=use_video_port) - self.camera.capture( - output, - resize=size, - format='rgb', - use_video_port=use_video_port) + # Resume stream + self.start_stream_recording() - # Resume stream - self.start_stream_recording() - - return output.array + return output.array # HANDLE STREAM FRAMES diff --git a/openflexure_microscope/lock.py b/openflexure_microscope/lock.py new file mode 100644 index 00000000..af040cd1 --- /dev/null +++ b/openflexure_microscope/lock.py @@ -0,0 +1,65 @@ +from threading import RLock, ThreadError + +class LockError(ThreadError): + ERROR_CODES = { + 'ACQUIRE_ERROR': "Unable to acquire. Lock in use by another thread.", + 'IN_USE_ERROR': "Lock in use by another thread." + } + + def __init__(self, code, lock): + self.code = code + if code in LockError.ERROR_CODES: + self.message = LockError.ERROR_CODES[code] + else: + self.message = "Unknown error." + print("{}: {}".format(self.code, self.message)) + + ThreadError.__init__(self) + + +class StrictLock(object): + def __init__(self, timeout=1): + self._lock = RLock() + self.timeout = timeout + + def acquire(self, blocking=True): + return self._lock.acquire(blocking, timeout=self.timeout) + + def __enter__(self): + result = self._lock.acquire(blocking=True, timeout=self.timeout) + if result: + return result + else: + raise LockError('ACQUIRE_ERROR', self) + + def __exit__(self, *args): + self._lock.release() + + def release(self): + self._lock.release() + + +class CompositeLock(object): + def __init__(self, locks, timeout=1): + self.locks = locks + self.timeout = timeout + + def acquire(self, blocking=True): + return (lock.acquire(blocking=blocking, timeout=self.timeout) for lock in self.locks) + + def __enter__(self): + result = (lock.acquire(blocking=True, timeout=self.timeout) for lock in self.locks) + if all(result): + return result + else: + msg = "Unable to acquire all locks {}.".format(self.locks) + logging.error(msg) + raise ThreadError(msg) + + def __exit__(self, *args): + for lock in self.locks: + lock.release() + + def release(self): + for lock in self.locks: + lock.release() \ No newline at end of file diff --git a/openflexure_microscope/plugins/example/api.py b/openflexure_microscope/plugins/example/api.py index a55254ed..f3da64ff 100644 --- a/openflexure_microscope/plugins/example/api.py +++ b/openflexure_microscope/plugins/example/api.py @@ -67,6 +67,7 @@ class LongRunningAPI(MicroscopeViewPlugin): # Extract a values from the JSON payload. time_to_run = payload.param('time', default=10, convert=int) + # Attach the long-running method as a microscope task try: task = self.microscope.task.start(self.plugin.long_running, time_to_run) return jsonify(task.state), 202 diff --git a/openflexure_microscope/plugins/example/plugin.py b/openflexure_microscope/plugins/example/plugin.py index 203aca19..3638a57f 100644 --- a/openflexure_microscope/plugins/example/plugin.py +++ b/openflexure_microscope/plugins/example/plugin.py @@ -39,10 +39,12 @@ class Plugin(MicroscopePlugin): print("Starting a long-running task...") n_array = [] - for _ in range(t_run): - n_array.append(random.random()) - time.sleep(1) + print("Acquiring camera lock...") + with self.microscope.camera.lock: + for _ in range(t_run): + n_array.append(random.random()) + time.sleep(1) - print("Long-running task finished!") + print("Long-running task finished! Releasing lock.") return n_array \ No newline at end of file diff --git a/openflexure_microscope/stage/__init__.py b/openflexure_microscope/stage/__init__.py new file mode 100644 index 00000000..5aab8d20 --- /dev/null +++ b/openflexure_microscope/stage/__init__.py @@ -0,0 +1 @@ +from . import openflexure \ No newline at end of file diff --git a/openflexure_microscope/stage/openflexure.py b/openflexure_microscope/stage/openflexure.py new file mode 100644 index 00000000..2703da35 --- /dev/null +++ b/openflexure_microscope/stage/openflexure.py @@ -0,0 +1,10 @@ +from openflexure_stage import OpenFlexureStage + +from openflexure_microscope.lock import StrictLock + +# TODO: Implement lock on movement +class Stage(OpenFlexureStage): + def __init__(self, *args, **kwargs): + self.lock = StrictLock(timeout=2) #: Strict lock controlling thread access to camera hardware + + OpenFlexureStage.__init__(self, *args, **kwargs) \ No newline at end of file diff --git a/openflexure_microscope/task.py b/openflexure_microscope/task.py index 83b32103..7684936c 100644 --- a/openflexure_microscope/task.py +++ b/openflexure_microscope/task.py @@ -57,10 +57,10 @@ class TaskOrchestrator: def start_condition(self, task_obj): """ Method returning bool describing if a particular task is allowed to run. - In the future, this will depend on the state of hardware locks? - Currently just allows a single task at a time. + Currently always allows tasks to start, as locks determine access to hardware. """ - return not any([task._running for task in self.tasks]) + #return not any([task._running for task in self.tasks]) + return True class Task: