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