feat(flows): add now() to the CEL expression environment (#7194)

* feat(flows): add now() to the CEL expression environment

CEL expressions in flow definitions had no way to produce the current
date: the environment was built bare, so date-dependent flows failed at
runtime. Register a now() function that returns the current UTC time as
a CEL timestamp. The value is frozen once per kickoff so every
expression in a run sees the same instant, even across midnight.

Standard CEL covers formatting from there: string(now()),
now().getFullYear(), now() - duration('24h').

* chore(flows): drop redundant comment on _cel_now

* refactor(flows): derive CEL env and functions from one registry

A function now lives in one _CelFunctionSpec entry: its annotation for
compile and its implementation factory for evaluate, so the two cannot
drift. Run-scoped values move into _CelRunContext; adding one is a
field, not a new parameter through every helper signature.

* chore(flows): drop _CelRunContext docstring

* fix(flows): freeze a fresh cel now() on human-feedback resume

resume_async never passes through kickoff_async, so a flow restored
with from_pending() had no frozen instant and now() fell back to live
wall-clock per expression. Freeze a fresh instant at resume instead of
persisting the kickoff one: a flow can pause on feedback for days, and
expressions after resume must see today.
This commit is contained in:
Thiago Moretto
2026-09-01 14:05:47 -03:00
committed by GitHub
parent 48cc5d4e5e
commit 1bc2e0722d
3 changed files with 216 additions and 19 deletions

View File

@@ -2,7 +2,9 @@
from __future__ import annotations
from collections.abc import Iterable
from collections.abc import Callable, Iterable
from dataclasses import dataclass
from datetime import datetime, timezone
from functools import lru_cache
import json
from typing import TYPE_CHECKING, Any, NamedTuple, TypeAlias, cast
@@ -21,6 +23,46 @@ _CEL_MACROS_WITH_LOCAL_BINDINGS = frozenset(
)
@dataclass(frozen=True)
class _CelRunContext:
now: datetime
@dataclass(frozen=True)
class _CelFunctionSpec:
annotation: Any
factory: Callable[[_CelRunContext], Callable[..., Any]]
@lru_cache(maxsize=1)
def _cel_function_registry() -> dict[str, _CelFunctionSpec]:
from celpy import celtypes
return {
"now": _CelFunctionSpec(
annotation=celtypes.FunctionType,
factory=lambda run: lambda: celtypes.TimestampType(run.now),
),
}
def _cel_environment() -> Any:
from celpy import Environment
return Environment(
annotations={
name: spec.annotation for name, spec in _cel_function_registry().items()
}
)
def _cel_functions(run_context: _CelRunContext) -> dict[str, Any]:
return {
name: spec.factory(run_context)
for name, spec in _cel_function_registry().items()
}
def _find_cel_eval_error(value: Any) -> Exception | None:
from celpy.evaluation import CELEvalError
@@ -103,6 +145,9 @@ FLOW_TEMPLATE_EXPRESSION_RULES: tuple[str, ...] = (
"Use this for numbers, booleans, objects, and lists.",
"If the string has other text, the final value is text. Non-text values "
"become JSON. `null` becomes empty text.",
"Use `now()` for the current UTC time as a CEL timestamp, frozen for the "
"whole run. Use standard CEL on it: `string(now())` for ISO text, "
"`now().getFullYear()`, or `now() - duration('24h')`.",
)
FLOW_TEMPLATE_EXPRESSION_CONTRACT = " ".join(FLOW_TEMPLATE_EXPRESSION_RULES)
FLOW_TEMPLATE_EXPRESSION_EXAMPLES: dict[str, tuple[dict[str, str], ...]] = {
@@ -176,10 +221,15 @@ class Expression:
"""CEL expression helper used for definition-time checks and runtime rendering."""
def __init__(
self, value: ExpressionData, *, context: dict[str, Any] | None = None
self,
value: ExpressionData,
*,
context: dict[str, Any] | None = None,
now: datetime | None = None,
) -> None:
self.value = value
self.context = context
self.now = now
@classmethod
def from_flow(
@@ -190,7 +240,11 @@ class Expression:
local_context: dict[str, Any] | None = None,
) -> Expression:
"""Build an expression with the standard Flow runtime context."""
return cls(value, context=cls._flow_context(flow, local_context=local_context))
return cls(
value,
context=cls._flow_context(flow, local_context=local_context),
now=getattr(flow, "_cel_now", None),
)
def validate_expression(
self,
@@ -231,6 +285,7 @@ class Expression:
return self._evaluate_cel(
self._require_cel_source(cast(str, self.value)),
resolved_context or {},
self._run_context(),
)
def render_template(self, context: dict[str, Any] | None = None) -> Any:
@@ -240,7 +295,12 @@ class Expression:
type; strings mixing literals and expressions render as text.
"""
resolved_context = self.context if context is None else context
return self._render_template_value(self.value, resolved_context or {})
return self._render_template_value(
self.value, resolved_context or {}, self._run_context()
)
def _run_context(self) -> _CelRunContext:
return _CelRunContext(now=self.now or datetime.now(timezone.utc))
@staticmethod
def _validate_template_value(
@@ -311,20 +371,27 @@ class Expression:
return context
@staticmethod
def _render_template_value(value: ExpressionData, context: dict[str, Any]) -> Any:
def _render_template_value(
value: ExpressionData, context: dict[str, Any], run_context: _CelRunContext
) -> Any:
if isinstance(value, str):
return Expression._render_template_string(value, context)
return Expression._render_template_string(value, context, run_context)
if isinstance(value, dict):
return {
key: Expression._render_template_value(item, context)
key: Expression._render_template_value(item, context, run_context)
for key, item in value.items()
}
if isinstance(value, list):
return [Expression._render_template_value(item, context) for item in value]
return [
Expression._render_template_value(item, context, run_context)
for item in value
]
return value
@staticmethod
def _render_template_string(value: str, context: dict[str, Any]) -> Any:
def _render_template_string(
value: str, context: dict[str, Any], run_context: _CelRunContext
) -> Any:
segments = _parse_template_segments(value)
expressions = [
segment for segment in segments if isinstance(segment, _ExpressionSegment)
@@ -333,26 +400,28 @@ class Expression:
return value
literals = [segment for segment in segments if isinstance(segment, str)]
if len(expressions) == 1 and all(not literal.strip() for literal in literals):
return Expression._evaluate_cel(expressions[0].source, context)
return Expression._evaluate_cel(expressions[0].source, context, run_context)
rendered: list[str] = []
for segment in segments:
if isinstance(segment, str):
rendered.append(segment)
continue
result = Expression._evaluate_cel(segment.source, context)
result = Expression._evaluate_cel(segment.source, context, run_context)
rendered.append("" if result is None else _stringify_cel_value(result))
return "".join(rendered)
@staticmethod
def _evaluate_cel(expression: str, context: dict[str, Any]) -> Any:
def _evaluate_cel(
expression: str, context: dict[str, Any], run_context: _CelRunContext
) -> Any:
try:
from celpy import Environment
from celpy.adapter import CELJSONEncoder, json_to_cel
from celpy.evaluation import Context
environment = Environment()
environment = _cel_environment()
program = environment.program(
Expression._compile_cel(expression, environment=environment)
Expression._compile_cel(expression, environment=environment),
functions=_cel_functions(run_context),
)
result = program.evaluate(cast(Context, json_to_cel(context)))
if (eval_error := _find_cel_eval_error(result)) is not None:
@@ -371,9 +440,7 @@ class Expression:
environment: Any | None = None,
) -> Any:
if environment is None:
from celpy import Environment
environment = Environment()
environment = _cel_environment()
try:
return environment.compile(expression)
except Exception as e:

View File

@@ -13,7 +13,7 @@ from collections.abc import Callable, Iterator, Sequence
from concurrent.futures import Future, ThreadPoolExecutor
import contextvars
import copy
from datetime import datetime
from datetime import datetime, timezone
import enum
import inspect
import logging
@@ -772,6 +772,7 @@ class Flow(BaseModel, Generic[T], metaclass=FlowMeta):
# duration span emitted at the end does not need to hold a span open for
# the life of the run.
_telemetry_started_at: float | None = PrivateAttr(default=None)
_cel_now: datetime | None = PrivateAttr(default=None)
_event_futures: list[Future[None]] = PrivateAttr(default_factory=list)
_pending_feedback_context: PendingFeedbackContext | None = PrivateAttr(default=None)
_human_feedback_method_outputs: dict[str, Any] = PrivateAttr(default_factory=dict)
@@ -1382,6 +1383,10 @@ class Flow(BaseModel, Generic[T], metaclass=FlowMeta):
"No pending feedback context. Use from_pending() to restore a paused flow."
)
# A fresh instant, not the persisted kickoff one: a flow can pause on
# feedback for days, and expressions after resume must see today.
self._cel_now = datetime.now(timezone.utc)
execution_token = begin_execution(self._pending_feedback_context.execution_uuid)
# Force `current_flow_id` to this flow's match id for the
@@ -2171,6 +2176,8 @@ class Flow(BaseModel, Generic[T], metaclass=FlowMeta):
restore_from_state_id=restore_from_state_id,
)
self._cel_now = datetime.now(timezone.utc)
ctx = baggage.set_baggage("flow_inputs", inputs or {})
ctx = baggage.set_baggage("flow_input_files", input_files or {}, context=ctx)
flow_token = attach(ctx)

View File

@@ -3047,6 +3047,94 @@ def test_expression_keeps_short_circuited_cel_errors():
)
def test_expression_now_evaluates_with_frozen_timestamp():
from datetime import datetime, timezone
from crewai.flow.expressions import Expression
frozen = datetime(2026, 1, 15, 12, 30, tzinfo=timezone.utc)
assert Expression("now().getFullYear()", context={}, now=frozen).evaluate() == 2026
assert (
Expression("string(now())", context={}, now=frozen).evaluate()
== "2026-01-15T12:30:00Z"
)
assert (
Expression("string(now() - duration('24h'))", context={}, now=frozen).evaluate()
== "2026-01-14T12:30:00Z"
)
def test_expression_now_defaults_to_current_time():
from datetime import datetime, timezone
from crewai.flow.expressions import Expression
year = Expression("now().getFullYear()", context={}).evaluate()
assert year == datetime.now(timezone.utc).year
def test_expression_now_renders_in_templates():
from datetime import datetime, timezone
from crewai.flow.expressions import Expression
frozen = datetime(2026, 1, 15, tzinfo=timezone.utc)
rendered = Expression(
{"query": "News from ${string(now().getFullYear())}"},
context={},
now=frozen,
).render_template()
assert rendered == {"query": "News from 2026"}
def test_expression_now_passes_root_validation():
from crewai.flow.expressions import Expression
Expression("string(now().getFullYear())").validate_expression(
allowed_roots=["state", "outputs"]
)
def test_expression_from_flow_uses_run_frozen_now():
from datetime import datetime, timezone
from crewai.flow.expressions import Expression
flow = Flow()
flow._cel_now = datetime(2026, 1, 15, tzinfo=timezone.utc)
assert (
Expression.from_flow("now().getFullYear()", flow).evaluate() == 2026
)
def test_expression_action_can_use_now():
definition = FlowDefinition.from_declaration(contents=
{
"schema": "crewai.flow/v1",
"name": "NowFlow",
"methods": {
"today": {
"start": True,
"do": {
"call": "expression",
"expr": "string(now().getFullYear())",
},
}
},
}
)
from datetime import datetime, timezone
result = Flow.from_declaration(contents=definition).kickoff()
assert result == str(datetime.now(timezone.utc).year)
def test_expression_action_can_route_like_if_else():
yaml_str = f"""
schema: crewai.flow/v1
@@ -3738,6 +3826,41 @@ def test_resume_synthetic_completion_persists():
assert _saved_methods("resume-synthetic") == ["generate"]
def test_resume_freezes_fresh_cel_now():
from crewai.flow.expressions import Expression
backend = DefinitionStoreBackend(store="resume-cel-now")
frozen_at_listener: list[Any] = []
class NowResumableFlow(Flow):
@start()
@human_feedback(message="Review:")
def generate(self):
return "content"
@listen(generate)
def process(self, result):
frozen_at_listener.append(self._cel_now)
return Expression.from_flow("string(now())", self).evaluate()
context = PendingFeedbackContext(
flow_id="resume-cel-now-1",
flow_class="NowResumableFlow",
method_name="generate",
method_output="content",
message="Review:",
)
backend.save_pending_feedback("resume-cel-now-1", context, {"id": "resume-cel-now-1"})
flow = NowResumableFlow.from_pending("resume-cel-now-1", backend)
assert flow._cel_now is None
result = flow.resume("looks good")
assert frozen_at_listener[0] is not None
assert result == frozen_at_listener[0].strftime("%Y-%m-%dT%H:%M:%SZ")
class ReviewFlow(Flow):
@start()
@human_feedback(