Skip to main content

systemprompt_analytics/feedback/
projection.rs

1//! Atomic fenced replacements, durable deltas and commit-ordered owner
2//! generations.
3//!
4//! Copyright (c) systemprompt.io — Business Source License 1.1.
5//! See <https://systemprompt.io> for licensing details.
6
7use 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    // JSON: the prior fact row is retained verbatim for delta history.
18    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}