Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,11 @@ to include examples, links to docs, or any other relevant information.

### Fixed

- Preserve empty activations in the shared External Workflow Streams input and
output replay schedule. Workflows that read input, publish decisions, and
schedule Activities now reproduce that schedule during replay. Inconsistent
prerelease markers are rejected explicitly rather than guessing where omitted
activations belonged.
- Resume external input waits when cold replay encounters a wake in an already
loaded History page, including Workers with workflow caching disabled.
- Avoid an unnecessary output replacement Workflow Task after stream input has
Expand Down
28 changes: 21 additions & 7 deletions temporalio/contrib/external_workflow_streams/_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -327,8 +327,10 @@ def __init__(
#: refused for: a fresh annotation is the most room there will ever be, so
#: refusing there rolls over to an annotation that refuses identically.
self._segments_in_annotation = 0
#: Set when a subscription is registered or a record delivered, so an
#: activation that changed nothing at all emits nothing.
#: A late first subscription must preserve the earlier empty drains in
#: the same task, without making unrelated workflows write markers.
self._unobserved_segments = 0
self._segment_pending = True
self._observed_this_activation = False
#: ``wait_id -> Future``, awaited by Workflow code and resolved by the
#: readiness activation. It lives here rather than on either side alone
Expand Down Expand Up @@ -428,12 +430,14 @@ def begin_activation(self, history_floor_event_id: int | None = None) -> None:
)
if self._output_history_floor_event_id != history_floor_event_id:
self._output_segments = []
self._unobserved_segments = 0
self._output_history_floor_event_id = history_floor_event_id
for waiter in self._output_capacity_waiters:
if not waiter.done():
waiter.set_result(None)
self._output_capacity_waiters = []
self._output_segments.append({})
self._segment_pending = True

def delivery_budget_remaining(self) -> int:
"""How many more records this activation may hand to Workflow code.
Expand Down Expand Up @@ -1638,11 +1642,21 @@ def close_segment(self, reason: SegmentEndReason | None = None) -> None:
nothing still ran one event-loop drain, and replay must reproduce that
drain or ``wait_condition`` predicates fire a different number of times.
"""
if not self._observed_this_activation:
if not self._observed_this_activation and not self._segment_pending:
return
self._segment_pending = False
if not self._observed_this_activation and not self._subscriptions:
self._unobserved_segments += 1
return
if reason is None:
reason = self._segment_end_reason()
accumulator = self._ensure_accumulator()
for _ in range(self._unobserved_segments):
self._pending_deltas.append(
accumulator.add_segment(Segment((), SegmentEndReason.NO_DATA_AVAILABLE))
)
self._segments_in_annotation += 1
self._unobserved_segments = 0
self._pending_deltas.append(
accumulator.add_segment(Segment(tuple(self._runs), reason))
)
Expand Down Expand Up @@ -1686,10 +1700,9 @@ def _segment_end_reason(self) -> SegmentEndReason:
def take_observation_delta(self) -> bytes | None:
"""The bytes to put on this completion's `WorkflowStreamProgress`.

``None`` means nothing replay-visible changed, which is the only case
where a completion legitimately carries no progress command -- and a
replay delivery is exactly that case, since everything it delivered is
already recorded in the marker being replayed.
Once an input subscription exists, even an empty activation belongs to
the shared input/output replay schedule. Workflows that never subscribe
still emit nothing, and replay does not rewrite its existing marker.
"""
if self._replay_ready is not None:
return None
Expand Down Expand Up @@ -1786,6 +1799,7 @@ def start_new_annotation(self) -> None:
self._run_sizes = []
self._max_run_bytes = 0
self._segments_in_annotation = 0
self._unobserved_segments = 0
self._observed_this_activation = False
self._annotation_start = {
wait_id: state.delivery_cursor
Expand Down
107 changes: 81 additions & 26 deletions temporalio/worker/_workflow_instance.py
Original file line number Diff line number Diff line change
Expand Up @@ -541,6 +541,14 @@ def activate(
job_sets: list[
list[temporalio.bridge.proto.workflow_activation.WorkflowActivationJob]
] = [[], [], [], []]
# An external stream replay marker is what this Workflow Task recorded,
# so it has to be installed before any of the task's code runs and
# publishes, or the install wipes those publishes and the manifest
# check fails them. It is applied right before the first drain of the
# activation, after every job that drain will act on.
replay_jobs: list[
temporalio.bridge.proto.workflow_activation.WorkflowActivationJob
] = []
for job in act.jobs:
if job.HasField("notify_has_patch"):
job_sets[0].append(job)
Expand All @@ -550,6 +558,8 @@ def activate(
# Ordered with the Signals, where Core puts them: what a
# task was woken for is known before anything it resolves.
job_sets[1].append(job)
elif job.HasField("replay_external_streams"):
replay_jobs.append(job)
elif not job.HasField("query_workflow"):
if job.HasField("initialize_workflow"):
start_job = job.initialize_workflow
Expand All @@ -567,32 +577,7 @@ def activate(
self._workflow_input = self._make_workflow_input(start_job)

try:
if self._single_batch_activation:
# Applying every job before giving workflow tasks a chance to
# run prevents their order in the activation from hiding state
# that arrived in the same workflow task.
for job_set in job_sets:
for job in job_set:
# Let errors bubble out of these to the caller to fail the task
self._apply(job)
if any(job_sets):
self._run_once(
check_conditions=bool(job_sets[1] or job_sets[2])
)
else:
# Preserve the legacy scheduling order for histories which do
# not contain the single-batch workflow logic flag.
for index, job_set in enumerate(job_sets):
if not job_set:
continue
for job in job_set:
# Let errors bubble out of these to the caller to fail the task
self._apply(job)

# Run one iteration of the loop. We do not allow conditions to
# be checked in patch jobs (first index) or query jobs (last
# index).
self._run_once(check_conditions=index == 1 or index == 2)
self._apply_activation_jobs(job_sets, replay_jobs)
except BaseException:
# An error is already on its way out, so the replay is *abandoned*
# rather than closed. Closing runs `verify_replay_consumed`, which
Expand Down Expand Up @@ -697,6 +682,66 @@ def activate(

return self._current_completion

def _apply_activation_jobs(
self,
job_sets: list[
list[temporalio.bridge.proto.workflow_activation.WorkflowActivationJob]
],
replay_jobs: list[
temporalio.bridge.proto.workflow_activation.WorkflowActivationJob
],
) -> None:
"""Apply one activation's jobs and run the drains they earn.

An external stream replay marker is installed in front of the first
drain that can publish. A drain that publishes ahead of the install has
its records wiped by it and fails the manifest check, and an install
that earns a drain of its own makes the activation run one more drain
than the task the marker was written for.
"""
if self._single_batch_activation:
# Applying every job before giving workflow tasks a chance to
# run prevents their order in the activation from hiding state
# that arrived in the same workflow task.
for job_set in job_sets:
for job in job_set:
# Let errors bubble out of these to the caller to fail the task
self._apply(job)
for job in replay_jobs:
self._apply(job)
if any(job_sets) or replay_jobs:
self._run_once(check_conditions=bool(job_sets[1] or job_sets[2]))
return

# Preserve the legacy scheduling order for histories which do
# not contain the single-batch workflow logic flag.
replay_pending = bool(replay_jobs)
# When nothing at index 1 or above will drain, the patch set's drain is
# the first one that can publish, so the install joins that set instead
# of adding a drain behind it.
first_draining_index = 1 if any(job_sets[1:]) else 0
for index, job_set in enumerate(job_sets):
if not job_set:
continue
for job in job_set:
# Let errors bubble out of these to the caller to fail the task
self._apply(job)
if replay_pending and index >= first_draining_index:
for job in replay_jobs:
self._apply(job)
replay_pending = False

# Run one iteration of the loop. We do not allow conditions to
# be checked in patch jobs (first index) or query jobs (last
# index).
self._run_once(check_conditions=index == 1 or index == 2)
if replay_pending:
# No job set drained, so the marker's drain is this one, under the
# conditions rule the single-batch branch gives the same activation.
for job in replay_jobs:
self._apply(job)
self._run_once(check_conditions=False)

def _apply(
self, job: temporalio.bridge.proto.workflow_activation.WorkflowActivationJob
) -> None:
Expand Down Expand Up @@ -1029,6 +1074,16 @@ def _apply_replay_external_streams(
runtime.begin_output_replay(output)
self._pending_output_replay_finish = True
plan = runtime.take_replay_plan()
if (
has_output
and plan is not None
and len(plan.segments) != len(output_segments)
):
raise temporalio.workflow.NondeterminismError(
"External stream History has incompatible input and output "
"activation schedules. This prerelease marker omitted empty "
"input activations, whose positions cannot be recovered safely."
)
if plan is None:
if has_output:
# An output-only marker deliberately has no input ReplayPlan.
Expand Down
132 changes: 132 additions & 0 deletions tests/contrib/external_workflow_streams/test_output_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
import temporalio.bridge.proto.workflow_completion
import temporalio.converter
import temporalio.workflow
from temporalio.contrib.external_workflow_streams._annotation import decode_annotation
from temporalio.contrib.external_workflow_streams._backend import StreamKey
from temporalio.contrib.external_workflow_streams._errors import (
ExternalStreamCapacityError,
Expand Down Expand Up @@ -1064,3 +1065,134 @@ def test_input_and_output_share_one_replay_segment_drain_schedule() -> None:
"abandon-output",
"end-input",
]


def test_ambiguous_legacy_combined_schedule_is_rejected_before_drain() -> None:
runtime = _CombinedReplayRuntime()
driver = _CombinedReplayDriver(runtime)
output = SimpleNamespace(segments=(object(), object(), object(), object()))
job = SimpleNamespace(output=output, HasField=lambda field: field == "output")

apply_replay = cast(Any, _WorkflowInstanceImpl._apply_replay_external_streams)
with pytest.raises(
temporalio.workflow.NondeterminismError, match="incompatible input and output"
):
apply_replay(driver, job)

assert runtime.events == ["begin-output"]


@pytest.mark.parametrize("subscription_exists", [False, True])
async def test_combined_marker_preserves_empty_input_activations(
runtime: WorkflowStreamRuntime, subscription_exists: bool
) -> None:
key = runtime.stream_key("inputs")
if subscription_exists:
runtime.begin_activation(1)
runtime.register(wait_id=1, stream_key=key)
runtime.take_observation_delta()
runtime.add_terminal()

deltas = []
try:
for index in range(4):
runtime.begin_activation(7)
if index == 1 and not subscription_exists:
runtime.register(wait_id=1, stream_key=key)
if index in (1, 3):
record = StreamRecord(
RecordKind.DATA, b"input", "producer", index
).placed_at(Offset(f"{index}-0"))
runtime.record_delivery(1, record)
await runtime.publish_output(
topic="events",
value=str(index),
value_type=str,
kind=RecordKind.DATA,
max_publish_latency=timedelta(seconds=1),
max_records=10,
max_logical_bytes=10_000,
)
delta = runtime.take_observation_delta()
if delta is not None:
deltas.append(delta)
deltas.append(runtime.add_terminal())
staged = await runtime.stage_output(7)
annotation = decode_annotation(b"".join(deltas))

assert [bool(segment.runs) for segment in annotation.segments] == [
False,
True,
False,
True,
]
assert staged.segment_record_counts == ((0,), (1,), (0,), (1,))
finally:
await runtime._manager.shutdown()


class _SchedulingStub:
"""The two calls the activation's job dispatch makes on the instance."""

def __init__(self, *, single_batch: bool) -> None:
self._single_batch_activation = single_batch
self.events: list[str] = []

def _apply(self, job: Any) -> None:
self.events.append(f"apply-{job.name}")

def _run_once(self, *, check_conditions: bool) -> None:
self.events.append(f"drain({check_conditions})")


def _job(name: str) -> Any:
return SimpleNamespace(name=name)


def _dispatch(stub: _SchedulingStub, job_sets: list[list[Any]], replay: list[Any]):
cast(Any, _WorkflowInstanceImpl._apply_activation_jobs)(stub, job_sets, replay)


@pytest.mark.parametrize("single_batch", [False, True])
def test_a_replay_marker_rides_the_patch_set_drain(single_batch: bool) -> None:
"""A marker beside patch jobs alone must not earn a second drain.

A patch job set drains once, and that drain is the first one that can
publish, so the install belongs in front of it. Behind it the activation
runs two drains where the recorded task ran one, and every
``wait_condition`` predicate fires an extra time.
"""
stub = _SchedulingStub(single_batch=single_batch)
_dispatch(stub, [[_job("patch")], [], [], []], [_job("marker")])

assert stub.events == ["apply-patch", "apply-marker", "drain(False)"]


def test_a_marker_still_waits_for_the_signal_set_it_precedes() -> None:
stub = _SchedulingStub(single_batch=False)
_dispatch(stub, [[_job("patch")], [_job("signal")], [], []], [_job("marker")])

assert stub.events == [
"apply-patch",
"drain(False)",
"apply-signal",
"apply-marker",
"drain(True)",
]


def test_a_query_only_activation_answers_after_the_marker() -> None:
stub = _SchedulingStub(single_batch=False)
_dispatch(stub, [[], [], [], [_job("query")]], [_job("marker")])

assert stub.events == ["apply-query", "apply-marker", "drain(False)"]


@pytest.mark.parametrize("single_batch", [False, True])
def test_a_marker_alone_drains_once_without_checking_conditions(
single_batch: bool,
) -> None:
stub = _SchedulingStub(single_batch=single_batch)
_dispatch(stub, [[], [], [], []], [_job("marker")])

assert stub.events == ["apply-marker", "drain(False)"]
10 changes: 7 additions & 3 deletions tests/contrib/external_workflow_streams/test_worker_handoff.py
Original file line number Diff line number Diff line change
Expand Up @@ -704,9 +704,6 @@ async def test_a_finalization_that_cannot_be_answered_writes_no_marker(
"no marker committed the first record, so there is no previous "
"marker for the retry to replay from",
)
before = markers(await history(handle))
committed_before = committed_boundary(before[-1], wait_id=1)
marker_before = marker_bytes(before[-1])

# Fed from here, and each record is *followed* to the Workflow rather
# than merely published on a schedule. The precondition this case needs
Expand Down Expand Up @@ -761,6 +758,13 @@ async def test_a_finalization_that_cannot_be_answered_writes_no_marker(
"since the last delivery is expected to fit inside"
)

# Read here rather than at the first marker: every Workflow Task that
# closes between the two commits its own marker, and the claim this
# case makes is about the one task that is open right now.
before = markers(await history(handle))
committed_before = committed_boundary(before[-1], wait_id=1)
marker_before = marker_bytes(before[-1])

# The Run's entry disappears underneath the finalization, exactly once.
original = workflow_worker_impl._handle_external_stream_jobs
sabotaged: list[str] = []
Expand Down
Loading