use cratestack_core::{AuditEvent, CratestackContext, CratestackError, ModelEventKind};
use crate::audit::ensure_audit_table;
use crate::descriptor::ensure_event_outbox_table;
use crate::query::support::evaluate_create_policies;
use crate::{ConflictTarget, ModelDescriptor, SqlxRuntime, UpsertModelInput, sqlx};
use super::upsert_do_nothing_authorize::authorize_existing_row;
use super::upsert_do_nothing_insert::run_insert_branch;
use super::upsert_do_nothing_probe::resolve_pre_probe;
use super::upsert_outcome::UpsertOutcome;
use super::upsert_prepare::prepare_upsert_insert;
pub(super) async fn run_upsert_do_nothing_in_tx<'tx, M, PK, I>(
tx: &mut sqlx::Transaction<'tx, sqlx::Postgres>,
runtime: &SqlxRuntime,
descriptor: &'static ModelDescriptor<M, PK>,
input: I,
conflict_target: ConflictTarget,
ctx: &CratestackContext,
) -> Result<(UpsertOutcome<M>, bool, Option<AuditEvent>), CratestackError>
where
I: UpsertModelInput<M>,
for<'r> M: Send + Unpin + sqlx::FromRow<'r, sqlx::postgres::PgRow> + serde::Serialize,
PK: Send + sqlx::Type<sqlx::Postgres> + for<'q> sqlx::Encode<'q, sqlx::Postgres>,
{
input.validate()?;
let (insert_values, conflict_columns) =
prepare_upsert_insert(descriptor, &input, ctx, conflict_target)?;
if !evaluate_create_policies(
runtime.pool(),
descriptor.create_allow_policies,
descriptor.create_deny_policies,
&insert_values,
ctx,
)
.await?
{
return Err(CratestackError::Forbidden(
"create policy denied this upsert".to_owned(),
));
}
let emits_created = descriptor.emits(ModelEventKind::Created);
let audit_enabled = descriptor.audit_enabled;
if emits_created {
ensure_event_outbox_table(&mut **tx).await?;
}
if audit_enabled {
ensure_audit_table(runtime).await?;
}
let pre_probe = resolve_pre_probe(
tx,
descriptor,
&conflict_columns,
&insert_values,
conflict_target,
)
.await?;
if let Some(existing) = pre_probe {
authorize_existing_row(
runtime,
descriptor,
&conflict_columns,
conflict_target.predicate(),
ctx,
)
.await?;
return Ok((UpsertOutcome::Existing(existing), false, None));
}
run_insert_branch(
tx,
runtime,
descriptor,
&insert_values,
conflict_target,
&conflict_columns,
ctx,
emits_created,
audit_enabled,
)
.await
}