valence_core/ownership/service/
read.rs1use 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 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 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 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 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 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}