use async_trait::async_trait;
use chrono::{DateTime, Utc};
use sqlx::PgConnection;
use turnframe_core::ids::AccountId;
use turnframe_store::commit::{CommitBundle, CommitReceipt, CommitStore};
use turnframe_store::error::StoreError;
use crate::codec::now;
use crate::store::{PgStores, commit as commit_transaction};
use crate::{conversations, events, interactions, journal, outbox, replay};
pub(crate) async fn apply_bundle(
conn: &mut PgConnection,
account: &AccountId,
bundle: CommitBundle,
at: DateTime<Utc>,
) -> Result<CommitReceipt, StoreError> {
bundle.validate(account)?;
let mut event_ids = Vec::new();
let mut inserted_interactions = Vec::new();
let mut invalidated_interactions = Vec::new();
for completion in bundle.journal_completions {
journal::complete(
&mut *conn,
account,
&completion.command_id,
completion.outcome,
at,
)
.await?;
}
for batch in bundle.events {
event_ids.extend(events::append(&mut *conn, batch).await?);
}
for finish in bundle.interaction_finishes {
interactions::finish_resolution(
&mut *conn,
account,
&finish.interaction_id,
finish.outcome,
)
.await?;
}
for invalidation in bundle.interaction_invalidations {
invalidated_interactions.extend(
interactions::invalidate_for_case(
&mut *conn,
account,
&invalidation.case_key,
invalidation.new_revision,
invalidation.reason,
at,
)
.await?,
);
}
for insert in bundle.interaction_inserts {
let id = insert.interaction.id;
invalidated_interactions.extend(
interactions::insert_interaction(
&mut *conn,
insert.interaction,
insert.replace_blocking,
at,
)
.await?,
);
inserted_interactions.push(id);
}
for entry in bundle.outbox_entries {
outbox::enqueue(&mut *conn, entry).await?;
}
if let Some(record) = bundle.replay_record {
replay::put(&mut *conn, record).await?;
}
if let Some(update) = bundle.turn_phase {
conversations::set_turn_phase(&mut *conn, account, &update.turn_id, update.phase, at)
.await?;
}
Ok(CommitReceipt {
event_ids,
inserted_interactions,
invalidated_interactions,
committed_at: at,
})
}
impl PgStores {
pub async fn commit_in(
&self,
conn: &mut PgConnection,
account: &AccountId,
bundle: CommitBundle,
) -> Result<CommitReceipt, StoreError> {
apply_bundle(conn, account, bundle, now()).await
}
}
#[async_trait]
impl CommitStore for PgStores {
async fn commit(
&self,
account: &AccountId,
bundle: CommitBundle,
) -> Result<CommitReceipt, StoreError> {
bundle.validate(account)?;
let mut transaction = self.transaction().await?;
let receipt = apply_bundle(&mut transaction, account, bundle, now()).await?;
commit_transaction(transaction).await?;
Ok(receipt)
}
}