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(
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 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 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 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 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}