Skip to content

Put the native stream commits on the unified Core. - #4

Closed
moedash wants to merge 42 commits into
moe/AI-198-core-unifiedfrom
moe/AI-198-core-all
Closed

moedash wants to merge 42 commits into
moe/AI-198-core-unifiedfrom
moe/AI-198-core-all

Conversation

@moedash

@moedash moedash commented Sep 16, 2026 •

Copy link
Copy Markdown
Owner

This PR replays the native stream commits onto the unified Core as linear cherry-picks.

What changed?

Review the stream content on #6 and #7. This PR is where the two lineages meet.

Part of AI-198 (epic AI-37).

Why?

One Core that speaks both prototypes means one Python build runs every provider. The merge base with #7's branch is upstream, so merging it would bring the whole native series in a second time, and cherry-picks avoid that. Changes the two chains share reach this branch by merging #2 forward.

How did you test it?

The full cargo suite ran, both stream test families included, and the workspace unit suite, cargo fmt --check and cargo test-lint ran again on the head that carries the deferred subscribe and the unsubscribe. The replay worker was fed a read-then-publish history with its slices, with a wrong slice and with none, and a missing re-supply fails before the first activation. The Python bridge builds against this head. moedash/sdk-python#6 pins it and runs its whole suite on it.

  • Unit Tests
  • Staging
  • End to End Tests

@moedash
moedash force-pushed the moe/AI-198-core-unified branch from ddd150d to 104a15d Compare September 16, 2026 21:52
@moedash
moedash force-pushed the moe/AI-198-core-all branch 2 times, most recently from bdf72f6 to d9c8bb4 Compare September 16, 2026 22:02
@moedash
moedash force-pushed the moe/AI-198-core-unified branch from 104a15d to f4d6ee1 Compare September 16, 2026 22:07
@moedash
moedash force-pushed the moe/AI-198-core-all branch from d9c8bb4 to 4f4c8a4 Compare September 16, 2026 22:07
@moedash
moedash force-pushed the moe/AI-198-core-unified branch from f4d6ee1 to ce28a0e Compare September 16, 2026 23:13
@moedash
moedash force-pushed the moe/AI-198-core-all branch 5 times, most recently from e9b3c4a to 02c0792 Compare September 17, 2026 06:24
History records the offsets a task consumed and never the payloads, so the
server sends the bytes on the poll response: untagged for the task about to
run, and tagged with a WorkflowTaskCompleted event id when re-supplying what an
earlier task consumed. Core partitions the two and emits recorded ranges in
event order before the live one, so a replaying workflow observes them exactly
as it did the first time.

An empty range is delivered rather than dropped: a task where the subscription
saw nothing is a fact replay has to reproduce.

The Rust SDKs have no stream API, so they fail loudly on this job instead of
ignoring it. The server has already recorded the range as consumed and will not
send it again, so dropping it would lose data silently.
The command has a history event, so it fits the matching every SDK's replay
depends on: commands are popped from a queue as command-generated events
arrive, and one producing no event would put that out of step.

The machine never resolves. The event records the subscription and hands
nothing back; the ranges arrive later as their own activation jobs.
The reporter panics on any machine name it has no visualizer for, so
every state machine has to be listed there.
The range a task consumed is recorded on the completion that closes it, and
the machines look ahead to that event while replaying the task. An update cut
at a WFT started event left the completion on the retained tail or on an
unfetched page, so the range arrived one activation late.
A publish is matched on its stream and message count, a subscription on its
stream. Replay sends no commands, so a reissued command that differs from the
record would otherwise be accepted and the workflow's state would diverge
silently.
A completion that records consumed stream cursors ends its task sequence. Such
a task ran as one activation live, and folding it into a heartbeat chain would
hand several ranges over at once on replay.
Ranges for one task are handed over sorted by stream and offset, so a workflow
waiting on several streams sees the same order live and on replay. A
re-supplied slice has to cover exactly the recorded offsets, and a range that
observed nothing is rebuilt from the cursor rather than demanded from the
server.
The external stream family, developed alongside this one, takes 17 to 20 in
the activation job and 23 to 28 in the command. Both trees now encode the
native stream messages at 21, 29 and 30, so an artifact generated from either
decodes correctly against the other.
@moedash
moedash force-pushed the moe/AI-198-core-all branch from 02c0792 to 419b548 Compare September 19, 2026 00:51
…lary.

The api names an entry a record and carries its kind, producer, attempt and sequence on it, so the vendored stream protos follow the api head and lang appends with AppendStreamRecords and receives DeliverStreamRecords. Core passes each record through untouched.
History records only the offsets a task consumed, so a language replayer
that fetched the records from the stream service needs a way to hand them
to Core with the history. The replay worker puts them on its synthetic poll
response, so the ordinary delivery path and identity checks run unchanged.
The bytes for a recorded range only travel on the response that carries the
task, so a sticky task, or a sticky legacy query, handed to a worker that no
longer holds the run and fetched the history itself cannot be replayed. Failing
when the lookahead sees the range keeps the workflow off less input than it
had, and lets a legacy query go unanswered so the server retries it on the
normal queue, where the records travel with it.
A recorded range with content and no bytes for it is the worker's failure
to reconstruct the run, not the workflow's nondeterminism, so it is its own
error kind and treated like a failed history fetch: the task fails as an
unhandled worker failure and a legacy query goes unanswered, so the server
retries both where the records travel. The broken run is evicted even when
the query failure is withheld, so the retry starts from history.
Both sides carried the stream protos: the unified branch from the repairs
branch, this one from the lookahead branch's earlier cut. The seven files
are the api branch's diff on this Core's own api tree, and the protos
module keeps this branch's arms, which already cover the incoming ones.
A subscription's explicit start offset and an unnamed append's resolved
stream both went unchecked, so a replay that asked for something else
passed. The offset check skips a repeat subscribe, which the server
records at the cursor rather than at what the command asked for.
Both command enums are uninhabited, so neither arm runs. Nondeterminism
is the wrong label for an internal invariant when the rest of the series
works to keep worker failures out of that bucket.
A task owed two streams and sent one came out as nondeterminism, which a
worker configured for it fails the execution over. Both sides of that
comparison come from the server, so no arm of it is the workflow's fault.
The flag that says a completion consumed a range was written into the two
that carry across the whole scan, which only works while every path after
it returns. A local one says what it means on its own.
The fetch costs a page for any run whose page boundary lands there. A
subscription made through the stream service records no event, so there
is no telling a stream run from any other before the ranges arrive.
The blank comment line folded into the line above it, so the summary and
the paragraph under it ran together.
Reordering two arms of a match on distinct variants changes nothing and
shows up as a deletion on a branch whose claim is that it only adds.
Comparing it is only right for a run's first subscribe to a stream, and a
subscription made through the stream service leaves no event, so which one
is first cannot be told. Failing a sound run costs more than the drift.
The append event carries the exclusive end offset instead of a count, so
the three range-carrying messages read the same way. The command fields
say name rather than id, which is what a Workflow actually addresses.
The lang command now says stream_name for an append and stream_name_or_id
for a subscribe, so one vocabulary runs from lang through to History. The
append machine reads the batch size off the event's offset range.
The delivery tests build their own commands and histories, so they name
the fields directly. One recorded append passed a count where the event
now wants an end offset.
# Conflicts:
#	crates/protos/protos/local/temporal/sdk/core/workflow_activation/workflow_activation.proto
#	crates/sdk-core/src/worker/workflow/managed_run.rs
The main merge introduced the released headings above them and the resolution left the entries inside 0.8.0. The checkpoint requires additions under Unreleased.
The api adds a StreamStartPosition message and a start_position field on the subscribe command, at a new field number, so a Workflow can ask for the earliest record, the tail or the last N.
Lang can now ask for the earliest record, the tail or the last N instead of a negative offset. Core only forwards it: the server resolves it and records the offset, which is all replay matches against.
# Conflicts:
#	crates/sdk-core/src/protosext/mod.rs
#	crates/sdk-core/src/worker/workflow/history_update.rs
#	crates/sdk-core/src/worker/workflow/mod.rs
# Conflicts:
#	crates/protos/protos/local/temporal/sdk/core/workflow_activation/workflow_activation.proto
#	crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs
# Conflicts:
#	crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs
# Conflicts:
#	crates/sdk-core/src/protosext/mod.rs
#	crates/sdk-core/src/worker/workflow/history_update.rs
#	crates/sdk-core/src/worker/workflow/mod.rs
# Conflicts:
#	crates/sdk-core/src/worker/workflow/machines/transition_coverage.rs
#	crates/sdk-core/src/worker/workflow/machines/workflow_machines.rs
#	crates/sdk-core/src/worker/workflow/mod.rs
Carries #2 7006981 through #3 093c55e: the five channel requests take a common.v1.Execution and both linked_to fields become Execution. The native stream protos are untouched.
Carries #2 677d468 through #3 2a516e1: the execution and linked_to fields reuse the old numbers and nothing is reserved. The native stream protos are untouched.
@moedash

moedash commented Oct 3, 2026

Copy link
Copy Markdown
Owner Author

Replaced by #8, #9, #10, #11, #12, #13, #14, #15, #16, #17, #18, #19, #20, #21, #22 and #23.

Same content, split into 16 PRs in the v3 series: the notification channel first, then the streaming interface, then native streams and the rest. The branch stays as a pin.

@moedash moedash closed this Oct 3, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant