1use super::registration::CreditRegistration;
5use super::state::{
6 ADMIN_ONLY_TARGETS, CreditRegistrationErrorCode, CreditRegistrationState,
7 PendingSupersessionEffect,
8};
9use crate::credit_registration_events::{
10 CreditRegistrationEventKind, NewCreditRegistrationEvent, SuotarAnswer,
11};
12use crate::error::missing_model_error;
13use crate::prelude::*;
14use crate::suotar_api_calls::SuotarEndpoint;
15use chrono::TimeDelta;
16use std::collections::HashMap;
17
18#[derive(Debug, Clone, Copy, PartialEq, Eq)]
20pub enum AdminAttention {
21 Raise,
22 Clear,
23}
24
25impl AdminAttention {
26 pub fn is_raised(self) -> bool {
28 self == Self::Raise
29 }
30}
31
32#[derive(Debug, Clone, PartialEq)]
33pub struct Transition {
34 pub to_state: CreditRegistrationState,
35 pub error_code: Option<CreditRegistrationErrorCode>,
36 pub error_message: Option<String>,
38 pub needs_admin_attention: Option<AdminAttention>,
40 pub event_kind: CreditRegistrationEventKind,
41 pub event_message: Option<String>,
42 pub actor_user_id: Option<Uuid>,
43 pub suotar_api_call_id: Option<Uuid>,
44 pub suotar_endpoint: Option<SuotarEndpoint>,
45 pub suotar_requested_at: Option<DateTime<Utc>>,
46 pub suotar_answered_at: Option<DateTime<Utc>>,
47 pub suotar_answer: Option<SuotarAnswer>,
48 pub event_details: Option<serde_json::Value>,
50 pub request_item_id: Option<String>,
52 pub expected_from_state: Option<CreditRegistrationState>,
57 pub policy: TransitionPolicy,
59 pub next_attempt_at: Option<DateTime<Utc>>,
62 pub keeps_enrolment_checked_at: bool,
65}
66
67#[derive(Debug, Clone, Copy, PartialEq, Eq)]
69pub enum TransitionPolicy {
70 Pipeline,
73 Admin,
78 Planted,
81}
82
83impl TransitionPolicy {
84 fn allows(self, from: CreditRegistrationState, to: CreditRegistrationState) -> bool {
85 match self {
86 Self::Pipeline => from.allowed_targets().contains(&to),
87 Self::Admin => from.allowed_targets().contains(&to) || ADMIN_ONLY_TARGETS.contains(&to),
88 Self::Planted => true,
89 }
90 }
91}
92
93impl Transition {
94 pub fn to(to_state: CreditRegistrationState) -> Self {
95 Self {
96 to_state,
97 error_code: None,
98 error_message: None,
99 needs_admin_attention: None,
100 event_kind: CreditRegistrationEventKind::StateChanged,
101 event_message: None,
102 actor_user_id: None,
103 suotar_api_call_id: None,
104 suotar_endpoint: None,
105 suotar_requested_at: None,
106 suotar_answered_at: None,
107 suotar_answer: None,
108 event_details: None,
109 request_item_id: None,
110 expected_from_state: None,
111 policy: TransitionPolicy::Pipeline,
112 next_attempt_at: None,
113 keeps_enrolment_checked_at: false,
114 }
115 }
116
117 pub fn by_hand(to_state: CreditRegistrationState) -> Self {
119 Self {
120 policy: TransitionPolicy::Admin,
121 ..Self::to(to_state)
122 }
123 }
124
125 pub fn planted(to_state: CreditRegistrationState) -> Self {
127 Self {
128 policy: TransitionPolicy::Planted,
129 ..Self::to(to_state)
130 }
131 }
132}
133
134pub async fn transition(
151 conn: &mut PgConnection,
152 id: Uuid,
153 transition: &Transition,
154) -> ModelResult<CreditRegistration> {
155 match transition_unless_moved_on(conn, id, transition).await? {
156 Transitioned::Written(after) => Ok(*after),
157 Transitioned::MovedOn { found } => Err(model_err!(
158 PreconditionFailed,
159 format!(
160 "Credit registration {id} is in {found:?}, not the expected {}: refusing to overwrite it.",
161 transition
162 .expected_from_state
163 .map(|expected| format!("{expected:?}"))
164 .unwrap_or_default()
165 )
166 )),
167 }
168}
169
170#[derive(Debug, Clone)]
172pub enum Transitioned {
173 Written(Box<CreditRegistration>),
175 MovedOn { found: CreditRegistrationState },
177}
178
179pub async fn transition_unless_moved_on(
183 conn: &mut PgConnection,
184 id: Uuid,
185 transition: &Transition,
186) -> ModelResult<Transitioned> {
187 let mut tx = conn.begin().await?;
188 let from_state = lock_for_moves(&mut tx, &[id])
189 .await?
190 .remove(&id)
191 .ok_or_else(missing_model_error(
192 ModelErrorType::RecordNotFound,
193 format!("Credit registration {id} does not exist."),
194 ))?;
195 if transition
196 .expected_from_state
197 .is_some_and(|expected| from_state != expected)
198 {
199 return Ok(Transitioned::MovedOn { found: from_state });
200 }
201 check_edge(id, from_state, transition.to_state, transition.policy)?;
202 let after = write_moves(&mut tx, &[(id, from_state, transition)])
203 .await?
204 .pop()
205 .ok_or_else(missing_model_error(
206 ModelErrorType::RecordNotFound,
207 format!("Credit registration {id} does not exist."),
208 ))?;
209 tx.commit().await?;
210 Ok(Transitioned::Written(Box::new(after)))
211}
212
213async fn lock_for_moves(
215 conn: &mut PgConnection,
216 ids: &[Uuid],
217) -> ModelResult<HashMap<Uuid, CreditRegistrationState>> {
218 let locked = sqlx::query!(
219 r#"
220SELECT id,
221 state
222FROM credit_registrations
223WHERE id = ANY($1)
224 AND deleted_at IS NULL
225ORDER BY id FOR
226UPDATE
227 "#,
228 ids
229 )
230 .fetch_all(conn)
231 .await?;
232 Ok(locked.into_iter().map(|row| (row.id, row.state)).collect())
233}
234
235async fn write_moves(
237 conn: &mut PgConnection,
238 moves: &[(Uuid, CreditRegistrationState, &Transition)],
239) -> ModelResult<Vec<CreditRegistration>> {
240 let events: Vec<NewCreditRegistrationEvent> = moves
241 .iter()
242 .map(|(id, from_state, transition)| NewCreditRegistrationEvent {
243 credit_registration_id: *id,
244 kind: transition.event_kind,
245 from_state: Some(*from_state),
246 to_state: Some(transition.to_state),
247 error_code: transition.error_code,
248 message: transition.event_message.clone(),
249 suotar_api_call_id: transition.suotar_api_call_id,
250 suotar_endpoint: transition.suotar_endpoint,
251 suotar_requested_at: transition.suotar_requested_at,
252 suotar_answered_at: transition.suotar_answered_at,
253 suotar_answer: transition.suotar_answer,
254 actor_user_id: transition.actor_user_id,
255 details: transition.event_details.clone(),
256 request_item_id: transition.request_item_id.clone(),
257 })
258 .collect();
259 let ids: Vec<Uuid> = moves.iter().map(|(id, _, _)| *id).collect();
260 let to_states: Vec<CreditRegistrationState> = moves
261 .iter()
262 .map(|(_, _, transition)| transition.to_state)
263 .collect();
264 let error_codes: Vec<Option<CreditRegistrationErrorCode>> = moves
265 .iter()
266 .map(|(_, _, transition)| transition.error_code)
267 .collect();
268 let error_messages: Vec<Option<String>> = moves
269 .iter()
270 .map(|(_, _, transition)| transition.error_message.clone())
271 .collect();
272 let needs_admin: Vec<Option<bool>> = moves
273 .iter()
274 .map(|(_, _, transition)| {
275 transition
276 .needs_admin_attention
277 .map(AdminAttention::is_raised)
278 })
279 .collect();
280 let terminal: Vec<bool> = to_states.iter().map(|state| state.is_terminal()).collect();
281 let failure: Vec<bool> = to_states
282 .iter()
283 .map(|state| state.is_failed_state())
284 .collect();
285 let next_attempts: Vec<Option<DateTime<Utc>>> = moves
286 .iter()
287 .map(|(_, _, transition)| transition.next_attempt_at)
288 .collect();
289 let default_delays: Vec<TimeDelta> = to_states
290 .iter()
291 .map(|state| state.default_attempt_delay())
292 .collect();
293 let keeps_checked_at: Vec<bool> = moves
294 .iter()
295 .map(|(_, _, transition)| transition.keeps_enrolment_checked_at)
296 .collect();
297 let keeps_schedule: Vec<bool> = to_states
298 .iter()
299 .map(|state| state.keeps_enrolment_check_schedule())
300 .collect();
301 let keeps_waiting_since: Vec<bool> = to_states
302 .iter()
303 .map(|state| {
304 state.keeps_enrolment_check_schedule()
305 && *state != CreditRegistrationState::NoUsableEnrolment
306 })
307 .collect();
308 let written = sqlx::query_as!(
309 CreditRegistration,
310 r#"
311UPDATE credit_registrations cr
312SET state = move.to_state,
313 -- clock_timestamp(), not now(): now() is the transaction timestamp, so several state changes in
314 -- one transaction would share an instant and the timeline would lose their order.
315 state_entered_at = clock_timestamp(),
316 error_code = move.error_code,
317 error_message = move.error_message,
318 needs_admin_attention = COALESCE(move.needs_admin_attention, cr.needs_admin_attention),
319 -- ELSE NULL: without it an admin retry stays invisible to every terminal_at IS NULL query.
320 terminal_at = CASE
321 WHEN move.terminal THEN COALESCE(cr.terminal_at, now())
322 ELSE NULL
323 END,
324 first_failed_at = CASE
325 WHEN move.failure THEN COALESCE(cr.first_failed_at, now())
326 WHEN move.to_state = 'no_usable_enrolment' THEN NULL
327 ELSE cr.first_failed_at
328 END,
329 submit_retry_count = CASE
330 WHEN move.to_state = 'no_usable_enrolment' THEN 0
331 ELSE cr.submit_retry_count
332 END,
333 registered_at = CASE
334 WHEN move.to_state = 'registered' THEN COALESCE(cr.registered_at, now())
335 ELSE cr.registered_at
336 END,
337 submitted_at = CASE
338 WHEN move.to_state = 'submitting' THEN now()
339 ELSE cr.submitted_at
340 END,
341 enrolment_checked_at = CASE
342 WHEN move.keeps_checked_at THEN cr.enrolment_checked_at
343 WHEN (
344 cr.state = 'resolving_enrolment'
345 OR cr.enrolment_check_claimed_until IS NOT NULL
346 )
347 AND move.to_state IN ('checking_enrolment', 'no_usable_enrolment') THEN now()
348 WHEN cr.state = 'checking_enrolment'
349 AND move.to_state <> 'checking_enrolment' THEN now()
350 ELSE cr.enrolment_checked_at
351 END,
352 -- Only on starting to wait, not on every check that finds no enrolment again.
353 enrolment_banner_dismissed_at = CASE
354 WHEN move.to_state = 'no_usable_enrolment'
355 AND cr.no_usable_enrolment_since IS NULL THEN NULL
356 ELSE cr.enrolment_banner_dismissed_at
357 END,
358 -- A retried lookup passes through the states that keep it on its way back to no_usable_enrolment.
359 no_usable_enrolment_since = CASE
360 WHEN move.to_state = 'no_usable_enrolment' THEN COALESCE(cr.no_usable_enrolment_since, now())
361 WHEN move.keeps_waiting_since THEN cr.no_usable_enrolment_since
362 ELSE NULL
363 END,
364 enrolment_check_anchor_at = CASE
365 WHEN move.keeps_schedule THEN cr.enrolment_check_anchor_at
366 END,
367 enrolment_check_step = CASE
368 WHEN move.keeps_schedule THEN cr.enrolment_check_step
369 END,
370 enrolment_check_due_at = CASE
371 WHEN move.keeps_schedule THEN cr.enrolment_check_due_at
372 END,
373 is_enrolment_check_batched = move.keeps_schedule
374 AND cr.is_enrolment_check_batched,
375 enrolment_checks_stopped_at = CASE
376 WHEN move.keeps_schedule THEN cr.enrolment_checks_stopped_at
377 END,
378 enrolment_check_source = CASE
379 WHEN move.keeps_schedule THEN cr.enrolment_check_source
380 ELSE 'schedule'
381 END,
382 enrolment_check_claimed_until = NULL,
383 next_attempt_at = COALESCE(
384 move.next_attempt_at,
385 now() + move.default_delay
386 )
387FROM UNNEST(
388 $1::uuid [],
389 $2::credit_registration_state [],
390 $3::credit_registration_error_code [],
391 $4::text [],
392 $5::boolean [],
393 $6::boolean [],
394 $7::boolean [],
395 $8::timestamptz [],
396 $9::interval [],
397 $10::boolean [],
398 $11::boolean [],
399 $12::boolean []
400 ) AS move(
401 id,
402 to_state,
403 error_code,
404 error_message,
405 needs_admin_attention,
406 terminal,
407 failure,
408 next_attempt_at,
409 default_delay,
410 keeps_checked_at,
411 keeps_schedule,
412 keeps_waiting_since
413 )
414WHERE cr.id = move.id
415 AND cr.deleted_at IS NULL
416RETURNING cr.*
417 "#,
418 &ids,
419 &to_states as &[CreditRegistrationState],
420 &error_codes as &[Option<CreditRegistrationErrorCode>],
421 &error_messages as &[Option<String>],
422 &needs_admin as &[Option<bool>],
423 &terminal,
424 &failure,
425 &next_attempts as &[Option<DateTime<Utc>>],
426 &default_delays as &[TimeDelta],
427 &keeps_checked_at,
428 &keeps_schedule,
429 &keeps_waiting_since,
430 )
431 .fetch_all(&mut *conn)
432 .await?;
433
434 let settled: Vec<_> = ids.iter().copied().zip(to_states.iter().copied()).collect();
435 settle_pending_supersessions(&mut *conn, &settled).await?;
436 crate::credit_registration_events::insert_batch(conn, &events).await?;
437 Ok(written)
438}
439
440async fn settle_pending_supersessions(
443 conn: &mut PgConnection,
444 moves: &[(Uuid, CreditRegistrationState)],
445) -> ModelResult<()> {
446 let mut completed = Vec::new();
447 let mut abandoned = Vec::new();
448 for &(id, to_state) in moves {
449 match to_state.pending_supersession_effect() {
450 PendingSupersessionEffect::Keep => {}
451 PendingSupersessionEffect::Complete => completed.push(id),
452 PendingSupersessionEffect::Abandon => abandoned.push(id),
453 }
454 }
455 if completed.is_empty() && abandoned.is_empty() {
456 return Ok(());
457 }
458 sqlx::query!(
459 r#"
460UPDATE credit_registrations
461SET superseded_by_id = CASE
462 WHEN pending_superseded_by_id = ANY($1::uuid []) THEN pending_superseded_by_id
463 ELSE superseded_by_id
464 END,
465 superseded_at = CASE
466 WHEN pending_superseded_by_id = ANY($1::uuid []) THEN now()
467 ELSE superseded_at
468 END,
469 pending_superseded_by_id = NULL
470WHERE (
471 pending_superseded_by_id = ANY($1::uuid [])
472 OR pending_superseded_by_id = ANY($2::uuid [])
473 )
474 AND deleted_at IS NULL
475 "#,
476 &completed,
477 &abandoned,
478 )
479 .execute(conn)
480 .await?;
481 Ok(())
482}
483
484fn check_edge(
487 id: Uuid,
488 from: CreditRegistrationState,
489 to: CreditRegistrationState,
490 policy: TransitionPolicy,
491) -> ModelResult<()> {
492 if from == to || policy.allows(from, to) {
494 return Ok(());
495 }
496 Err(model_err!(
497 InvalidRequest,
498 format!("Credit registration {id} may not move from {from:?} to {to:?} under {policy:?}.")
499 ))
500}
501
502#[derive(Debug, Clone, PartialEq)]
504pub struct BatchMove {
505 pub id: Uuid,
506 pub transition: Transition,
507}
508
509pub async fn transition_batch(conn: &mut PgConnection, moves: &[BatchMove]) -> ModelResult<i64> {
516 if moves.is_empty() {
517 return Ok(0);
518 }
519 let mut tx = conn.begin().await?;
520 let ids: Vec<Uuid> = moves.iter().map(|batch_move| batch_move.id).collect();
521 let locked = lock_for_moves(&mut tx, &ids).await?;
522 let mut writes = Vec::new();
523 for batch_move in moves {
524 let Some(&from) = locked.get(&batch_move.id) else {
525 continue;
526 };
527 if batch_move
528 .transition
529 .expected_from_state
530 .is_some_and(|expected| expected != from)
531 {
532 continue;
533 }
534 check_edge(
535 batch_move.id,
536 from,
537 batch_move.transition.to_state,
538 batch_move.transition.policy,
539 )?;
540 writes.push((batch_move.id, from, &batch_move.transition));
541 }
542 if !writes.is_empty() {
543 write_moves(&mut tx, &writes).await?;
544 }
545 tx.commit().await?;
546 Ok(i64::try_from(writes.len()).unwrap_or(i64::MAX))
547}