Skip to content

Commit 5958a23

Browse files
authored
Merge pull request #592 from obeli-sk/ignore-pause-status-of-delay-requests
refactor(db,grpc,webapi): Ignore `delay.paused` flag on advance requests
2 parents 0a779c1 + 7a0f8ed commit 5958a23

17 files changed

Lines changed: 133 additions & 14 deletions

assets/db/schema.json

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1216,6 +1216,10 @@
12161216
"schedule_at": {
12171217
"$ref": "#/$defs/HistoryEventScheduleAt"
12181218
},
1219+
"paused": {
1220+
"type": "boolean",
1221+
"default": false
1222+
},
12191223
"type": {
12201224
"type": "string",
12211225
"const": "delay_request"

crates/concepts/src/storage.rs

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -810,6 +810,8 @@ pub enum JoinSetRequest {
810810
delay_id: DelayId,
811811
expires_at: DateTime<Utc>,
812812
schedule_at: HistoryEventScheduleAt,
813+
#[serde(default)]
814+
paused: bool,
813815
},
814816
// Must be created by the executor in `PendingState::Locked`.
815817
#[display("ChildExecutionRequest({child_execution_id}, {target_ffqn}, params: {params})")]

crates/db-postgres/src/postgres_dao.rs

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -341,6 +341,7 @@ struct DelayReq {
341341
join_set_id: JoinSetId,
342342
delay_id: DelayId,
343343
expires_at: DateTime<Utc>,
344+
paused: bool,
344345
}
345346

346347
async fn fetch_created_event(
@@ -882,16 +883,18 @@ async fn bump_state_next_version(
882883
join_set_id,
883884
delay_id,
884885
expires_at,
886+
paused,
885887
}) = delay_req
886888
{
887889
debug!("Inserting delay to `t_delay`");
888890
tx.execute(
889-
"INSERT INTO t_delay (execution_id, join_set_id, delay_id, expires_at) VALUES ($1, $2, $3, $4)",
891+
"INSERT INTO t_delay (execution_id, join_set_id, delay_id, expires_at, is_paused) VALUES ($1, $2, $3, $4, $5)",
890892
&[
891893
&execution_id_str,
892894
&join_set_id.to_string(),
893895
&delay_id.to_string(),
894896
&expires_at,
897+
&paused,
895898
],
896899
)
897900
.await?;
@@ -2104,6 +2107,7 @@ async fn append(
21042107
JoinSetRequest::DelayRequest {
21052108
delay_id,
21062109
expires_at,
2110+
paused,
21072111
..
21082112
},
21092113
},
@@ -2117,6 +2121,7 @@ async fn append(
21172121
join_set_id: join_set_id.clone(),
21182122
delay_id: delay_id.clone(),
21192123
expires_at: *expires_at,
2124+
paused: *paused,
21202125
}),
21212126
)
21222127
.await?,

crates/db-sqlite/src/sqlite_dao.rs

Lines changed: 7 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -63,6 +63,7 @@ struct DelayReq {
6363
join_set_id: JoinSetId,
6464
delay_id: DelayId,
6565
expires_at: DateTime<Utc>,
66+
paused: bool,
6667
}
6768
/*
6869
mmap_size = 128MB - Set the global memory map so all processes can share some data
@@ -1326,19 +1327,21 @@ impl SqlitePool {
13261327
join_set_id,
13271328
delay_id,
13281329
expires_at,
1330+
paused,
13291331
}) = delay_req
13301332
{
13311333
debug!("Inserting delay to `t_delay`");
13321334
let mut stmt = tx.prepare_cached(
1333-
"INSERT INTO t_delay (execution_id, join_set_id, delay_id, expires_at) \
1335+
"INSERT INTO t_delay (execution_id, join_set_id, delay_id, expires_at, is_paused) \
13341336
VALUES \
1335-
(:execution_id, :join_set_id, :delay_id, :expires_at)",
1337+
(:execution_id, :join_set_id, :delay_id, :expires_at, :is_paused)",
13361338
)?;
13371339
stmt.execute(named_params! {
13381340
":execution_id": execution_id_str,
13391341
":join_set_id": join_set_id.to_string(),
13401342
":delay_id": delay_id.to_string(),
13411343
":expires_at": expires_at,
1344+
":is_paused": paused,
13421345
})?;
13431346
}
13441347
Ok(appending_version.increment())
@@ -2120,6 +2123,7 @@ impl SqlitePool {
21202123
JoinSetRequest::DelayRequest {
21212124
delay_id,
21222125
expires_at,
2126+
paused,
21232127
..
21242128
},
21252129
},
@@ -2133,6 +2137,7 @@ impl SqlitePool {
21332137
join_set_id: join_set_id.clone(),
21342138
delay_id: delay_id.clone(),
21352139
expires_at: *expires_at,
2140+
paused: *paused,
21362141
}),
21372142
)?,
21382143
AppendNotifier::default(),

crates/grpc/src/grpc_mapping.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1182,6 +1182,7 @@ fn history_event_from_grpc(
11821182
schedule_at: delay_schedule_at_from_grpc(
11831183
delay.scheduled_at.argument_must_exist("scheduled_at")?,
11841184
)?,
1185+
paused: delay.paused,
11851186
}
11861187
}
11871188
history_event::join_set_request::JoinSetRequest::ChildExecutionRequest(child) => {
@@ -1386,12 +1387,14 @@ pub fn history_event_to_grpc(event: HistoryEvent) -> grpc_gen::execution_event::
13861387
delay_id,
13871388
expires_at,
13881389
schedule_at,
1390+
paused,
13891391
} => Some(
13901392
history_event::join_set_request::JoinSetRequest::DelayRequest(
13911393
history_event::join_set_request::DelayRequest {
13921394
delay_id: Some(delay_id.into()),
13931395
expires_at: Some(prost_wkt_types::Timestamp::from(expires_at)),
13941396
scheduled_at: Some(delay_schedule_at_to_grpc(schedule_at)),
1397+
paused,
13951398
},
13961399
),
13971400
),

crates/testing/db-tests/tests/lifecycle.rs

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1314,6 +1314,7 @@ async fn test_lock_inner(db_connection: &dyn DbConnection, sim_clock: SimClock)
13141314
delay_id: DelayId::new(&execution_id, &join_set_id),
13151315
expires_at: sim_clock.now(),
13161316
schedule_at: HistoryEventScheduleAt::Now,
1317+
paused: false,
13171318
},
13181319
join_set_id,
13191320
},
@@ -1494,6 +1495,7 @@ async fn get_expired_delay(db_connection: &dyn DbConnection, sim_clock: SimClock
14941495
delay_id: delay_id.clone(),
14951496
expires_at: sim_clock.now() + lock_expiry,
14961497
schedule_at: HistoryEventScheduleAt::In(lock_expiry),
1498+
paused: false,
14971499
},
14981500
},
14991501
},
@@ -1702,6 +1704,7 @@ async fn append_same_delay_id_twice_should_fail(
17021704
delay_id: delay_id.clone(),
17031705
expires_at: sim_clock.now() + lock_expiry,
17041706
schedule_at: HistoryEventScheduleAt::In(lock_expiry),
1707+
paused: false,
17051708
},
17061709
},
17071710
},
@@ -1723,6 +1726,7 @@ async fn append_same_delay_id_twice_should_fail(
17231726
delay_id: delay_id.clone(),
17241727
expires_at: sim_clock.now() + lock_expiry,
17251728
schedule_at: HistoryEventScheduleAt::In(lock_expiry),
1729+
paused: false,
17261730
},
17271731
},
17281732
},
@@ -1774,6 +1778,7 @@ async fn test_append_response_with_same_delay_id_twice_should_fail(database: Dat
17741778
delay_id: delay_id.clone(),
17751779
expires_at: sim_clock.now(),
17761780
schedule_at: HistoryEventScheduleAt::Now,
1781+
paused: false,
17771782
};
17781783
let response = JoinSetResponse::DelayFinished {
17791784
delay_id: delay_id.clone(),
@@ -1982,6 +1987,7 @@ async fn delay_cancellation_should_be_idempotent(database: Database) {
19821987
delay_id: delay_id.clone(),
19831988
expires_at: sim_clock.now() + lock_expiry,
19841989
schedule_at: HistoryEventScheduleAt::In(lock_expiry),
1990+
paused: false,
19851991
},
19861992
},
19871993
},
@@ -2692,6 +2698,7 @@ async fn pause_with_pending_delay_then_response_then_unpause_should_be_pending(d
26922698
delay_id: delay_id.clone(),
26932699
expires_at: sim_clock.now() + delay_duration,
26942700
schedule_at: HistoryEventScheduleAt::In(delay_duration),
2701+
paused: false,
26952702
},
26962703
},
26972704
},
@@ -3143,6 +3150,7 @@ async fn test_list_responses(database: Database) {
31433150
delay_id: delay_id.clone(),
31443151
expires_at: delay_expires_at,
31453152
schedule_at: HistoryEventScheduleAt::Now,
3153+
paused: false,
31463154
},
31473155
},
31483156
},
@@ -3268,6 +3276,7 @@ async fn test_list_responses_pagination_direction(database: Database) {
32683276
delay_id: delay_id.clone(),
32693277
expires_at: delay_expires_at,
32703278
schedule_at: HistoryEventScheduleAt::Now,
3279+
paused: false,
32713280
},
32723281
},
32733282
},
@@ -3451,6 +3460,7 @@ async fn test_list_execution_events_pagination_direction(database: Database) {
34513460
delay_id,
34523461
expires_at: delay_expires_at,
34533462
schedule_at: HistoryEventScheduleAt::Now,
3463+
paused: false,
34543464
},
34553465
},
34563466
},

crates/wasm-workers/src/workflow/event_history.rs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -794,6 +794,7 @@ impl EventHistory {
794794
delay_id: found_delay_id,
795795
expires_at,
796796
schedule_at: found_schedule_at,
797+
..
797798
},
798799
},
799800
) if *join_set_id == *found_join_set_id
@@ -1300,6 +1301,7 @@ impl EventHistory {
13001301
delay_id,
13011302
expires_at: expires_at_if_new,
13021303
schedule_at,
1304+
paused: false,
13031305
},
13041306
};
13051307
let history_event = (event.clone(), db_connection.version().clone());
@@ -1779,6 +1781,7 @@ impl EventHistory {
17791781
delay_id,
17801782
expires_at: expires_at_if_new,
17811783
schedule_at,
1784+
paused: false,
17821785
},
17831786
};
17841787
version = version.increment();

crates/wasm-workers/src/workflow/replay_advance.rs

Lines changed: 75 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -298,10 +298,12 @@ fn normalize_join_set_request_for_matching(request: JoinSetRequest) -> JoinSetRe
298298
delay_id,
299299
expires_at: _,
300300
schedule_at,
301+
paused: _, // Ignore for comparison, user's flag will make it to the database in `merge_requested_overrides_into_fresh_write`
301302
} => JoinSetRequest::DelayRequest {
302303
delay_id,
303304
expires_at: DateTime::UNIX_EPOCH,
304305
schedule_at: normalize_schedule_at_for_matching(schedule_at),
306+
paused: false,
305307
},
306308
JoinSetRequest::ChildExecutionRequest {
307309
child_execution_id,
@@ -346,13 +348,15 @@ pub(crate) fn merge_requested_overrides_into_fresh_write(
346348
(
347349
CapturedDbWrite::AppendBatchCreateNewExecution {
348350
child_req: requested_child_req,
351+
batch: requested_batch,
349352
..
350353
},
351354
CapturedDbWrite::AppendBatchCreateNewExecution { child_req: _, .. },
352355
) => {
353356
let mut merged = fresh.clone();
354357
let CapturedDbWrite::AppendBatchCreateNewExecution {
355358
child_req: merged_child_req,
359+
batch: merged_batch,
356360
..
357361
} = &mut merged.write
358362
else {
@@ -361,8 +365,79 @@ pub(crate) fn merge_requested_overrides_into_fresh_write(
361365
for (requested, fresh) in requested_child_req.iter().zip(merged_child_req.iter_mut()) {
362366
fresh.paused = requested.paused; // Allow users to specify the paused behavior.
363367
}
368+
merge_delay_paused_flags(requested_batch, merged_batch);
369+
merged
370+
}
371+
(
372+
CapturedDbWrite::Append {
373+
req: requested_req, ..
374+
},
375+
CapturedDbWrite::Append { .. },
376+
) => {
377+
let mut merged = fresh.clone();
378+
let CapturedDbWrite::Append {
379+
req: merged_req, ..
380+
} = &mut merged.write
381+
else {
382+
unreachable!("matched variant must stay matched")
383+
};
384+
merge_delay_paused_flag(requested_req, merged_req);
385+
merged
386+
}
387+
(
388+
CapturedDbWrite::AppendBatch {
389+
batch: requested_batch,
390+
..
391+
},
392+
CapturedDbWrite::AppendBatch { .. },
393+
) => {
394+
let mut merged = fresh.clone();
395+
let CapturedDbWrite::AppendBatch {
396+
batch: merged_batch,
397+
..
398+
} = &mut merged.write
399+
else {
400+
unreachable!("matched variant must stay matched")
401+
};
402+
merge_delay_paused_flags(requested_batch, merged_batch);
364403
merged
365404
}
366405
_ => fresh.clone(),
367406
}
368407
}
408+
409+
fn merge_delay_paused_flags(requested: &[AppendRequest], merged: &mut [AppendRequest]) {
410+
for (requested, merged) in requested.iter().zip(merged.iter_mut()) {
411+
merge_delay_paused_flag(requested, merged);
412+
}
413+
}
414+
415+
fn merge_delay_paused_flag(requested: &AppendRequest, merged: &mut AppendRequest) {
416+
if let (
417+
ExecutionRequest::HistoryEvent {
418+
event:
419+
HistoryEvent::JoinSetRequest {
420+
request:
421+
JoinSetRequest::DelayRequest {
422+
paused: requested_paused,
423+
..
424+
},
425+
..
426+
},
427+
},
428+
ExecutionRequest::HistoryEvent {
429+
event:
430+
HistoryEvent::JoinSetRequest {
431+
request:
432+
JoinSetRequest::DelayRequest {
433+
paused: merged_paused,
434+
..
435+
},
436+
..
437+
},
438+
},
439+
) = (&requested.event, &mut merged.event)
440+
{
441+
*merged_paused = *requested_paused;
442+
}
443+
}

crates/wasm-workers/src/workflow/snapshots/execute_workflow_fn_with_delays@execute_workflow_fn_with_single_delay-testing_stub-workflow__workflow.join-next-in-scope.snap

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -100,7 +100,8 @@ expression: "ExecutionLogSanitized::from(execution_log.clone())"
100100
"secs": 86400,
101101
"nanos": 0
102102
}
103-
}
103+
},
104+
"paused": false
104105
}
105106
}
106107
}
@@ -248,7 +249,8 @@ expression: "ExecutionLogSanitized::from(execution_log.clone())"
248249
"secs": 86400,
249250
"nanos": 0
250251
}
251-
}
252+
},
253+
"paused": false
252254
}
253255
}
254256
}

crates/wasm-workers/src/workflow/snapshots/execute_workflow_fn_with_delays@join_next_produces_all_processed_error-testing_sleep-workflow__workflow.join-next-produces-all-processed-error.snap

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -77,7 +77,8 @@ expression: "ExecutionLogSanitized::from(execution_log.clone())"
7777
"secs": 0,
7878
"nanos": 10000000
7979
}
80-
}
80+
},
81+
"paused": false
8182
}
8283
}
8384
}

0 commit comments

Comments
 (0)