use cratestack_core::{
AuditEvent, AuditOperation, CratestackContext, CratestackError, ModelEventKind,
};
use crate::audit::{build_audit_event, enqueue_audit_event};
use crate::descriptor::enqueue_event_outbox;
use crate::{ConflictTarget, ModelDescriptor, SqlColumnValue, SqlValue, SqlxRuntime, sqlx};
use super::upsert_do_nothing_authorize::authorize_existing_row;
use super::upsert_do_nothing_sql::upsert_returning_record_do_nothing;
use super::upsert_outcome::UpsertOutcome;
use super::upsert_sql::select_for_update_by_conflict_target;
#[allow(clippy::too_many_arguments)]
pub(super) async fn run_insert_branch<'tx, M, PK>(
tx: &mut sqlx::Transaction<'tx, sqlx::Postgres>,
runtime: &SqlxRuntime,
descriptor: &'static ModelDescriptor<M, PK>,
insert_values: &[SqlColumnValue],
conflict_target: ConflictTarget,
conflict_columns: &[(&'static str, SqlValue)],
ctx: &CratestackContext,
emits_created: bool,
audit_enabled: bool,
) -> Result<(UpsertOutcome<M>, bool, Option<AuditEvent>), CratestackError>
where
for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
{
match upsert_returning_record_do_nothing(&mut **tx, descriptor, insert_values, conflict_target)
.await?
{
Some(record) => {
if emits_created {
enqueue_event_outbox(
&mut **tx,
descriptor.schema_name,
ModelEventKind::Created,
&record,
)
.await?;
}
let mut audit_event = None;
if audit_enabled {
let after = serde_json::to_value(&record).ok();
let event = build_audit_event(descriptor, AuditOperation::Create, None, after, ctx);
enqueue_audit_event(&mut **tx, &event).await?;
audit_event = Some(event);
}
Ok((UpsertOutcome::Inserted(record), emits_created, audit_event))
}
None => {
let existing = select_for_update_by_conflict_target(
&mut **tx,
descriptor,
conflict_columns,
conflict_target.predicate(),
)
.await?
.ok_or_else(|| {
CratestackError::Conflict(format!(
"upsert do_nothing on `{}` lost a conflict race and the \
conflicting row was deleted before it could be read back; retry the call",
descriptor.table_name,
))
})?;
authorize_existing_row(
runtime,
descriptor,
conflict_columns,
conflict_target.predicate(),
ctx,
)
.await?;
Ok((UpsertOutcome::Existing(existing), false, None))
}
}
}