use std::collections::HashMap;
use std::sync::Arc;
use async_trait::async_trait;
use chat_engine_sdk::models::{
FileCitation, LinkCitation, LinkReference, MessagePart, MessagePartInput,
};
use sea_orm::sea_query::Expr;
use sea_orm::{ActiveValue::Set, ColumnTrait, Condition, EntityTrait, QueryOrder};
use serde_json::Value as JsonValue;
use time::OffsetDateTime;
use toolkit_db::secure::{
AccessScope, DBRunner, SecureDeleteExt, SecureEntityExt, SecureInsertExt, SecureUpdateExt,
TxConfig,
};
use uuid::Uuid;
use crate::domain::error::ChatEngineError;
use crate::domain::message::Message;
use crate::domain::ports::{
FinalizeOutcome, InsertedPair, MessageRepo, NewUserMessage, PartCitations,
};
use crate::infra::db::conversions::part_type_to_entity;
use crate::infra::db::entity::message::{self as message_entity, Entity as MessageEntity};
use crate::infra::db::entity::message_part::{
self as message_part_entity, Entity as MessagePartEntity, compute_next_part_number,
};
use crate::infra::db::entity::{
file_citation as file_citation_entity, link_citation as link_citation_entity,
link_reference as link_reference_entity,
};
use crate::infra::db::repo::ChatEngineDb;
use crate::infra::db::{
VARIANT_INDEX_MAX_RETRIES, compute_next_variant_index, is_variant_unique_violation,
};
pub struct SeaMessageRepo {
db: Arc<ChatEngineDb>,
}
impl SeaMessageRepo {
#[must_use]
pub fn new(db: Arc<ChatEngineDb>) -> Self {
Self { db }
}
async fn attach_parts(&self, mut msgs: Vec<Message>) -> Result<Vec<Message>, ChatEngineError> {
if msgs.is_empty() {
return Ok(msgs);
}
let ids: Vec<Uuid> = msgs.iter().map(|m| m.message_id).collect();
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let rows = MessagePartEntity::find()
.order_by_asc(message_part_entity::Column::MessageId)
.order_by_asc(message_part_entity::Column::Number)
.secure()
.scope_with(&scope)
.filter(Condition::all().add(message_part_entity::Column::MessageId.is_in(ids)))
.all(&conn)
.await?;
let mut parts: Vec<MessagePart> = rows.into_iter().map(MessagePart::from).collect();
self.attach_citations(&mut parts).await?;
let mut by_msg: HashMap<Uuid, Vec<MessagePart>> = HashMap::new();
for part in parts {
by_msg.entry(part.message_id).or_default().push(part);
}
for m in &mut msgs {
if let Some(parts) = by_msg.remove(&m.message_id) {
m.parts = parts;
}
}
Ok(msgs)
}
async fn attach_citations(&self, parts: &mut [MessagePart]) -> Result<(), ChatEngineError> {
if parts.is_empty() {
return Ok(());
}
let ids: Vec<Uuid> = parts.iter().map(|p| p.id).collect();
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let file_rows = file_citation_entity::Entity::find()
.order_by_asc(file_citation_entity::Column::Number)
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(file_citation_entity::Column::MessagePartId.is_in(ids.clone())),
)
.all(&conn)
.await?;
let link_rows = link_citation_entity::Entity::find()
.order_by_asc(link_citation_entity::Column::Number)
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(link_citation_entity::Column::MessagePartId.is_in(ids.clone())),
)
.all(&conn)
.await?;
let ref_rows = link_reference_entity::Entity::find()
.order_by_asc(link_reference_entity::Column::Number)
.secure()
.scope_with(&scope)
.filter(Condition::all().add(link_reference_entity::Column::MessagePartId.is_in(ids)))
.all(&conn)
.await?;
let mut files: HashMap<Uuid, Vec<FileCitation>> = HashMap::new();
for row in file_rows {
if let Ok(c) = serde_json::from_value::<FileCitation>(row.content) {
files.entry(row.message_part_id).or_default().push(c);
}
}
let mut links: HashMap<Uuid, Vec<LinkCitation>> = HashMap::new();
for row in link_rows {
if let Ok(c) = serde_json::from_value::<LinkCitation>(row.content) {
links.entry(row.message_part_id).or_default().push(c);
}
}
let mut refs: HashMap<Uuid, Vec<LinkReference>> = HashMap::new();
for row in ref_rows {
if let Ok(r) = serde_json::from_value::<LinkReference>(row.content) {
refs.entry(row.message_part_id).or_default().push(r);
}
}
for p in parts.iter_mut() {
if let Some(v) = files.remove(&p.id) {
p.file_citations = v;
}
if let Some(v) = links.remove(&p.id) {
p.link_citations = v;
}
if let Some(v) = refs.remove(&p.id) {
p.references = v;
}
}
Ok(())
}
}
#[async_trait]
impl MessageRepo for SeaMessageRepo {
async fn insert_user_and_assistant_stub(
&self,
req: NewUserMessage,
) -> Result<InsertedPair, ChatEngineError> {
let session_id = req.session_id;
let parent = req.parent_message_id;
let file_ids_json = req
.file_ids
.as_ref()
.filter(|ids| !ids.is_empty())
.and_then(|ids| serde_json::to_value(ids).ok());
let parts = req.parts;
let metadata = req.metadata;
let tenant_id = req.tenant_id;
let author_id = req.user_id;
let mut last_err: Option<ChatEngineError> = None;
for _attempt in 0..VARIANT_INDEX_MAX_RETRIES {
let user_message_id = Uuid::new_v4();
let assistant_message_id = Uuid::new_v4();
let now = OffsetDateTime::now_utc();
let parts_attempt = parts.clone();
let metadata_attempt = metadata.clone();
let file_ids_attempt = file_ids_json.clone();
let user_tenant = tenant_id.clone();
let author = author_id.clone();
let assistant_tenant = tenant_id.clone();
let outcome: Result<i32, ChatEngineError> = self
.db
.transaction_with_config(TxConfig::serializable(), move |tx| {
Box::pin(async move {
let scope = AccessScope::allow_all();
let user_variant_index =
compute_next_variant_index(tx, session_id, parent).await?;
let user_active = message_entity::ActiveModel {
message_id: Set(user_message_id),
session_id: Set(session_id),
tenant_id: Set(user_tenant),
user_id: Set(author),
parent_message_id: Set(parent),
role: Set(message_entity::MessageRole::User),
file_ids: Set(file_ids_attempt),
variant_index: Set(user_variant_index),
is_active: Set(true),
is_complete: Set(true),
is_hidden_from_user: Set(false),
is_hidden_from_backend: Set(false),
metadata: Set(metadata_attempt),
created_at: Set(now),
};
let assistant_active = message_entity::ActiveModel {
message_id: Set(assistant_message_id),
session_id: Set(session_id),
tenant_id: Set(assistant_tenant),
user_id: Set(None),
parent_message_id: Set(Some(user_message_id)),
role: Set(message_entity::MessageRole::Assistant),
file_ids: Set(None),
variant_index: Set(0),
is_active: Set(true),
is_complete: Set(false),
is_hidden_from_user: Set(false),
is_hidden_from_backend: Set(false),
metadata: Set(None),
created_at: Set(now),
};
MessageEntity::insert(user_active)
.secure()
.scope_unchecked(&scope)?
.exec(tx)
.await?;
insert_message_parts(tx, &scope, user_message_id, &parts_attempt).await?;
MessageEntity::insert(assistant_active)
.secure()
.scope_unchecked(&scope)?
.exec(tx)
.await?;
Ok(user_variant_index)
})
})
.await;
match outcome {
Ok(user_variant_index) => {
return Ok(InsertedPair {
user_message_id,
assistant_message_id,
user_variant_index,
});
}
Err(e) => {
if let Some(db_err) = chat_engine_db_err(&e) {
if !is_variant_unique_violation(db_err) {
return Err(e);
}
} else {
return Err(e);
}
last_err = Some(e);
}
}
}
Err(retry_exhausted_conflict(last_err))
}
async fn finalize_assistant(
&self,
session_id: Uuid,
assistant_message_id: Uuid,
outcome: FinalizeOutcome,
) -> Result<(), ChatEngineError> {
let (text, metadata, is_complete, citations, extra_parts) = match outcome {
FinalizeOutcome::Complete {
text,
metadata,
citations,
extra_parts,
} => (text, metadata, true, citations, extra_parts),
FinalizeOutcome::Cancelled { text } => {
let mut meta = serde_json::Map::new();
meta.insert("cancelled".into(), JsonValue::Bool(true));
meta.insert("partial".into(), JsonValue::Bool(true));
(
text,
Some(JsonValue::Object(meta)),
false,
PartCitations::default(),
Vec::new(),
)
}
FinalizeOutcome::Errored {
text,
error,
finish_reason,
} => {
let mut meta = serde_json::Map::new();
meta.insert(
"finish_reason".into(),
JsonValue::String(finish_reason.to_string()),
);
meta.insert("error".into(), JsonValue::String(error));
meta.insert("partial".into(), JsonValue::Bool(true));
(
text,
Some(JsonValue::Object(meta)),
false,
PartCitations::default(),
Vec::new(),
)
}
};
let rows_affected = self
.db
.transaction(move |tx| {
Box::pin(async move {
let scope = AccessScope::allow_all();
let result = MessageEntity::update_many()
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::MessageId.eq(assistant_message_id))
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::IsComplete.eq(false))
.add(message_entity::Column::Metadata.is_null()),
)
.col_expr(message_entity::Column::IsComplete, Expr::value(is_complete))
.col_expr(message_entity::Column::Metadata, Expr::value(metadata))
.exec(tx)
.await?;
if result.rows_affected == 1 {
let part_id =
append_text_part(tx, &scope, assistant_message_id, &text).await?;
insert_part_citations(tx, &scope, part_id, &citations).await?;
for part in &extra_parts {
let extra_id =
append_part(tx, &scope, assistant_message_id, part).await?;
let part_citations = PartCitations {
file_citations: part.file_citations.clone(),
link_citations: part.link_citations.clone(),
references: part.references.clone(),
};
if !part_citations.is_empty() {
insert_part_citations(tx, &scope, extra_id, &part_citations)
.await?;
}
}
}
Ok::<u64, ChatEngineError>(result.rows_affected)
})
})
.await?;
if rows_affected == 0 {
tracing::debug!(
session_id = %session_id,
assistant_message_id = %assistant_message_id,
"finalize_assistant no-op: stub already terminated or not in session",
);
}
Ok(())
}
async fn fetch_active_history(
&self,
session_id: Uuid,
depth: Option<u32>,
) -> Result<Vec<Message>, ChatEngineError> {
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let mut query = MessageEntity::find()
.order_by_asc(message_entity::Column::CreatedAt)
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::IsActive.eq(true))
.add(message_entity::Column::IsHiddenFromBackend.eq(false))
.add(message_entity::Column::IsComplete.eq(true)),
);
if let Some(d) = depth {
query = query.limit(u64::from(d));
}
let rows = query.all(&conn).await?;
self.attach_parts(rows.into_iter().map(Message::from).collect())
.await
}
async fn find_message_in_session(
&self,
session_id: Uuid,
message_id: Uuid,
) -> Result<Option<Message>, ChatEngineError> {
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let row = MessageEntity::find()
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::MessageId.eq(message_id))
.add(message_entity::Column::SessionId.eq(session_id)),
)
.one(&conn)
.await?;
match row {
Some(row) => Ok(self.attach_parts(vec![Message::from(row)]).await?.pop()),
None => Ok(None),
}
}
async fn find_message_by_id(
&self,
message_id: Uuid,
) -> Result<Option<Message>, ChatEngineError> {
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let row = MessageEntity::find()
.secure()
.scope_with(&scope)
.filter(Condition::all().add(message_entity::Column::MessageId.eq(message_id)))
.one(&conn)
.await?;
match row {
Some(row) => Ok(self.attach_parts(vec![Message::from(row)]).await?.pop()),
None => Ok(None),
}
}
async fn list_active_path(&self, session_id: Uuid) -> Result<Vec<Message>, ChatEngineError> {
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let rows = MessageEntity::find()
.order_by_asc(message_entity::Column::CreatedAt)
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::IsActive.eq(true))
.add(message_entity::Column::IsComplete.eq(true)),
)
.all(&conn)
.await?;
self.attach_parts(rows.into_iter().map(Message::from).collect())
.await
}
async fn list_non_root_messages_chrono(
&self,
session_id: Uuid,
) -> Result<Vec<Message>, ChatEngineError> {
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let rows = MessageEntity::find()
.order_by_asc(message_entity::Column::CreatedAt)
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::ParentMessageId.is_not_null()),
)
.all(&conn)
.await?;
self.attach_parts(rows.into_iter().map(Message::from).collect())
.await
}
async fn list_non_root_messages_older_than(
&self,
session_id: Uuid,
older_than: OffsetDateTime,
) -> Result<Vec<Message>, ChatEngineError> {
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let rows = MessageEntity::find()
.order_by_asc(message_entity::Column::CreatedAt)
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::ParentMessageId.is_not_null())
.add(message_entity::Column::CreatedAt.lt(older_than)),
)
.all(&conn)
.await?;
self.attach_parts(rows.into_iter().map(Message::from).collect())
.await
}
async fn count_non_root_messages(&self, session_id: Uuid) -> Result<u64, ChatEngineError> {
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let n = MessageEntity::find()
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::ParentMessageId.is_not_null()),
)
.count(&conn)
.await?;
Ok(n)
}
async fn list_oldest_non_root_message_ids(
&self,
session_id: Uuid,
limit: u32,
) -> Result<Vec<Uuid>, ChatEngineError> {
if limit == 0 {
return Ok(Vec::new());
}
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let rows = MessageEntity::find()
.order_by_asc(message_entity::Column::CreatedAt)
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::ParentMessageId.is_not_null()),
)
.limit(u64::from(limit))
.all(&conn)
.await?;
Ok(rows.into_iter().map(|m| m.message_id).collect())
}
async fn list_non_root_message_ids_older_than(
&self,
session_id: Uuid,
older_than: OffsetDateTime,
limit: u32,
) -> Result<Vec<Uuid>, ChatEngineError> {
if limit == 0 {
return Ok(Vec::new());
}
let conn = self.db.conn()?;
let scope = AccessScope::allow_all();
let rows = MessageEntity::find()
.order_by_asc(message_entity::Column::CreatedAt)
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::ParentMessageId.is_not_null())
.add(message_entity::Column::CreatedAt.lt(older_than)),
)
.limit(u64::from(limit))
.all(&conn)
.await?;
Ok(rows.into_iter().map(|m| m.message_id).collect())
}
async fn insert_summary_message(
&self,
session_id: Uuid,
text: String,
metadata: Option<JsonValue>,
summarized_ids: Vec<Uuid>,
tenant_id: Option<String>,
) -> Result<Uuid, ChatEngineError> {
let summary_id = Uuid::new_v4();
let now = OffsetDateTime::now_utc();
let summarized = summarized_ids.clone();
self.db
.transaction_with_config(TxConfig::serializable(), move |tx| {
Box::pin(async move {
let scope = AccessScope::allow_all();
let variant_index = compute_next_variant_index(tx, session_id, None).await?;
let summary_active = message_entity::ActiveModel {
message_id: Set(summary_id),
session_id: Set(session_id),
tenant_id: Set(tenant_id),
user_id: Set(None),
parent_message_id: Set(None),
role: Set(message_entity::MessageRole::System),
file_ids: Set(None),
variant_index: Set(variant_index),
is_active: Set(true),
is_complete: Set(true),
is_hidden_from_user: Set(true),
is_hidden_from_backend: Set(false),
metadata: Set(metadata),
created_at: Set(now),
};
MessageEntity::insert(summary_active)
.secure()
.scope_unchecked(&scope)?
.exec(tx)
.await?;
append_text_part(tx, &scope, summary_id, &text).await?;
if !summarized.is_empty() {
MessageEntity::update_many()
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(
message_entity::Column::MessageId.is_in(summarized.clone()),
),
)
.col_expr(
message_entity::Column::IsHiddenFromBackend,
Expr::value(true),
)
.exec(tx)
.await?;
}
Ok(summary_id)
})
})
.await
}
async fn delete_message_subtree(
&self,
session_id: Uuid,
root_id: Uuid,
) -> Result<u64, ChatEngineError> {
self.db
.transaction_with_config(TxConfig::serializable(), move |tx| {
Box::pin(async move {
let scope = AccessScope::allow_all();
let mut levels: Vec<Vec<Uuid>> = Vec::new();
let mut frontier: Vec<Uuid> = vec![root_id];
loop {
let children: Vec<Uuid> = MessageEntity::find()
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(
message_entity::Column::ParentMessageId
.is_in(frontier.clone()),
),
)
.all(tx)
.await?
.into_iter()
.map(|m| m.message_id)
.collect();
levels.push(frontier);
if children.is_empty() {
break;
}
frontier = children;
}
let mut removed: u64 = 0;
for level in levels.into_iter().rev() {
let res = MessageEntity::delete_many()
.secure()
.scope_with(&scope)
.filter(
Condition::all()
.add(message_entity::Column::SessionId.eq(session_id))
.add(message_entity::Column::MessageId.is_in(level)),
)
.exec(tx)
.await?;
removed += res.rows_affected;
}
Ok(removed)
})
})
.await
}
async fn insert_assistant_variant_stub(
&self,
session_id: Uuid,
parent_message_id: Uuid,
tenant_id: Option<String>,
) -> Result<InsertedPair, ChatEngineError> {
let mut last_err: Option<ChatEngineError> = None;
for _attempt in 0..VARIANT_INDEX_MAX_RETRIES {
let new_message_id = Uuid::new_v4();
let now = OffsetDateTime::now_utc();
let variant_tenant = tenant_id.clone();
let outcome: Result<i32, ChatEngineError> = self
.db
.transaction_with_config(TxConfig::serializable(), move |tx| {
Box::pin(async move {
let scope = AccessScope::allow_all();
let new_variant_index =
compute_next_variant_index(tx, session_id, Some(parent_message_id))
.await?;
let assistant_active = message_entity::ActiveModel {
message_id: Set(new_message_id),
session_id: Set(session_id),
tenant_id: Set(variant_tenant),
user_id: Set(None),
parent_message_id: Set(Some(parent_message_id)),
role: Set(message_entity::MessageRole::Assistant),
file_ids: Set(None),
variant_index: Set(new_variant_index),
is_active: Set(true),
is_complete: Set(false),
is_hidden_from_user: Set(false),
is_hidden_from_backend: Set(false),
metadata: Set(None),
created_at: Set(now),
};
MessageEntity::insert(assistant_active)
.secure()
.scope_unchecked(&scope)?
.exec(tx)
.await?;
Ok(new_variant_index)
})
})
.await;
match outcome {
Ok(new_variant_index) => {
return Ok(InsertedPair {
user_message_id: parent_message_id,
assistant_message_id: new_message_id,
user_variant_index: new_variant_index,
});
}
Err(e) => {
if let Some(db_err) = chat_engine_db_err(&e) {
if !is_variant_unique_violation(db_err) {
return Err(e);
}
} else {
return Err(e);
}
last_err = Some(e);
}
}
}
Err(retry_exhausted_conflict(last_err))
}
}
fn chat_engine_db_err(err: &ChatEngineError) -> Option<&sea_orm::DbErr> {
let ChatEngineError::Internal { source, .. } = err else {
return None;
};
let source = source.as_ref()?;
source.downcast_ref::<sea_orm::DbErr>().or_else(|| {
source
.downcast_ref::<toolkit_db::DbError>()
.and_then(|dbe| match dbe {
toolkit_db::DbError::Sea(inner) => Some(inner),
_ => None,
})
})
}
fn retry_exhausted_conflict(last_err: Option<ChatEngineError>) -> ChatEngineError {
let base = format!(
"variant index allocation contended; exhausted {VARIANT_INDEX_MAX_RETRIES} retries"
);
ChatEngineError::conflict(match last_err {
Some(e) => format!("{base}: {e}"),
None => base,
})
}
fn text_part_content(text: &str) -> JsonValue {
serde_json::json!({ "text": text })
}
pub(crate) async fn insert_message_parts<R>(
runner: &R,
scope: &AccessScope,
message_id: Uuid,
parts: &[MessagePartInput],
) -> Result<(), ChatEngineError>
where
R: DBRunner,
{
for (idx, part) in parts.iter().enumerate() {
let am = message_part_entity::ActiveModel {
id: Set(Uuid::new_v4()),
message_id: Set(message_id),
r#type: Set(part_type_to_entity(&part.part_type)),
content: Set(part.content.clone()),
number: Set(i32::try_from(idx).unwrap_or(i32::MAX)),
};
MessagePartEntity::insert(am)
.secure()
.scope_unchecked(scope)?
.exec(runner)
.await?;
}
Ok(())
}
async fn append_text_part<R>(
runner: &R,
scope: &AccessScope,
message_id: Uuid,
text: &str,
) -> Result<Uuid, ChatEngineError>
where
R: DBRunner,
{
let number = compute_next_part_number(runner, message_id).await?;
let part_id = Uuid::new_v4();
let am = message_part_entity::ActiveModel {
id: Set(part_id),
message_id: Set(message_id),
r#type: Set(message_part_entity::MessagePartType::Text),
content: Set(text_part_content(text)),
number: Set(number),
};
MessagePartEntity::insert(am)
.secure()
.scope_unchecked(scope)?
.exec(runner)
.await?;
Ok(part_id)
}
async fn append_part<R>(
runner: &R,
scope: &AccessScope,
message_id: Uuid,
part: &MessagePartInput,
) -> Result<Uuid, ChatEngineError>
where
R: DBRunner,
{
let number = compute_next_part_number(runner, message_id).await?;
let part_id = Uuid::new_v4();
let am = message_part_entity::ActiveModel {
id: Set(part_id),
message_id: Set(message_id),
r#type: Set(part_type_to_entity(&part.part_type)),
content: Set(part.content.clone()),
number: Set(number),
};
MessagePartEntity::insert(am)
.secure()
.scope_unchecked(scope)?
.exec(runner)
.await?;
Ok(part_id)
}
async fn insert_part_citations<R>(
runner: &R,
scope: &AccessScope,
part_id: Uuid,
citations: &PartCitations,
) -> Result<(), ChatEngineError>
where
R: DBRunner,
{
for (idx, c) in citations.file_citations.iter().enumerate() {
let am = file_citation_entity::ActiveModel {
id: Set(Uuid::new_v4()),
message_part_id: Set(part_id),
content: Set(serde_json::to_value(c).unwrap_or(JsonValue::Null)),
number: Set(i32::try_from(idx).unwrap_or(i32::MAX)),
};
file_citation_entity::Entity::insert(am)
.secure()
.scope_unchecked(scope)?
.exec(runner)
.await?;
}
for (idx, c) in citations.link_citations.iter().enumerate() {
let am = link_citation_entity::ActiveModel {
id: Set(Uuid::new_v4()),
message_part_id: Set(part_id),
content: Set(serde_json::to_value(c).unwrap_or(JsonValue::Null)),
number: Set(i32::try_from(idx).unwrap_or(i32::MAX)),
};
link_citation_entity::Entity::insert(am)
.secure()
.scope_unchecked(scope)?
.exec(runner)
.await?;
}
for (idx, r) in citations.references.iter().enumerate() {
let am = link_reference_entity::ActiveModel {
id: Set(Uuid::new_v4()),
message_part_id: Set(part_id),
content: Set(serde_json::to_value(r).unwrap_or(JsonValue::Null)),
number: Set(i32::try_from(idx).unwrap_or(i32::MAX)),
};
link_reference_entity::Entity::insert(am)
.secure()
.scope_unchecked(scope)?
.exec(runner)
.await?;
}
Ok(())
}