use rusqlite::{params, OptionalExtension};
use crate::error::{Result, YantrikDbError};
use super::now;
pub(crate) struct ClaimRow<'a> {
pub origin_actor: &'a str,
pub namespace: &'a str,
pub idempotency_key: &'a str,
pub rid: &'a str,
pub payload_digest: &'a [u8; 32],
pub op_id: &'a str,
pub route: &'static str,
pub generation: i64,
}
pub(crate) struct PendingClaim<'a> {
pub namespace: &'a str,
pub idempotency_key: &'a str,
pub payload_digest: &'a [u8; 32],
pub rid: &'a str,
pub generation: i64,
}
#[must_use]
pub(crate) enum ClaimAttempt {
Won,
Hit { existing_rid: String },
}
fn resolve_existing(
namespace: &str,
payload_digest: &[u8; 32],
existing_rid: String,
existing_digest: &[u8],
state: &str,
) -> Result<ClaimAttempt> {
if existing_digest != &payload_digest[..] {
return Err(YantrikDbError::IdempotencyConflict {
namespace: namespace.to_string(),
existing_rid,
reason: "same idempotency key with a DIFFERENT payload — the first \
write's content stands; change the key or the payload"
.to_string(),
});
}
if state != "committed" {
return Err(YantrikDbError::IdempotencyConflict {
namespace: namespace.to_string(),
existing_rid,
reason: format!(
"claim is in state '{state}', not 'committed' — a crashed \
prior attempt; the recovery sweep that reconciles this is \
not yet implemented"
),
});
}
Ok(ClaimAttempt::Hit { existing_rid })
}
pub(crate) fn probe_committed_claim(
conn: &rusqlite::Connection,
origin_actor: &str,
namespace: &str,
idempotency_key: &str,
payload_digest: &[u8; 32],
) -> Result<Option<String>> {
let existing: Option<(String, Vec<u8>, String)> = conn
.query_row(
"SELECT rid, payload_digest, state FROM idempotency_claims \
WHERE origin_actor = ?1 AND namespace = ?2 AND idempotency_key = ?3",
params![origin_actor, namespace, idempotency_key],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.optional()?;
match existing {
None => Ok(None),
Some((rid, digest, state)) => {
match resolve_existing(namespace, payload_digest, rid, &digest, &state)? {
ClaimAttempt::Hit { existing_rid } => Ok(Some(existing_rid)),
ClaimAttempt::Won => unreachable!("resolve_existing cannot yield Won"),
}
}
}
}
pub(crate) fn claim_in_tx(
conn: &rusqlite::Connection,
claim: &ClaimRow<'_>,
) -> Result<ClaimAttempt> {
let inserted = conn.execute(
"INSERT INTO idempotency_claims \
(origin_actor, namespace, idempotency_key, rid, payload_digest, \
op_id, route, generation, state, created_at) \
VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, 'committed', ?9) \
ON CONFLICT(origin_actor, namespace, idempotency_key) DO NOTHING",
params![
claim.origin_actor,
claim.namespace,
claim.idempotency_key,
claim.rid,
&claim.payload_digest[..],
claim.op_id,
claim.route,
claim.generation,
now(),
],
)?;
if inserted == 1 {
return Ok(ClaimAttempt::Won);
}
let existing: Option<(String, Vec<u8>, String)> = conn
.query_row(
"SELECT rid, payload_digest, state FROM idempotency_claims \
WHERE origin_actor = ?1 AND namespace = ?2 AND idempotency_key = ?3",
params![claim.origin_actor, claim.namespace, claim.idempotency_key],
|r| Ok((r.get(0)?, r.get(1)?, r.get(2)?)),
)
.optional()?;
let Some((existing_rid, existing_digest, state)) = existing else {
return Err(YantrikDbError::IdempotencyConflict {
namespace: claim.namespace.to_string(),
existing_rid: String::new(),
reason: "claim row vanished between ON CONFLICT and read-back \
(engine invariant violation — retry; if it persists, the \
claims table is being mutated outside the engine)"
.to_string(),
});
};
resolve_existing(
claim.namespace,
claim.payload_digest,
existing_rid,
&existing_digest,
&state,
)
}