diff --git a/CHANGELOG.md b/CHANGELOG.md index 4479e2b33..3e87bcade 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/temporalio/contrib/external_workflow_streams/_runtime.py b/temporalio/contrib/external_workflow_streams/_runtime.py index 739e8111b..e57b3aa91 100644 --- a/temporalio/contrib/external_workflow_streams/_runtime.py +++ b/temporalio/contrib/external_workflow_streams/_runtime.py @@ -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 @@ -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. @@ -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)) ) @@ -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 @@ -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 diff --git a/temporalio/worker/_workflow_instance.py b/temporalio/worker/_workflow_instance.py index 7924d5ca0..c7a8ae7ca 100644 --- a/temporalio/worker/_workflow_instance.py +++ b/temporalio/worker/_workflow_instance.py @@ -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) @@ -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 @@ -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 @@ -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: @@ -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. diff --git a/tests/contrib/external_workflow_streams/test_output_runtime.py b/tests/contrib/external_workflow_streams/test_output_runtime.py index 1c58854a7..c194bcf0d 100644 --- a/tests/contrib/external_workflow_streams/test_output_runtime.py +++ b/tests/contrib/external_workflow_streams/test_output_runtime.py @@ -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, @@ -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)"] diff --git a/tests/contrib/external_workflow_streams/test_worker_handoff.py b/tests/contrib/external_workflow_streams/test_worker_handoff.py index 3e52796c0..cf6b212a0 100644 --- a/tests/contrib/external_workflow_streams/test_worker_handoff.py +++ b/tests/contrib/external_workflow_streams/test_worker_handoff.py @@ -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 @@ -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] = []