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
72 changes: 32 additions & 40 deletions dev-tools/omdb/src/bin/omdb/db/saga.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ use internal_dns_resolver::ResolveError;
use internal_dns_types::names::ServiceName;
use nexus_db_lookup::DataStoreConnection;
use nexus_db_model::Saga;
use nexus_db_model::SagaExecState;
use nexus_db_model::SagaNodeEvent;
use nexus_db_model::SagaState;
use nexus_db_model::SecId;
Expand Down Expand Up @@ -185,7 +186,6 @@ impl From<Saga> for SagaRow {
current_sec,
adopt_generation: _,
adopt_time: _,
abandon_metadata: _,
} = saga;
Self {
id: id.0.into(),
Expand All @@ -196,7 +196,7 @@ impl From<Saga> for SagaRow {
},
time_created,
name,
state: format!("{saga_state:?}"),
state: format!("{:?}", SagaState::from(saga_state)),
}
}
}
Expand Down Expand Up @@ -240,12 +240,11 @@ You should only do this if:
if !args.bypass_sec_check {
let saga: Saga = {
use nexus_db_schema::schema::saga::dsl;
let row = dsl::saga
dsl::saga
.filter(dsl::id.eq(args.saga_id))
.select(nexus_db_model::SagaRow::as_select())
.first_async(&*conn)
.await?;
Saga::try_from(row)?
.select(Saga::as_select())
.first_async::<Saga>(&*conn)
.await?
};

let status = get_saga_sec_status(omdb, opctx, &saga).await;
Expand Down Expand Up @@ -377,23 +376,20 @@ async fn cmd_sagas_abandon(
let should_print_color =
should_colorize(omdb.output.color, supports_color::Stream::Stdout);
let conn = datastore.pool_connection_for_tests().await?;
let saga: Saga = {
let row = dsl::saga
.filter(dsl::id.eq(args.saga_id))
.select(nexus_db_model::SagaRow::as_select())
.first_async(&*conn)
.await?;
Saga::try_from(row)?
};
let saga: Saga = dsl::saga
.filter(dsl::id.eq(args.saga_id))
.select(Saga::as_select())
.first_async::<Saga>(&*conn)
.await?;

match saga.saga_state {
SagaState::Done => {
SagaExecState::Done => {
bail!("saga {} is already done executing", args.saga_id);
}
SagaState::Abandoned => {
SagaExecState::Abandoned(_) => {
bail!("saga {} is already abandoned", args.saga_id);
}
SagaState::Running | SagaState::Unwinding => {}
SagaExecState::Running | SagaExecState::Unwinding => {}
}

let text = r#"
Expand Down Expand Up @@ -424,14 +420,11 @@ execute even if it is abandoned. You should only proceed if:
// Before doing anything: find the current SEC for the saga, and ping it to
// ensure that the Nexus is down.
if !args.bypass_sec_check {
let saga = {
let row = dsl::saga
.filter(dsl::id.eq(args.saga_id))
.select(nexus_db_model::SagaRow::as_select())
.first_async(&*conn)
.await?;
Saga::try_from(row)?
};
let saga = dsl::saga
.filter(dsl::id.eq(args.saga_id))
.select(Saga::as_select())
.first_async::<Saga>(&*conn)
.await?;

let status = get_saga_sec_status(omdb, opctx, &saga).await;
status.display_message(should_print_color);
Expand Down Expand Up @@ -481,22 +474,22 @@ async fn get_all_sagas_in_state(
let records_batch =
paginated(dsl::saga, dsl::id, &p.current_pagparams())
.filter(dsl::saga_state.eq(state))
.select(nexus_db_model::SagaRow::as_select())
.load_async::<nexus_db_model::SagaRow>(&**conn)
.select(nexus_db_model::LoadedSaga::as_select())
.load_async::<nexus_db_model::LoadedSaga>(&**conn)
.await
.context("fetching sagas")?;

paginator = p
.found_batch(&records_batch, &|s: &nexus_db_model::SagaRow| s.id());
.found_batch(&records_batch, &|s: &nexus_db_model::LoadedSaga| {
s.id()
});

for row in records_batch {
match Saga::try_from(row.clone()) {
Ok(saga_row) => sagas.push(saga_row),
let saga_id = row.id();
match Saga::try_from(row) {
Ok(saga) => sagas.push(saga),
Err(e) => {
eprintln!(
"WARNING: Skipping saga with id {}: {e}",
row.id()
)
eprintln!("WARNING: Skipping saga with id {saga_id}: {e}")
}
};
}
Expand Down Expand Up @@ -785,13 +778,12 @@ async fn cmd_sagas_show(

let saga = {
use nexus_db_schema::schema::saga::dsl;
let row = dsl::saga
dsl::saga
.filter(dsl::id.eq(saga_id))
.select(nexus_db_model::SagaRow::as_select())
.first_async(&*conn)
.select(Saga::as_select())
.first_async::<Saga>(&*conn)
.await
.with_context(|| format!("error fetching saga {saga_id}"))?;
Saga::try_from(row)?
.with_context(|| format!("error fetching saga {saga_id}"))?
};

print_saga_nodes(Some(saga), nodes);
Expand Down
Loading
Loading