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
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@ import "temporal/sdk/core/nexus/nexus.proto";
// * Signal and update handlers should be invoked before workflow routines are iterated. That is to
// say before the users' main workflow function and anything spawned by it is allowed to continue.
// * Channel notifications are input from outside the workflow, like signals, so they go with
// them.
// them and ahead of the stream ranges among the other jobs.
// * Local activities resolutions go after other normal jobs because while *not* replaying, they
// will always take longer than anything else that produces an immediate job (which is
// effectively instant). When *replaying* we need to scan ahead for LA markers so that we can
Expand Down
11 changes: 11 additions & 0 deletions crates/sdk-core/CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,17 @@ relevant information.
command and end a subscription with `UnsubscribeNotificationChannel`. The notifications the
server folds for a Workflow Task arrive as one `NotificationsReceived` activation job, taken from
the task's scheduled event, so replay yields the same job at the same point.
* Workflows can subscribe to server-side streams and append batches of records to them with the
`SubscribeStream` and `AppendStreamRecords` commands. Consumed ranges reach the workflow as
`DeliverStreamRecords` activation jobs, and replay hands each recorded range back in the
activation of the task that consumed it.
* A history fed to a replay worker can carry the stream records its tasks consumed
(`HistoryForReplay::with_stream_slices`), so a language replayer that fetched them from the
stream service can replay a consuming workflow. History alone holds only the offsets.
* A task whose history records a consumed range with content that the response carried no
records for fails before the workflow runs, rather than after it ran on less input. A legacy
query dispatched that way to a worker that no longer holds the run goes unanswered, so the
server retries it on the normal task queue, where the records travel with it.

### Fixed
* Task-poll targets no longer decrease after cancelled or timed-out polls. Affected pollers still
Expand Down
33 changes: 31 additions & 2 deletions crates/sdk-core/src/core_tests/channels.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,8 +2,8 @@ use crate::{
init_replay_worker,
replay::{HistoryFeeder, HistoryForReplay, ReplayWorkerInput, TestHistoryBuilder},
test_help::{
MockPollCfg, ResponseType, WorkerTestHelpers, build_mock_pollers, mock_worker,
test_worker_cfg,
MockPollCfg, PollWFTRespExt, ResponseType, WorkerTestHelpers, build_mock_pollers,
hist_to_poll_resp, mock_worker, test_worker_cfg,
},
worker::client::mocks::mock_worker_client,
};
Expand Down Expand Up @@ -351,6 +351,7 @@ fn job_kinds(task: &WorkflowActivation) -> Vec<&'static str> {
.map(|j| match j.variant.as_ref().unwrap() {
workflow_activation_job::Variant::InitializeWorkflow(_) => "init",
workflow_activation_job::Variant::NotificationsReceived(_) => "notifications",
workflow_activation_job::Variant::DeliverStreamRecords(_) => "stream",
_ => "other",
})
.collect()
Expand All @@ -368,6 +369,34 @@ fn one_task_with_notifications() -> TestHistoryBuilder {
t
}

/// The notifications on the scheduled event reach the activation of that task,
/// as one job, ahead of the stream ranges the same task was handed.
#[tokio::test]
async fn notifications_on_the_scheduled_event_reach_the_live_activation() {
let t = one_task_with_notifications();
let mut poll_resp = hist_to_poll_resp(&t, "wfid".to_owned(), ResponseType::AllHistory);
poll_resp.add_stream_slice("s1", 0, 0, &["alpha"]);

let mock = MockPollCfg::from_resp_batches(
"wfid",
t,
[ResponseType::Raw(poll_resp.resp)],
mock_worker_client(),
);
let core = mock_worker(build_mock_pollers(mock));

let task = core.poll_workflow_activation().await.unwrap();
assert!(!task.is_replaying);
assert_eq!(job_kinds(&task), vec!["init", "notifications", "stream"]);
assert_eq!(
received(&task),
vec![vec![notification("orders", 3), notification("invoices", 7)]]
);
core.complete_workflow_activation(WorkflowActivationCompletion::empty(task.run_id))
.await
.unwrap();
}

/// Replaying the same history yields the same job. Nothing is re-supplied: the
/// notifications are on the event, so History alone carries them.
#[tokio::test]
Expand Down
Loading
Loading