Skip to main content

valence_core/ownership/service/
read.rs

1//! Ownership read paths: lookups, pending-deletion subsets, and owner rollups.
2
3use std::collections::HashSet;
4use std::sync::Arc;
5
6use serde_json::Value;
7
8use crate::error::{Error, Result};
9use crate::query::{QueryCore, RecordPredicate, SortDirection};
10use crate::runtime::Valence;
11use crate::schema::SchemaRegistry;
12
13use super::helpers::{
14    normalize_pending_deletion_query_value, owner_id_query_values, ownership_colocate_enabled,
15    ownership_row_id, parse_count_from_row, schema_skipped_for_owner_summary,
16    skip_ownership_for_table, system_valence, OwnerDataSummary, OwnerSchemaRowCount,
17    OWNER_SUMMARY_CONCURRENCY,
18};
19use super::OwnershipService;
20
21impl OwnershipService {
22    /// Resolve the backend that stores `valence_data_ownership` rows for `valence_model`.
23    pub(crate) fn ownership_backend(
24        valence_model: &str,
25        v: &Valence,
26    ) -> Result<Arc<dyn crate::backend::DatabaseBackend>> {
27        if ownership_colocate_enabled()
28            && !skip_ownership_for_table(valence_model)
29            && SchemaRegistry::global().has_schema(valence_model)
30        {
31            v.backend_for_table(valence_model)
32        } else {
33            v.backend_for_table("valence_data_ownership")
34        }
35    }
36
37    async fn pending_deletion_ids_via_in_query(
38        valence_model: &str,
39        bare_record_ids: &[String],
40        v: &Valence,
41    ) -> Result<HashSet<String>> {
42        let sys = system_valence(v);
43        let backend = Self::ownership_backend(valence_model, &sys)?;
44        let q = concat!(
45            "SELECT VALUE record_id FROM valence_data_ownership ",
46            "WHERE valence_model = $model AND status = 'pending_deletion' AND record_id IN $ids"
47        );
48        let compiled = crate::compiled_query::CompiledQuery::new(
49            q.to_string(),
50            vec![
51                (
52                    "model".to_string(),
53                    Value::String(valence_model.to_string()),
54                ),
55                (
56                    "ids".to_string(),
57                    Value::Array(bare_record_ids.iter().cloned().map(Value::String).collect()),
58                ),
59            ],
60        );
61        let rows = backend
62            .execute_compiled_query(&compiled)
63            .await
64            .map_err(|e| Error::Database(e.to_string()))?;
65        Ok(rows
66            .into_iter()
67            .filter_map(|v| v.as_str().map(normalize_pending_deletion_query_value))
68            .collect())
69    }
70
71    /// Load ownership JSON for a row, if present.
72    pub async fn get_ownership_json(
73        valence_model: &str,
74        record_id: &str,
75        v: &Valence,
76    ) -> Result<Option<Value>> {
77        let id = ownership_row_id(valence_model, record_id);
78        let sys = system_valence(v);
79        let backend = Self::ownership_backend(valence_model, &sys)?;
80        backend
81            .get_record("valence_data_ownership", &id)
82            .await
83            .map_err(|e| Error::Database(e.to_string()))
84    }
85
86    /// Subset of `bare_record_ids` that currently have `status = pending_deletion` in ownership.
87    pub async fn pending_deletion_bare_ids_subset(
88        valence_model: &str,
89        bare_record_ids: &[String],
90        v: &Valence,
91    ) -> Result<HashSet<String>> {
92        if skip_ownership_for_table(valence_model) || bare_record_ids.is_empty() {
93            return Ok(HashSet::new());
94        }
95        if ownership_colocate_enabled() {
96            return Self::pending_deletion_ids_via_in_query(valence_model, bare_record_ids, v)
97                .await;
98        }
99
100        let sys = system_valence(v);
101        const POINT_LOOKUP_MAX: usize = 64;
102        if bare_record_ids.len() <= POINT_LOOKUP_MAX {
103            let mut out = HashSet::with_capacity(bare_record_ids.len());
104            let backend = Self::ownership_backend(valence_model, &sys)?;
105            for bare_id in bare_record_ids {
106                let id = ownership_row_id(valence_model, bare_id);
107                let Some(json) = backend
108                    .get_record("valence_data_ownership", &id)
109                    .await
110                    .map_err(|e| Error::Database(e.to_string()))?
111                else {
112                    continue;
113                };
114                if json
115                    .get("status")
116                    .and_then(|s| s.as_str())
117                    .is_some_and(|s| s == "pending_deletion")
118                {
119                    out.insert(bare_id.clone());
120                }
121            }
122            return Ok(out);
123        }
124
125        Self::pending_deletion_ids_via_in_query(valence_model, bare_record_ids, v).await
126    }
127
128    async fn count_ownership_rows_for_schema(
129        valence_model: &str,
130        owner_id: &str,
131        owner_type: &str,
132        status: &str,
133        v: &Valence,
134    ) -> Result<u64> {
135        let sys = system_valence(v);
136        let backend = Self::ownership_backend(valence_model, &sys)?;
137        let owner_ids = owner_id_query_values(owner_id, owner_type);
138        let q = concat!(
139            "SELECT count() AS n FROM valence_data_ownership ",
140            "WHERE valence_model = $model AND owner_id IN $owner_ids ",
141            "AND owner_type = $owner_type AND status = $status GROUP ALL"
142        );
143        let compiled = crate::compiled_query::CompiledQuery::new(
144            q.to_string(),
145            vec![
146                (
147                    "model".to_string(),
148                    Value::String(valence_model.to_string()),
149                ),
150                (
151                    "owner_ids".to_string(),
152                    Value::Array(owner_ids.into_iter().map(Value::String).collect()),
153                ),
154                (
155                    "owner_type".to_string(),
156                    Value::String(owner_type.to_string()),
157                ),
158                ("status".to_string(), Value::String(status.to_string())),
159            ],
160        );
161        let rows = backend
162            .execute_compiled_query(&compiled)
163            .await
164            .map_err(|e| Error::Database(e.to_string()))?;
165        Ok(rows.first().map_or(0, parse_count_from_row))
166    }
167
168    /// Roll up ownership sidecar rows for `owner_id` / `owner_type` across registered schemas.
169    pub async fn owner_data_summary(
170        owner_id: &str,
171        owner_type: &str,
172        v: &Valence,
173    ) -> Result<OwnerDataSummary> {
174        let registry = SchemaRegistry::global();
175        let schemas: Vec<&str> = registry
176            .list_schemas()
177            .into_iter()
178            .filter(|table| {
179                registry
180                    .get_full_schema(table)
181                    .is_some_and(|s| !schema_skipped_for_owner_summary(table, s))
182            })
183            .collect();
184
185        let semaphore = Arc::new(tokio::sync::Semaphore::new(OWNER_SUMMARY_CONCURRENCY));
186        let mut handles = Vec::with_capacity(schemas.len());
187
188        for model in schemas {
189            let owner_id = owner_id.to_string();
190            let owner_type = owner_type.to_string();
191            let v = v.clone();
192            let permit = semaphore
193                .clone()
194                .acquire_owned()
195                .await
196                .map_err(|e| Error::Internal(e.to_string()))?;
197            handles.push(tokio::spawn(async move {
198                let _permit = permit;
199                let active = Self::count_ownership_rows_for_schema(
200                    model,
201                    &owner_id,
202                    &owner_type,
203                    "active",
204                    &v,
205                )
206                .await?;
207                let pending = Self::count_ownership_rows_for_schema(
208                    model,
209                    &owner_id,
210                    &owner_type,
211                    "pending_deletion",
212                    &v,
213                )
214                .await?;
215                Ok::<_, Error>(OwnerSchemaRowCount {
216                    valence_model: model.to_string(),
217                    active_rows: active,
218                    pending_deletion_rows: pending,
219                })
220            }));
221        }
222
223        let mut rows_by_schema = Vec::new();
224        for handle in handles {
225            let row = handle.await.map_err(|e| Error::Internal(e.to_string()))??;
226            if row.active_rows > 0 || row.pending_deletion_rows > 0 {
227                rows_by_schema.push(row);
228            }
229        }
230
231        rows_by_schema.sort_by(|a, b| {
232            b.active_rows
233                .cmp(&a.active_rows)
234                .then_with(|| a.valence_model.cmp(&b.valence_model))
235        });
236
237        let owned_rows: u64 = rows_by_schema.iter().map(|r| r.active_rows).sum();
238        let pending_deletion_rows: u64 =
239            rows_by_schema.iter().map(|r| r.pending_deletion_rows).sum();
240        let tables_with_data = rows_by_schema.iter().filter(|r| r.active_rows > 0).count() as u64;
241
242        Ok(OwnerDataSummary {
243            owned_rows,
244            tables_with_data,
245            pending_deletion_rows,
246            rows_by_schema,
247        })
248    }
249
250    /// Recent transfer rows for the ownership row of `valence_model` / `record_id`.
251    pub async fn transfer_history(
252        valence_model: &str,
253        record_id: &str,
254        v: &Valence,
255        limit: u32,
256    ) -> Result<Vec<Value>> {
257        let oid = ownership_row_id(valence_model, record_id);
258        let sys = system_valence(v);
259        let q = QueryCore::new("valence_ownership_transfer".to_string())
260            .select(vec!["*".to_string()])
261            .where_record(
262                "ownership_id".to_string(),
263                RecordPredicate::Equals(crate::RecordId::new("valence_data_ownership", &oid)),
264            )
265            .order_by("transferred_at".to_string(), SortDirection::Desc)
266            .limit(limit);
267        let rows: Vec<Value> = q.execute(&sys).await?;
268        Ok(rows)
269    }
270}