Skip to content

How to write a plan that runs until stopped

Here you write a plan that loops until the user stops it, declare actions the user can trigger while it runs, and start, pause and stop it from its plan widget. Continuous plans explains how actions are offered and taken.

Prerequisites

You need a presenter that runs the plans of the session and a view that shows their plan widgets, as in How to run a plan from a presenter. This guide adds to both. The blocks below are parts of one script, and the whole script is at the end. The text before each block names the class it belongs to.

The plans take a camera that can be triggered and read and has a shutter signal, through this protocol:

@runtime_checkable
class Camera(Readable[Any], Triggerable, Protocol):
    shutter: SignalRW[bool]

Mark the plan continuous

In the component that offers the plan, decorate it with continuous and loop until it is stopped. In MyController:

@continuous(pausable=True)
def live(self, camera: Camera) -> MsgGenerator[None]:
    yield from bps.open_run()
    while True:
        yield from bps.checkpoint()
        yield from bps.trigger_and_read([camera])

The Run button of its plan widget becomes a toggle that starts and stops the plan, and pausable=True adds a button to pause and resume it. After a pause, the plan starts again from the checkpoint.

Declare actions

You declare an action as a PlanAction given as the default of a parameter annotated PlanAction:

SNAP = PlanAction(name="snap", description="Take one frame")
SHUTTER = PlanAction(name="shutter", toggle_states=("Open", "Close"))

A PlanAction with toggle_states gets a button that stays pressed. The states are the labels shown while the button is released and while it is pressed.

The component that offers the plan owns an ActionManager in self.actions, and the plan waits on it. In MyController:

@continuous
def snapshots(
    self, camera: Camera, snap: PlanAction = SNAP, shutter: PlanAction = SHUTTER
) -> MsgGenerator[None]:
    yield from bps.open_run()
    while True:
        name = yield from self.actions.wait(snap, shutter)
        try:
            if name == snap.name:
                yield from bps.trigger_and_read([camera])
            else:
                yield from bps.mv(camera.shutter, True)
                yield from self.actions.wait_released(shutter)
        finally:
            if name == shutter.name:
                yield from bps.mv(camera.shutter, False)
            self.actions.done(name)

wait returns the name of the action the user asked for, and doesn't time out. wait_released waits until the user releases a toggle button. Call done in a finally block, so that an action running when the plan is stopped goes back to idle too. A finally block may yield messages, so a plan stopped with the shutter open closes it.

Two actions with one name raise

create_plan_spec raises ValueError when a plan declares two actions of the same name. Give each action of a plan a name of its own.

Start, pause and stop it

In PlanPresenter, keep the futures of the plan that runs, add a slot for the toggle and one for the pause button, and stop the plan at shutdown:

@slot
def run(self, plan: str, values: dict[str, Any]) -> None:
    if self.futures:
        self.logger.warning(f"A plan is running; {plan!r} not started")
        return
    resolved = resolve_arguments(self.specs[plan], values, self.devices)
    args, kwargs = collect_arguments(self.specs[plan], resolved)
    self.watch(self.engine(self.plans[plan]["plan"](*args, **kwargs)))

@slot
def toggle(self, plan: str, on: bool, values: dict[str, Any]) -> None:
    if on:
        self.run(plan, values)
    elif self.engine.state != "idle":
        self.watch(self.engine.stop())

@slot
def pause(self, paused: bool) -> None:
    if paused:
        self.engine.request_pause(defer=True)
    else:
        self.watch(self.engine.resume())

def watch(self, future: Future[Any]) -> None:
    self.futures.add(future)
    future.add_done_callback(self.finished)

def finished(self, future: Future[Any]) -> None:
    self.futures.discard(future)
    if not self.futures and self.engine.state != "paused":
        self.sig_finished.emit()

def shutdown(self) -> None:
    if self.futures or self.engine.state == "paused":
        wait([self.engine.stop()], timeout=10)

The engine returns a Future for each plan it starts, and watch keeps it in self.futures, an empty set made in __init__, until it completes. Step through what each button does to the engine and to the futures kept:

Start, pause and stop a continuous plan

A plan starts only from idle, because run refuses a second plan while a Future is kept.

idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here.
idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here.

Toggle on calls run(). The engine runs the plan and watch keeps the Future it returned.

idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here.runningits Future keptThe engine runs the plan and watch keeps the Future it returned. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here. The engine runs the plan and watch keeps the Future it returned.
idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here.runningits Future keptThe engine runs the plan and watch keeps the Future it returned. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here. The engine runs the plan and watch keeps the Future it returned.

Pause calls request_pause(defer=True), and the plan pauses at its next checkpoint. Its Future completes with RunEngineInterrupted, but the engine is paused, so finished sends nothing.

idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here.runningits Future keptThe engine runs the plan and watch keeps the Future it returned.pausedno Future keptThe plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here. The engine runs the plan and watch keeps the Future it returned. The plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing.
idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here.runningits Future keptThe engine runs the plan and watch keeps the Future it returned.pausedno Future keptThe plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here. The engine runs the plan and watch keeps the Future it returned. The plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing.

Resume calls resume(), which returns a new Future for watch to keep. The plan carries on from the checkpoint.

idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here.runningits Future keptresume returns a new Future, which watch keeps. The plan starts again from the checkpoint.pausedno Future keptThe plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here. resume returns a new Future, which watch keeps. The plan starts again from the checkpoint. The plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing.
idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here.runningits Future keptresume returns a new Future, which watch keeps. The plan starts again from the checkpoint.pausedno Future keptThe plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here. resume returns a new Future, which watch keeps. The plan starts again from the checkpoint. The plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing.

Toggle off calls stop(), from running or paused. finished sends sig_finished once the last Future kept is complete and the engine isn't paused, so once for each plan.

idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here.runningits Future keptresume returns a new Future, which watch keeps. The plan starts again from the checkpoint.pausedno Future keptThe plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing.sig_finishedsent oncestop returns a Future that completes once the plan has cleaned up. The stop button stays enabled while the plan is paused, so you stop a paused plan the same way. finished sends sig_finished when the last Future kept is complete and the engine is not paused, so once for each plan. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here. resume returns a new Future, which watch keeps. The plan starts again from the checkpoint. The plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing. stop returns a Future that completes once the plan has cleaned up. The stop button stays enabled while the plan is paused, so you stop a paused plan the same way. finished sends sig_finished when the last Future kept is complete and the engine is not paused, so once for each plan.
idleno Future keptrun refuses a second plan while a Future is kept, so a plan starts only from here.runningits Future keptresume returns a new Future, which watch keeps. The plan starts again from the checkpoint.pausedno Future keptThe plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing.sig_finishedsent oncestop returns a Future that completes once the plan has cleaned up. The stop button stays enabled while the plan is paused, so you stop a paused plan the same way. finished sends sig_finished when the last Future kept is complete and the engine is not paused, so once for each plan. toggle on: run()Pause: request_pause(defer=True)Resume: resume()toggle off: stop(),or the plan ends or failstoggle off: stop()run refuses a second plan while a Future is kept, so a plan starts only from here. resume returns a new Future, which watch keeps. The plan starts again from the checkpoint. The plan pauses at its next checkpoint. Its Future completes with a RunEngineInterrupted exception, and the engine stays in the state paused, so finished sends nothing. stop returns a Future that completes once the plan has cleaned up. The stop button stays enabled while the plan is paused, so you stop a paused plan the same way. finished sends sig_finished when the last Future kept is complete and the engine is not paused, so once for each plan.

request_pause comes from bluesky without type annotations, so mypy --strict reports it as no-untyped-call.

In PlanView, pass create_plan_widget a callback for each button, and update the plan widget when they are pressed and when the plan ends:

def add_plan(self, spec: PlanSpec) -> None:
    widget = create_plan_widget(
        spec,
        run_callback=lambda: self.ask_to_run(spec.name),
        toggle_callback=lambda on: self.ask_to_toggle(spec.name, on),
        pause_callback=lambda paused: self.ask_to_pause(spec.name, paused),
        action_clicked_callback=self.ask,
        action_toggled_callback=self.ask_or_release,
    )
    self.widgets[spec.name] = widget
    self.chooser.addItem(spec.name)
    self.pages.addWidget(widget.group_box)

def ask_to_run(self, plan: str) -> None:
    self.setEnabled(False)
    self.sig_run.emit(plan, self.widgets[plan].parameters)

def ask_to_toggle(self, plan: str, on: bool) -> None:
    self.chooser.setEnabled(not on)
    self.widgets[plan].toggle(on)
    self.sig_toggle.emit(plan, on, self.widgets[plan].parameters)

def ask_to_pause(self, plan: str, paused: bool) -> None:
    self.widgets[plan].pause(paused)
    self.sig_pause.emit(paused)

@slot
def on_finished(self) -> None:
    self.setEnabled(True)
    self.chooser.setEnabled(True)
    self.widgets[self.chooser.currentText()].toggle(False)

PlanWidget.toggle sets the label of the toggle, enables the action buttons and the pause button while the plan runs, and locks the inputs of the parameters. on_finished calls toggle(False) so that a plan that fails or ends by itself shows as stopped. Disabling the combo box keeps the user from starting a second plan meanwhile.

Follow the actions

In PlanView, ask for an action when its button is pressed, and set each button from the state the ActionManager reports:

def ask(self, name: str) -> None:
    self.sig_action_request.emit(name, True)

def ask_or_release(self, checked: bool, name: str) -> None:
    self.sig_action_request.emit(name, checked)

@slot
def on_action_changed(self, name: str, state: str) -> None:
    for widget in self.widgets.values():
        if name in widget.action_buttons:
            self.set_action_button(widget.action_buttons[name], state)

def set_action_button(self, button: ActionButton, state: str) -> None:
    match state:
        case ActionState.IDLE:
            button.setEnabled(False)
            button.release()
        case ActionState.OFFERED:
            button.setEnabled(True)
        case ActionState.RUNNING:
            button.setEnabled(button.isCheckable())

How to follow a plan action from a view explains the three states and why the view calls release.

Link the view to the presenter, and to the ActionManager of the component that offers the plan:

class MyApp(QtSession):
    camera: AsDevice[MyCamera]
    ctrl: AsPresenter[MyController]
    plan_ctrl: AsPresenter[PlanPresenter]
    plan_view: AsView[PlanView]

    def wire(self) -> Iterator[Link]:
        yield self.plan_view.sig_run, self.plan_ctrl.run
        yield self.plan_view.sig_toggle, self.plan_ctrl.toggle
        yield self.plan_view.sig_pause, self.plan_ctrl.pause
        yield self.plan_ctrl.sig_finished, self.plan_view.on_finished
        yield self.plan_view.sig_action_request, self.ctrl.actions.request
        yield self.ctrl.actions.sig_changed, self.plan_view.on_action_changed
        yield self.plan_ctrl.sig_progress, self.plan_view.on_progress


if __name__ == "__main__":
    MyApp().run()

The example in full

The whole script
"""The session of the guide "How to write a plan that runs until stopped"."""

from __future__ import annotations

from collections.abc import Iterator, Mapping  # noqa: TC003
from concurrent.futures import Future, wait
from typing import Any, Protocol, runtime_checkable

import bluesky.plan_stubs as bps
from bluesky.protocols import Readable, Triggerable
from bluesky.utils import MsgGenerator  # noqa: TC002
from ophyd_async.core import AsyncStatus, SignalRW, StandardReadable, soft_signal_rw
from psygnal import Signal
from qtpy.QtWidgets import QComboBox, QStackedWidget, QVBoxLayout, QWidget

import redsun.engine.plan_stubs as rps
from redsun import (
    AsDevice,
    AsPresenter,
    AsView,
    DeviceMapping,
    HasPlans,
    Link,
    Placement,
    PlanEntry,
    slot,
)
from redsun.engine import ProgressState, RunEngine
from redsun.engine.actions import ActionManager, ActionState, PlanAction, continuous
from redsun.log import Loggable
from redsun.presenter.plan_spec import (
    PlanSpec,
    UnresolvableAnnotationError,
    collect_arguments,
    create_plan_spec,
    resolve_arguments,
)
from redsun.qt import Dock, QtSession
from redsun.view.qt.utils import ActionButton, PlanWidget, create_plan_widget


class MyCamera(StandardReadable):
    def __init__(self, name: str = "") -> None:
        with self.add_children_as_readables():
            self.frames = soft_signal_rw(int)
        self.shutter = soft_signal_rw(bool)
        super().__init__(name=name)

    @AsyncStatus.wrap
    async def trigger(self) -> None:
        await self.frames.set(await self.frames.get_value() + 1)


@runtime_checkable
class Camera(Readable[Any], Triggerable, Protocol):
    shutter: SignalRW[bool]


SNAP = PlanAction(name="snap", description="Take one frame")
SHUTTER = PlanAction(name="shutter", toggle_states=("Open", "Close"))


class MyController:
    def __init__(self, name: str) -> None:
        self.name = name
        self.actions = ActionManager()

    @continuous(pausable=True)
    def live(self, camera: Camera) -> MsgGenerator[None]:
        yield from bps.open_run()
        while True:
            yield from bps.checkpoint()
            yield from bps.trigger_and_read([camera])

    @continuous
    def snapshots(
        self, camera: Camera, snap: PlanAction = SNAP, shutter: PlanAction = SHUTTER
    ) -> MsgGenerator[None]:
        yield from bps.open_run()
        while True:
            name = yield from self.actions.wait(snap, shutter)
            try:
                if name == snap.name:
                    yield from bps.trigger_and_read([camera])
                else:
                    yield from bps.mv(camera.shutter, True)
                    yield from self.actions.wait_released(shutter)
            finally:
                if name == shutter.name:
                    yield from bps.mv(camera.shutter, False)
                self.actions.done(name)

    def series(
        self, camera: Camera, frames: int = 10, repeats: int = 2
    ) -> MsgGenerator[None]:
        yield from bps.open_run()
        yield from rps.declare_progress("repeats")
        for repeat in range(repeats):
            status = yield from bps.abs_set(
                camera.shutter, True, wait=False, group="shutter"
            )
            yield from rps.monitor_progress("shutter", status, parent="repeats")
            yield from bps.wait(group="shutter")
            yield from rps.declare_progress("series", parent="repeats")
            for frame in range(frames):
                yield from bps.trigger_and_read([camera])
                yield from rps.update_progress(
                    "series", current=frame + 1, initial=0, target=frames, unit="frames"
                )
            yield from rps.update_progress("series", done=True)
            yield from rps.update_progress(
                "repeats", current=repeat + 1, initial=0, target=repeats, unit="repeats"
            )
        yield from rps.update_progress("repeats", done=True)
        yield from bps.close_run()

    def plan_map(self) -> Mapping[str, PlanEntry]:
        return {
            "live": {"plan": self.live},
            "snapshots": {"plan": self.snapshots},
            "series": {"plan": self.series},
        }


class PlanPresenter(Loggable):
    sig_finished = Signal()
    sig_progress = Signal(tuple)

    def __init__(self, name: str, *, devices: DeviceMapping) -> None:
        self.name = name
        self.devices = devices
        self.engine = RunEngine()
        self.plans: dict[str, PlanEntry] = {}
        self.specs: dict[str, PlanSpec] = {}
        self.futures: set[Future[Any]] = set()
        self.engine.sig_progress.connect(self.sig_progress.emit)

    def setup(self, plan_sources: Mapping[str, HasPlans]) -> None:
        for component in plan_sources.values():
            for plan, entry in component.plan_map().items():
                try:
                    self.specs[plan] = create_plan_spec(entry["plan"], self.devices)
                except (UnresolvableAnnotationError, ValueError) as error:
                    self.logger.warning(error)
                    continue
                self.plans[plan] = entry

    @slot
    def run(self, plan: str, values: dict[str, Any]) -> None:
        if self.futures:
            self.logger.warning(f"A plan is running; {plan!r} not started")
            return
        resolved = resolve_arguments(self.specs[plan], values, self.devices)
        args, kwargs = collect_arguments(self.specs[plan], resolved)
        self.watch(self.engine(self.plans[plan]["plan"](*args, **kwargs)))

    @slot
    def toggle(self, plan: str, on: bool, values: dict[str, Any]) -> None:
        if on:
            self.run(plan, values)
        elif self.engine.state != "idle":
            self.watch(self.engine.stop())

    @slot
    def pause(self, paused: bool) -> None:
        if paused:
            self.engine.request_pause(defer=True)
        else:
            self.watch(self.engine.resume())

    def watch(self, future: Future[Any]) -> None:
        self.futures.add(future)
        future.add_done_callback(self.finished)

    def finished(self, future: Future[Any]) -> None:
        self.futures.discard(future)
        if not self.futures and self.engine.state != "paused":
            self.sig_finished.emit()

    def shutdown(self) -> None:
        if self.futures or self.engine.state == "paused":
            wait([self.engine.stop()], timeout=10)



class PlanView(QWidget):
    placement: Placement = Dock("right")
    sig_run = Signal(str, dict)
    sig_toggle = Signal(str, bool, dict)
    sig_pause = Signal(bool)
    sig_action_request = Signal(str, bool)

    def __init__(self, name: str, parent: QWidget) -> None:
        super().__init__(parent)
        self.name = name
        self.chooser = QComboBox()
        self.pages = QStackedWidget()
        self.chooser.currentIndexChanged.connect(self.pages.setCurrentIndex)
        layout = QVBoxLayout(self)
        layout.addWidget(self.chooser)
        layout.addWidget(self.pages)
        self.widgets: dict[str, PlanWidget] = {}

    def setup(
        self, plan_sources: Mapping[str, HasPlans], devices: DeviceMapping
    ) -> None:
        for component in plan_sources.values():
            for entry in component.plan_map().values():
                try:
                    self.add_plan(create_plan_spec(entry["plan"], devices))
                except (UnresolvableAnnotationError, ValueError):
                    continue

    def add_plan(self, spec: PlanSpec) -> None:
        widget = create_plan_widget(
            spec,
            run_callback=lambda: self.ask_to_run(spec.name),
            toggle_callback=lambda on: self.ask_to_toggle(spec.name, on),
            pause_callback=lambda paused: self.ask_to_pause(spec.name, paused),
            action_clicked_callback=self.ask,
            action_toggled_callback=self.ask_or_release,
        )
        self.widgets[spec.name] = widget
        self.chooser.addItem(spec.name)
        self.pages.addWidget(widget.group_box)

    def ask_to_run(self, plan: str) -> None:
        self.setEnabled(False)
        self.sig_run.emit(plan, self.widgets[plan].parameters)

    def ask_to_toggle(self, plan: str, on: bool) -> None:
        self.chooser.setEnabled(not on)
        self.widgets[plan].toggle(on)
        self.sig_toggle.emit(plan, on, self.widgets[plan].parameters)

    def ask_to_pause(self, plan: str, paused: bool) -> None:
        self.widgets[plan].pause(paused)
        self.sig_pause.emit(paused)

    @slot
    def on_finished(self) -> None:
        self.setEnabled(True)
        self.chooser.setEnabled(True)
        self.widgets[self.chooser.currentText()].toggle(False)

    def ask(self, name: str) -> None:
        self.sig_action_request.emit(name, True)

    def ask_or_release(self, checked: bool, name: str) -> None:
        self.sig_action_request.emit(name, checked)

    @slot
    def on_action_changed(self, name: str, state: str) -> None:
        for widget in self.widgets.values():
            if name in widget.action_buttons:
                self.set_action_button(widget.action_buttons[name], state)

    def set_action_button(self, button: ActionButton, state: str) -> None:
        match state:
            case ActionState.IDLE:
                button.setEnabled(False)
                button.release()
            case ActionState.OFFERED:
                button.setEnabled(True)
            case ActionState.RUNNING:
                button.setEnabled(button.isCheckable())

    @slot
    def on_progress(self, scopes: tuple[ProgressState, ...]) -> None:
        self.widgets[self.chooser.currentText()].show_progress(scopes)



class MyApp(QtSession):
    camera: AsDevice[MyCamera]
    ctrl: AsPresenter[MyController]
    plan_ctrl: AsPresenter[PlanPresenter]
    plan_view: AsView[PlanView]

    def wire(self) -> Iterator[Link]:
        yield self.plan_view.sig_run, self.plan_ctrl.run
        yield self.plan_view.sig_toggle, self.plan_ctrl.toggle
        yield self.plan_view.sig_pause, self.plan_ctrl.pause
        yield self.plan_ctrl.sig_finished, self.plan_view.on_finished
        yield self.plan_view.sig_action_request, self.ctrl.actions.request
        yield self.ctrl.actions.sig_changed, self.plan_view.on_action_changed
        yield self.plan_ctrl.sig_progress, self.plan_view.on_progress


if __name__ == "__main__":
    MyApp().run()