Skip to content

Delay

Delay

Bases: Node

Waits before the flow continues: for a duration, or until a time taken from the input.

The time to continue at is fixed on the first run and kept in the checkpoint, so neither a pause, an early resume nor a repeated one moves it. A wait of up to a minute happens in place. A longer one pauses the run: nodes that depend on the Delay wait with it, independent branches finish, and the run resumes when the time comes. Pausing needs checkpointing and a Delay at the top level of the flow; without them a longer wait fails with a message saying so, rather than holding a worker for hours.

Everything the node receives passes through to its output, with waited_until (the time it waited for) and waited_seconds (how long it waited) added.

Source code in dynamiq/nodes/operators/delay.py
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
class Delay(Node):
    """Waits before the flow continues: for a duration, or until a time taken from the input.

    The time to continue at is fixed on the first run and kept in the checkpoint, so neither a pause, an
    early resume nor a repeated one moves it. A wait of up to a minute happens in place. A longer one pauses
    the run: nodes that depend on the Delay wait with it, independent branches finish, and the run resumes
    when the time comes. Pausing needs checkpointing and a Delay at the top level of the flow; without them a
    longer wait fails with a message saying so, rather than holding a worker for hours.

    Everything the node receives passes through to its output, with `waited_until` (the time it waited for)
    and `waited_seconds` (how long it waited) added.
    """

    name: str | None = "delay"
    group: Literal[NodeGroup.OPERATORS] = NodeGroup.OPERATORS
    duration_seconds: float | None = Field(
        default=None, ge=0, description="How long to wait when the input gives neither `until` nor `duration_seconds`."
    )
    input_schema: ClassVar[type[DelayInputSchema]] = DelayInputSchema

    _wait_started_at: datetime | None = PrivateAttr(default=None)
    _resume_at: datetime | None = PrivateAttr(default=None)

    def to_checkpoint_state(self) -> DelayCheckpointState:
        state = super().to_checkpoint_state()
        return DelayCheckpointState(
            **state.model_dump(), wait_started_at=self._wait_started_at, resume_at=self._resume_at
        )

    def from_checkpoint_state(self, state: BaseCheckpointState | dict[str, Any]) -> None:
        super().from_checkpoint_state(state)
        restored = DelayCheckpointState.model_validate(state if isinstance(state, dict) else state.model_dump())
        self._wait_started_at = restored.wait_started_at
        self._resume_at = restored.resume_at

    def execute(self, input_data: DelayInputSchema, config: RunnableConfig = None, **kwargs) -> dict[str, Any]:
        """Waits until the time the first run fixed, in place or by pausing the run, then passes the input on."""
        config = ensure_config(config)
        self.run_on_node_execute_run(config.callbacks, **kwargs)

        now = datetime.now(timezone.utc)
        if not self.is_resumed or self._resume_at is None:
            self._wait_started_at = now
            self._resume_at = self._resolve_resume_at(input_data, now)

        remaining = (self._resume_at - now).total_seconds()
        if remaining > INLINE_WAIT_LIMIT_SECONDS:
            self._pause_run(config)
        if remaining > 0:
            self._wait_in_place(remaining, config)

        # The extra fields as received, like Pass: dumping them would turn files into iterators.
        passed_on = dict(input_data.model_extra or {})
        return passed_on | {
            "waited_until": self._resume_at.isoformat(),
            "waited_seconds": round((datetime.now(timezone.utc) - self._wait_started_at).total_seconds(), 3),
        }

    def _resolve_resume_at(self, input_data: DelayInputSchema, now: datetime) -> datetime:
        if input_data.until is not None:
            until = input_data.until
            return until if until.tzinfo else until.replace(tzinfo=timezone.utc)

        duration = input_data.duration_seconds if input_data.duration_seconds is not None else self.duration_seconds
        if duration is None:
            raise ValueError(
                f"Delay '{self.name or self.id}' has nothing to wait for: set `duration_seconds` on the node, "
                "or pass `until` or `duration_seconds` as input."
            )
        try:
            return now + timedelta(seconds=duration)
        except OverflowError as e:
            raise ValueError(f"Delay '{self.name or self.id}': {duration} seconds is too far in the future.") from e

    def _pause_run(self, config: RunnableConfig) -> None:
        """Pauses the run until the wait is over, or explains why this run cannot pause."""
        context = config.checkpoint.context if config.checkpoint else None
        if context is not None and context.pause_run(self.id, resume_at=self._resume_at):
            raise RunPausedException(
                f"Delay '{self.name or self.id}' waits until {self._resume_at.isoformat()}", resume_at=self._resume_at
            )
        raise ValueError(
            f"Delay '{self.name or self.id}' would wait until {self._resume_at.isoformat()}, longer than the "
            f"{INLINE_WAIT_LIMIT_SECONDS} seconds that can be waited in place. A longer wait pauses the run, which "
            "needs checkpointing and the Delay at the top level of the flow, not inside a Map, an agent or a "
            "sub-workflow. To test the flow without waiting, mock this node."
        )

    @staticmethod
    def _wait_in_place(seconds: float, config: RunnableConfig) -> None:
        deadline = time.monotonic() + seconds
        while (remaining := deadline - time.monotonic()) > 0:
            check_cancellation(config)
            time.sleep(min(INLINE_WAIT_POLL_SECONDS, remaining))

execute(input_data, config=None, **kwargs)

Waits until the time the first run fixed, in place or by pausing the run, then passes the input on.

Source code in dynamiq/nodes/operators/delay.py
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
def execute(self, input_data: DelayInputSchema, config: RunnableConfig = None, **kwargs) -> dict[str, Any]:
    """Waits until the time the first run fixed, in place or by pausing the run, then passes the input on."""
    config = ensure_config(config)
    self.run_on_node_execute_run(config.callbacks, **kwargs)

    now = datetime.now(timezone.utc)
    if not self.is_resumed or self._resume_at is None:
        self._wait_started_at = now
        self._resume_at = self._resolve_resume_at(input_data, now)

    remaining = (self._resume_at - now).total_seconds()
    if remaining > INLINE_WAIT_LIMIT_SECONDS:
        self._pause_run(config)
    if remaining > 0:
        self._wait_in_place(remaining, config)

    # The extra fields as received, like Pass: dumping them would turn files into iterators.
    passed_on = dict(input_data.model_extra or {})
    return passed_on | {
        "waited_until": self._resume_at.isoformat(),
        "waited_seconds": round((datetime.now(timezone.utc) - self._wait_started_at).total_seconds(), 3),
    }

DelayCheckpointState

Bases: BaseCheckpointState

The wait a Delay fixed on its first run, kept so a resumed run waits for the same time.

Source code in dynamiq/nodes/operators/delay.py
35
36
37
38
39
class DelayCheckpointState(BaseCheckpointState):
    """The wait a Delay fixed on its first run, kept so a resumed run waits for the same time."""

    wait_started_at: datetime | None = None
    resume_at: datetime | None = None

DelayInputSchema

Bases: BaseModel

What a Delay waits for. Fields other than these pass through to its output.

Source code in dynamiq/nodes/operators/delay.py
22
23
24
25
26
27
28
29
30
31
32
class DelayInputSchema(BaseModel):
    """What a Delay waits for. Fields other than these pass through to its output."""

    model_config = ConfigDict(extra="allow")

    until: datetime | None = Field(
        default=None, description="Time to continue at, as ISO-8601. A time without a zone is taken as UTC."
    )
    duration_seconds: float | None = Field(
        default=None, ge=0, description="How long to wait. Overrides the node's own duration_seconds."
    )