1use super::columns::Columns;
8use super::{ApplyOutcome, FactLease, FeedbackFactsRepository, validation};
9use crate::Result;
10use systemprompt_identifiers::{
11 AnalyticsFactId, DeviceId, ManagedResourceId, NativeSessionId, ResourceRevisionId, UserId,
12};
13use systemprompt_models::feedback::analytics::{AnalyticsChange, AnalyticsChangeOperation};
14
15struct Replacement {
16 generation: i64,
17 before: Option<serde_json::Value>,
19}
20
21impl FeedbackFactsRepository {
22 pub async fn apply(&self, owner: &UserId, lease: &FactLease) -> Result<ApplyOutcome> {
23 let mut tx = self.pool.begin().await?;
24 sqlx::query!(
25 "INSERT INTO analytics_fact_checkpoints(owner_id) VALUES($1) ON CONFLICT DO NOTHING",
26 owner.as_str()
27 )
28 .execute(&mut *tx)
29 .await?;
30 let checkpoint = sqlx::query!(
31 "SELECT generation FROM analytics_fact_checkpoints WHERE owner_id=$1 FOR UPDATE",
32 owner.as_str()
33 )
34 .fetch_one(&mut *tx)
35 .await?;
36 let row = sqlx::query!("SELECT payload FROM analytics_fact_changes WHERE owner_id=$1 AND change_id=$2 AND state='leased' AND lease_worker=$3 AND lease_epoch=$4 AND lease_until>clock_timestamp() FOR UPDATE", owner.as_str(), lease.change_id.as_str(), lease.worker_id.as_str(), lease.epoch).fetch_optional(&mut *tx).await?.ok_or_else(validation::invalid)?;
37 let change: AnalyticsChange =
38 serde_json::from_value(row.payload.ok_or_else(validation::invalid)?)?;
39 validation::validate(&change)?;
40 let kind = validation::kind(change.key.kind);
41 let revision = i64::try_from(change.revision).map_err(|_error| validation::invalid())?;
42 let before = sqlx::query!("SELECT revision,fact FROM analytics_normalized_facts WHERE owner_id=$1 AND fact_kind=$2 AND source=$3 AND fact_id=$4 FOR UPDATE", owner.as_str(), kind, &change.key.source, change.key.id.as_str()).fetch_optional(&mut *tx).await?;
43 let replaced = before.as_ref().is_none_or(|row| row.revision < revision);
44 let generation = if replaced {
45 checkpoint
46 .generation
47 .checked_add(1)
48 .ok_or_else(validation::invalid)?
49 } else {
50 checkpoint.generation
51 };
52 if replaced {
53 let before_fact = before.and_then(|row| row.fact);
54 self.replace_in(
55 &mut tx,
56 owner,
57 &change,
58 Replacement {
59 generation,
60 before: before_fact,
61 },
62 )
63 .await?;
64 }
65 let state = if replaced { "applied" } else { "superseded" };
66 let updated = sqlx::query!("UPDATE analytics_fact_changes SET state=$5,applied_at=clock_timestamp(),lease_worker=NULL,lease_until=NULL,last_error=NULL WHERE owner_id=$1 AND change_id=$2 AND state='leased' AND lease_worker=$3 AND lease_epoch=$4 AND lease_until>clock_timestamp()", owner.as_str(), lease.change_id.as_str(), lease.worker_id.as_str(), lease.epoch, state).execute(&mut *tx).await?;
67 if updated.rows_affected() != 1 {
68 return Err(validation::invalid());
69 }
70 sqlx::query!("UPDATE analytics_fact_checkpoints SET generation=$2,applied_changes=applied_changes+$3,superseded_changes=superseded_changes+$4,last_applied_at=clock_timestamp(),last_recorded_at=GREATEST(last_recorded_at,$5) WHERE owner_id=$1", owner.as_str(), generation, i64::from(replaced), i64::from(!replaced), change.recorded_at).execute(&mut *tx).await?;
71 tx.commit().await?;
72 Ok(ApplyOutcome {
73 change_id: lease.change_id.clone(),
74 generation,
75 replaced,
76 })
77 }
78
79 async fn replace_in(
80 &self,
81 tx: &mut sqlx::Transaction<'_, sqlx::Postgres>,
82 owner: &UserId,
83 change: &AnalyticsChange,
84 replacement: Replacement,
85 ) -> Result<()> {
86 let Replacement { generation, before } = replacement;
87 let fact = match &change.operation {
88 AnalyticsChangeOperation::Replace { fact } => Some(fact),
89 AnalyticsChangeOperation::Tombstone => None,
90 };
91 let columns = Columns::of(&change.key.source, fact)?;
92 let after = fact.map(serde_json::to_value).transpose()?;
93 let kind = validation::kind(change.key.kind);
94 let revision = i64::try_from(change.revision).map_err(|_error| validation::invalid())?;
95 sqlx::query!("INSERT INTO analytics_normalized_facts(owner_id,fact_kind,source,fact_id,revision,occurred_at,deleted,fact,consumer_id,device_id,host,session_id,resource_id,resource_revision_id,invocation_source,invocation_id,request_source,request_id,succeeded,currency,amount_micros,input_tokens,output_tokens,latency_micros,assessment_status,score_millionths,generation,conversation_source,conversation_id) VALUES($1,$2,$3,$4,$5,$6,$7,$8,$9,$10,$11,$12,$13,$14,$15,$16,$17,$18,$19,$20,$21,$22,$23,$24,$25,$26,$27,$28,$29) ON CONFLICT(owner_id,fact_kind,source,fact_id) DO UPDATE SET revision=EXCLUDED.revision,occurred_at=EXCLUDED.occurred_at,deleted=EXCLUDED.deleted,fact=EXCLUDED.fact,consumer_id=EXCLUDED.consumer_id,device_id=EXCLUDED.device_id,host=EXCLUDED.host,session_id=EXCLUDED.session_id,resource_id=EXCLUDED.resource_id,resource_revision_id=EXCLUDED.resource_revision_id,invocation_source=EXCLUDED.invocation_source,invocation_id=EXCLUDED.invocation_id,request_source=EXCLUDED.request_source,request_id=EXCLUDED.request_id,succeeded=EXCLUDED.succeeded,currency=EXCLUDED.currency,amount_micros=EXCLUDED.amount_micros,input_tokens=EXCLUDED.input_tokens,output_tokens=EXCLUDED.output_tokens,latency_micros=EXCLUDED.latency_micros,assessment_status=EXCLUDED.assessment_status,score_millionths=EXCLUDED.score_millionths,generation=EXCLUDED.generation,conversation_source=EXCLUDED.conversation_source,conversation_id=EXCLUDED.conversation_id,updated_at=clock_timestamp()", owner.as_str(), kind, &change.key.source, change.key.id.as_str(), revision, change.occurred_at, fact.is_none(), after, columns.consumer.as_ref().map(UserId::as_str), columns.device.as_ref().map(DeviceId::as_str), columns.host, columns.session.as_ref().map(NativeSessionId::as_str), columns.resource.as_ref().map(ManagedResourceId::as_str), columns.resource_revision.as_ref().map(ResourceRevisionId::as_str), columns.invocation_source, columns.invocation.as_ref().map(AnalyticsFactId::as_str), columns.request_source, columns.request.as_ref().map(AnalyticsFactId::as_str), columns.succeeded, columns.currency, columns.amount, columns.input, columns.output, columns.latency, columns.assessment, columns.score, generation, columns.conversation_source, columns.conversation.as_ref().map(AnalyticsFactId::as_str)).execute(&mut **tx).await?;
96 sqlx::query!("INSERT INTO analytics_fact_deltas(owner_id,generation,fact_kind,source,fact_id,before_fact,after_fact,occurred_at) VALUES($1,$2,$3,$4,$5,$6,$7,$8)", owner.as_str(), generation, kind, &change.key.source, change.key.id.as_str(), before, after,change.occurred_at).execute(&mut **tx).await?;
97 Ok(())
98 }
99}