1use crate::dataset_admission::ReceiptRow;
2use crate::{MetadataExecutor, PgMetaError, PgMetadataRepository};
3use ahri_tre_core::{CoreError, OperationRepository, StoredOperation};
4use ahri_tre_observability::{FailureCategory, LifecycleSpan, Outcome, Phase};
5use ahri_tre_protocol::operation::OperationStatus;
6use uuid::Uuid;
7
8impl PgMetadataRepository<'static> {
9 pub fn operation_stage_authorized(&self, id: Uuid) -> Result<bool, CoreError> {
12 self.with_session_connection(
13 "authorize operation stage",
14 |c| stage_authorized(c, id),
15 |c| stage_authorized(c, id),
16 )
17 }
18
19 pub fn operation_control(
22 connection: ahri_tre_libpq_oauth::LibpqOAuthConnection,
23 ) -> Result<Self, CoreError> {
24 Ok(Self::from_oauth_shared(std::sync::Arc::new(
25 std::sync::Mutex::new(connection),
26 )))
27 }
28}
29
30impl OperationRepository for PgMetadataRepository<'_> {
31 fn create_operation(
32 &self,
33 attempt_id: Uuid,
34 record: &StoredOperation,
35 ) -> Result<(), CoreError> {
36 let summary = &record.detail.summary;
37 if summary.status != OperationStatus::Pending
38 || !summary.operation_kind.is_supported()
39 || record.detail.events.items.len() != 1
40 || record.detail.events.items[0].event_sequence != 1
41 || record.result.is_some()
42 {
43 return Err(unavailable());
44 }
45 let id = operation_id(record)?;
46 let body = serde_json::to_string(record).map_err(|_| unavailable())?;
47 self.with_session_connection(
48 "accept operation",
49 |c| create(c, id, attempt_id, &body),
50 |c| create(c, id, attempt_id, &body),
51 )
52 .map_err(|_| unavailable())
53 }
54 fn get_operation(&self, id: Uuid) -> Result<Option<StoredOperation>, CoreError> {
55 self.with_session_connection("inspect operation", |c| read(c, id), |c| read(c, id))
56 .map_err(|_| unavailable())
57 }
58 fn operation_history(
59 &self,
60 query: &ahri_tre_core::OperationHistoryQuery,
61 ) -> Result<Vec<ahri_tre_core::OperationHistoryEntry>, CoreError> {
62 self.with_session_connection(
63 "browse operation history",
64 |c| history(c, query),
65 |c| history(c, query),
66 )
67 .map_err(|_| unavailable())
68 }
69 fn transition_operation(
70 &self,
71 expected: OperationStatus,
72 sequence: u64,
73 record: &StoredOperation,
74 ) -> Result<(), CoreError> {
75 if !(expected.can_transition_to(record.detail.summary.status)
76 || (expected == record.detail.summary.status && !expected.is_terminal()))
77 || record.detail.events.items.last().map(|e| e.event_sequence) != Some(sequence + 1)
78 {
79 return Err(unavailable());
80 }
81 let id = operation_id(record)?;
82 let body = serde_json::to_string(record).map_err(|_| unavailable())?;
83 let expected = serde_json::to_value(expected).map_err(|_| unavailable())?;
84 let expected = expected.as_str().ok_or_else(unavailable)?;
85 let sequence = i64::try_from(sequence).map_err(|_| unavailable())?;
86 let changed = self
87 .with_session_connection(
88 "transition operation",
89 |c| transition(c, id, expected, sequence, &body),
90 |c| transition(c, id, expected, sequence, &body),
91 )
92 .map_err(|_| unavailable())?;
93 if changed != 1 {
94 return Err(CoreError::Conflict("Operation state changed".into()));
95 }
96 Ok(())
97 }
98}
99fn stage_authorized<E: MetadataExecutor>(c: &mut E, id: Uuid) -> Result<bool, PgMetaError>
100where
101 E::Row: ReceiptRow,
102{
103 let row = c.query_optional("authorize operation stage", "SELECT public.tre_user_has_study_access((record#>>'{scope,study_id}')::uuid)::text FROM public.operations WHERE operation_id=$1 AND actor=CURRENT_USER", &[&id])?;
104 Ok(row.is_some_and(|row| row.text(0).is_ok_and(|value| value == "true")))
105}
106
107fn operation_id(record: &StoredOperation) -> Result<Uuid, CoreError> {
108 let reference = record.detail.summary.operation;
109 if reference.kind != ahri_tre_protocol::refs::ObjectKind::Operation
110 || reference.datastore_id.as_uuid() != record.scope.datastore_id
111 {
112 return Err(unavailable());
113 }
114 Ok(reference.id.as_uuid())
115}
116
117fn unavailable() -> CoreError {
118 CoreError::Infrastructure("Operation persistence is unavailable".into())
119}
120fn create(
121 c: &mut impl MetadataExecutor,
122 id: Uuid,
123 attempt: Uuid,
124 body: &str,
125) -> Result<(), PgMetaError> {
126 let changed=c.execute_command(
127 "accept operation",
128 "INSERT INTO public.operations(operation_id, attempt_id, created_at, record) SELECT $1,a.attempt_id,($3::text::jsonb#>>'{detail,summary,timestamps,created_at}')::timestamptz,$3::text::jsonb FROM public.dataset_output_attempts a JOIN public.dataset_output_reservations r USING(attempt_id) WHERE a.attempt_id=$2 AND a.actor=CURRENT_USER AND NOT a.admitted AND a.study_id=($3::text::jsonb->'scope'->>'study_id')::uuid",
129 &[&id, &attempt, &body],
130 )?;
131 if changed != 1 {
132 return Err(PgMetaError::Decode {
133 field: "operation",
134 value: "<unavailable>".into(),
135 message: "acceptance does not own its output".into(),
136 });
137 }
138 Ok(())
139}
140fn read<E: MetadataExecutor>(c: &mut E, id: Uuid) -> Result<Option<StoredOperation>, PgMetaError>
141where
142 E::Row: ReceiptRow,
143{
144 c.query_optional(
145 "inspect operation",
146 "SELECT record::text FROM public.operations WHERE operation_id=$1 AND actor=CURRENT_USER",
147 &[&id],
148 )?
149 .map(|r| decode_operation_record(r.text(0)?))
150 .transpose()
151}
152fn decode_operation_record(text: &str) -> Result<StoredOperation, PgMetaError> {
155 use ahri_tre_protocol::{
156 PublicUuid,
157 refs::{ObjectKind, ObjectRef},
158 };
159 let mut value: serde_json::Value = serde_json::from_str(text).map_err(|_| recovery_error())?;
160 if value
161 .pointer("/result/data/lineage/transformation/transformation_id")
162 .is_some_and(serde_json::Value::is_i64)
163 {
164 let datastore: uuid::Uuid = serde_json::from_value(
165 value
166 .pointer("/scope/datastore_id")
167 .cloned()
168 .ok_or_else(recovery_error)?,
169 )
170 .map_err(|_| recovery_error())?;
171 let lineage: ahri_tre_types::TransformationLineageRecord = serde_json::from_value(
172 value
173 .pointer("/result/data/lineage")
174 .cloned()
175 .ok_or_else(recovery_error)?,
176 )
177 .map_err(|_| recovery_error())?;
178 let version_ref = |id: ahri_tre_types::VersionId| ObjectRef {
179 datastore_id: PublicUuid::from_uuid(datastore),
180 kind: ObjectKind::AssetVersion,
181 id: PublicUuid::from_uuid(id.0),
182 };
183 let projected = ahri_tre_protocol::transformation::TransformationSummary {
184 transformation: ahri_tre_protocol::transformation::TransformationDetail {
185 transformation_id: ObjectRef {
186 datastore_id: PublicUuid::from_uuid(datastore),
187 kind: ObjectKind::Transformation,
188 id: ahri_tre_protocol::refs::encode_scoped_integer_ref(
189 "transformation",
190 lineage.transformation.transformation_id.0,
191 ),
192 },
193 transformation_type: lineage.transformation.transformation_type,
194 description: lineage.transformation.description,
195 date_created: lineage.transformation.date_created,
196 created_by: lineage.transformation.created_by,
197 },
198 inputs: lineage
199 .inputs
200 .into_iter()
201 .map(|i| version_ref(i.version_id))
202 .collect(),
203 outputs: lineage
204 .outputs
205 .into_iter()
206 .map(|i| version_ref(i.version_id))
207 .collect(),
208 };
209 value["result"]["data"]["lineage"] =
210 serde_json::to_value(projected).map_err(|_| recovery_error())?;
211 }
212 serde_json::from_value(value).map_err(|_| recovery_error())
213}
214
215fn transition(
216 c: &mut impl MetadataExecutor,
217 id: Uuid,
218 expected: &str,
219 sequence: i64,
220 body: &str,
221) -> Result<u64, PgMetaError> {
222 c.execute_command("transition operation", "UPDATE public.operations SET record=$4::text::jsonb, event_sequence=event_sequence+1 WHERE operation_id=$1 AND actor=CURRENT_USER AND record->'detail'->'summary'->>'status'=$2 AND event_sequence=$3 AND record->'scope'=$4::text::jsonb->'scope' AND record#>'{detail,summary,started_by_request_id}'=$4::text::jsonb#>'{detail,summary,started_by_request_id}' AND record#>'{detail,summary,timestamps,created_at}'=$4::text::jsonb#>'{detail,summary,timestamps,created_at}' AND record#>'{detail,events,items}'=($4::text::jsonb#>'{detail,events,items}') - $3::int AND record->'detail'->'summary'->>'status' NOT IN ('completed','failed','cancelled')", &[&id,&expected,&sequence,&body])
223}
224
225impl crate::PgDatasetMaintenance {
229 pub fn verify_operation_recovery(&self) -> Result<(), CoreError> {
230 self.repository.with_session_connection(
231 "verify operation recovery",
232 verify_recovery,
233 verify_recovery,
234 )
235 }
236
237 pub fn reconcile_operations(
238 &self,
239 executor: std::sync::Arc<dyn ahri_tre_core::DatasetExecutorLease>,
240 interrupt: impl Fn(&mut StoredOperation),
241 ) -> Result<(), CoreError> {
242 self.verify_operation_recovery()?;
243 let candidates = self.repository.with_session_connection(
244 "inspect operation recovery",
245 recovery_candidates,
246 recovery_candidates,
247 )?;
248 for (attempt, owner, operation_id) in candidates {
249 if !executor.can_recover(owner, attempt) {
250 continue;
251 }
252 let Ok(_exclusion) = executor.clone().begin_attempt(attempt) else {
253 continue;
254 };
255 let recovery =
256 LifecycleSpan::start_operation(Phase::Recovery, None, Some(operation_id));
257 let result = self.repository.with_session_connection(
258 "reconcile operation",
259 |c| reconcile(c, attempt, owner, &interrupt),
260 |c| reconcile(c, attempt, owner, &interrupt),
261 );
262 match &result {
263 Ok(true) => {
264 recovery.finish(Outcome::Interrupted, Some(FailureCategory::Interrupted))
265 }
266 Ok(false) => recovery.finish(Outcome::Skipped, None),
267 Err(_) => recovery.finish(Outcome::Unavailable, Some(FailureCategory::Metadata)),
268 }
269 result?;
270 }
271 self.prune_expired_operations()?;
272 Ok(())
273 }
274
275 pub fn prune_expired_operations(&self) -> Result<u64, CoreError> {
279 self.verify_operation_recovery()?;
280 self.repository.with_session_connection(
281 "prune expired operations",
282 prune_expired,
283 prune_expired,
284 )
285 }
286}
287
288fn prune_expired(c: &mut impl MetadataExecutor) -> Result<u64, PgMetaError> {
289 c.execute_command("prune expired operations", r#"
290 WITH expired AS (
291 SELECT o.operation_id FROM public.operations o
292 WHERE o.record#>>'{detail,summary,status}' IN ('completed','failed','cancelled')
293 AND public.operation_now() >= (o.record#>>'{detail,summary,retention,operation_expires_at}')::timestamptz
294 AND NOT EXISTS (
295 SELECT 1 FROM public.operation_idempotency k
296 WHERE k.operation_id=o.operation_id
297 AND public.operation_now() < k.accepted_at + INTERVAL '24 hours'
298 )
299 ORDER BY o.created_at, o.operation_id
300 LIMIT 100 FOR UPDATE OF o SKIP LOCKED
301 )
302 DELETE FROM public.operations o USING expired e WHERE o.operation_id=e.operation_id
303 "#, &[])
304}
305
306fn verify_recovery<E: MetadataExecutor>(c: &mut E) -> Result<(), PgMetaError>
307where
308 E::Row: ReceiptRow,
309{
310 let row = c.query_optional("verify operation recovery", "SELECT (count(*)=2 AND bool_and(r.rolsuper OR (has_table_privilege(c.oid, 'SELECT') AND (c.relname <> 'operations' OR (has_table_privilege(c.oid, 'UPDATE') AND has_table_privilege(c.oid, 'DELETE'))) AND (r.rolbypassrls OR (NOT c.relforcerowsecurity AND pg_has_role(c.relowner, 'USAGE'))))))::text FROM pg_class c CROSS JOIN pg_roles r WHERE r.rolname = CURRENT_USER AND c.oid IN (to_regclass('public.operations'), to_regclass('public.operation_idempotency'))", &[])?;
311 if row.is_some_and(|r| r.text(0).is_ok_and(|v| v == "true")) {
312 Ok(())
313 } else {
314 Err(recovery_error())
315 }
316}
317
318fn recovery_candidates<E: MetadataExecutor>(
319 c: &mut E,
320) -> Result<Vec<(Uuid, ahri_tre_core::DatasetExecutorIdentity, Uuid)>, PgMetaError>
321where
322 E::Row: ReceiptRow,
323{
324 c.query_many("inspect operation recovery", "SELECT a.attempt_id::text, a.coordinator_id::text, a.generation_id::text, o.operation_id::text FROM public.dataset_output_attempts a JOIN public.operations o USING(attempt_id) WHERE o.record#>>'{detail,summary,status}' NOT IN ('completed','failed','cancelled') OR (NOT a.admitted AND NOT EXISTS (SELECT 1 FROM public.dataset_output_reservations r WHERE r.attempt_id=a.attempt_id)) ORDER BY a.created_at, a.attempt_id", &[])?.into_iter().map(|r| Ok((
325 r.text(0)?.parse().map_err(|_| recovery_error())?,
326 ahri_tre_core::DatasetExecutorIdentity {
327 coordinator_id: r.text(1)?.parse().map_err(|_| recovery_error())?,
328 generation_id: r.text(2)?.parse().map_err(|_| recovery_error())?,
329 },
330 r.text(3)?.parse().map_err(|_| recovery_error())?
331 ))).collect()
332}
333
334fn reconcile<E: MetadataExecutor>(
335 c: &mut E,
336 attempt: Uuid,
337 owner: ahri_tre_core::DatasetExecutorIdentity,
338 interrupt: &impl Fn(&mut StoredOperation),
339) -> Result<bool, PgMetaError>
340where
341 E::Row: ReceiptRow,
342{
343 c.with_transaction("reconcile stopped operation", |t| {
344 t.query_optional("lock operation reservation", "SELECT attempt_id FROM public.dataset_output_reservations WHERE attempt_id=$1 FOR UPDATE", &[&attempt])?;
346 let row = t.query_optional("inspect stopped operation", "SELECT o.record::text FROM public.operations o JOIN public.dataset_output_attempts a USING(attempt_id) WHERE a.attempt_id=$1 AND a.coordinator_id=$2 AND a.generation_id=$3 AND NOT a.admitted AND NOT EXISTS (SELECT 1 FROM public.dataset_admission_receipts d WHERE d.version_id=a.version_id AND d.admitted) FOR UPDATE OF o, a", &[&attempt, &owner.coordinator_id, &owner.generation_id])?;
347 let Some(row) = row else { return Ok(false); };
348 let mut record: StoredOperation = decode_operation_record(row.text(0)?)?;
349 if record.detail.summary.status.is_terminal() {
350 t.execute_command("remove reconciled clean attempt", "DELETE FROM public.dataset_output_attempts a WHERE a.attempt_id=$1 AND NOT EXISTS (SELECT 1 FROM public.dataset_output_reservations r WHERE r.attempt_id=a.attempt_id)", &[&attempt])?;
351 return Ok(false);
352 }
353 let before = record.clone();
354 interrupt(&mut record);
355 let sequence = before.detail.events.items.len();
357 if record.scope != before.scope || record.result != before.result
358 || record.detail.summary.status != OperationStatus::Failed
359 || record.detail.summary.operation != before.detail.summary.operation
360 || record.detail.summary.started_by_request_id != before.detail.summary.started_by_request_id
361 || record.detail.summary.timestamps.created_at != before.detail.summary.timestamps.created_at
362 || record.detail.events.items.len() != sequence + 1
363 || record.detail.events.items[..sequence] != before.detail.events.items
364 || record.detail.events.items[sequence].event_sequence != sequence as u64 + 1 {
365 return Err(recovery_error());
366 }
367 let body = serde_json::to_string(&record).map_err(|_| recovery_error())?;
368 t.execute_command("interrupt stopped operation", "UPDATE public.operations SET record=$2::text::jsonb, event_sequence=event_sequence+1 WHERE attempt_id=$1", &[&attempt, &body])?;
369 t.execute_command("remove reconciled clean attempt", "DELETE FROM public.dataset_output_attempts a WHERE a.attempt_id=$1 AND NOT EXISTS (SELECT 1 FROM public.dataset_output_reservations r WHERE r.attempt_id=a.attempt_id)", &[&attempt])?;
372 Ok(true)
373 })
374}
375fn recovery_error() -> PgMetaError {
376 PgMetaError::Decode {
377 field: "operation",
378 value: "<unavailable>".into(),
379 message: "Operation recovery is unavailable".into(),
380 }
381}
382
383pub struct PgOperationKey {
387 repository: PgMetadataRepository<'static>,
388 scope_digest: String,
389 intent_digest: String,
390}
391
392impl PgMetadataRepository<'static> {
393 pub fn lock_operation_key(
394 &self,
395 scope_digest: String,
396 intent_digest: String,
397 ) -> Result<PgOperationKey, CoreError> {
398 self.with_session_connection(
399 "serialize operation admission",
400 |c| key_lock(c, &scope_digest, true),
401 |c| key_lock(c, &scope_digest, true),
402 )?;
403 Ok(PgOperationKey {
404 repository: self.clone(),
405 scope_digest,
406 intent_digest,
407 })
408 }
409}
410
411impl PgOperationKey {
412 pub fn retained(&self) -> Result<Option<(StoredOperation, bool)>, CoreError> {
415 self.repository.with_session_connection(
416 "recover idempotent operation",
417 |c| retained_key(c, &self.scope_digest, &self.intent_digest),
418 |c| retained_key(c, &self.scope_digest, &self.intent_digest),
419 )
420 }
421
422 pub fn bind(
425 &self,
426 transaction: &PgMetadataRepository<'_>,
427 operation: Uuid,
428 ) -> Result<(), CoreError> {
429 transaction.with_session_connection(
430 "bind operation key",
431 |c| bind_key(c, &self.scope_digest, &self.intent_digest, operation),
432 |c| bind_key(c, &self.scope_digest, &self.intent_digest, operation),
433 )
434 }
435}
436
437impl Drop for PgOperationKey {
438 fn drop(&mut self) {
439 let _ = self.repository.with_session_connection(
440 "release operation admission",
441 |c| key_lock(c, &self.scope_digest, false),
442 |c| key_lock(c, &self.scope_digest, false),
443 );
444 }
445}
446
447fn key_lock<E: MetadataExecutor>(
448 c: &mut E,
449 digest: &str,
450 acquire: bool,
451) -> Result<(), PgMetaError> {
452 c.query_optional(
453 "serialize operation admission",
454 if acquire {
455 "SELECT pg_advisory_lock(hashtextextended(CURRENT_USER || $1, 606))"
456 } else {
457 "SELECT pg_advisory_unlock(hashtextextended(CURRENT_USER || $1, 606))"
458 },
459 &[&digest],
460 )?;
461 Ok(())
462}
463
464fn retained_key<E: MetadataExecutor>(
465 c: &mut E,
466 scope: &str,
467 intent: &str,
468) -> Result<Option<(StoredOperation, bool)>, PgMetaError>
469where
470 E::Row: ReceiptRow,
471{
472 c.query_optional("recover idempotent operation", "SELECT o.record::text, (k.intent_digest=$2)::text FROM public.operation_idempotency k JOIN public.operations o USING(operation_id) WHERE k.actor=CURRENT_USER AND o.actor=CURRENT_USER AND k.scope_digest=$1 AND (o.record#>>'{detail,summary,status}' NOT IN ('completed','failed','cancelled') OR o.record#>>'{detail,summary,retention,operation_expires_at}' IS NULL OR public.operation_now() < GREATEST(k.accepted_at + INTERVAL '24 hours', (o.record#>>'{detail,summary,retention,operation_expires_at}')::timestamptz))", &[&scope, &intent])?.map(|row| {
473 let record = decode_operation_record(row.text(0)?)?;
474 Ok((record, row.text(1)? == "true"))
475 }).transpose()
476}
477
478fn bind_key(
479 c: &mut impl MetadataExecutor,
480 scope: &str,
481 intent: &str,
482 operation: Uuid,
483) -> Result<(), PgMetaError> {
484 let changed = c.execute_command("bind operation key", "INSERT INTO public.operation_idempotency(scope_digest,intent_digest,operation_id) VALUES($1,$2,$3) ON CONFLICT(actor,scope_digest) DO UPDATE SET intent_digest=EXCLUDED.intent_digest, operation_id=EXCLUDED.operation_id, accepted_at=public.operation_now() WHERE EXISTS (SELECT 1 FROM public.operations o WHERE o.operation_id=operation_idempotency.operation_id AND o.record#>>'{detail,summary,status}' IN ('completed','failed','cancelled') AND o.record#>>'{detail,summary,retention,operation_expires_at}' IS NOT NULL AND public.operation_now() >= GREATEST(operation_idempotency.accepted_at + INTERVAL '24 hours', (o.record#>>'{detail,summary,retention,operation_expires_at}')::timestamptz))", &[&scope, &intent, &operation])?;
485 if changed != 1 {
486 return Err(recovery_error());
487 }
488 Ok(())
489}
490
491impl PgMetadataRepository<'_> {
492 pub fn operation_precondition_satisfied(
494 &self,
495 scope: &ahri_tre_core::OperationResourceScope,
496 precondition: &ahri_tre_protocol::mutation::Precondition,
497 ) -> Result<bool, CoreError> {
498 self.with_session_connection(
499 "validate operation precondition",
500 |c| check_precondition(c, scope, precondition),
501 |c| check_precondition(c, scope, precondition),
502 )
503 }
504}
505
506fn check_precondition<E: MetadataExecutor>(
507 c: &mut E,
508 scope: &ahri_tre_core::OperationResourceScope,
509 precondition: &ahri_tre_protocol::mutation::Precondition,
510) -> Result<bool, PgMetaError>
511where
512 E::Row: ReceiptRow,
513{
514 use ahri_tre_protocol::mutation::Precondition;
515 let exists = |c: &mut E| -> Result<bool, PgMetaError> {
516 Ok(c.query_optional("inspect Dataset precondition", "SELECT asset_id FROM public.assets WHERE study_id=$1 AND name=$2 AND asset_type='dataset'", &[&scope.study_id.0, &scope.dataset_name])?.is_some())
517 };
518 match precondition {
519 Precondition::Exists { target } if target == "dataset" => exists(c),
520 Precondition::DoesNotExist { target } if target == "dataset" => Ok(!exists(c)?),
521 Precondition::Exists { target } if target == "source" => Ok(c.query_optional("inspect source precondition", "SELECT datafile_id FROM public.datafiles WHERE datafile_id=$1", &[&scope.source_version_id.0])?.is_some()),
522 Precondition::ResourceVersion { target, version } if target == "source" => Ok(c.query_optional("inspect source version precondition", "SELECT version_id FROM public.asset_versions WHERE version_id=$1 AND major::text || '.' || minor::text || '.' || patch::text=$2", &[&scope.source_version_id.0, &version.trim_start_matches('v')])?.is_some()),
523 _ => Err(recovery_error()),
524 }
525}
526
527fn history<E: MetadataExecutor>(
528 c: &mut E,
529 query: &ahri_tre_core::OperationHistoryQuery,
530) -> Result<Vec<ahri_tre_core::OperationHistoryEntry>, PgMetaError>
531where
532 E::Row: ReceiptRow,
533{
534 let lower = query.created_after.map(|value| value.to_rfc3339());
535 let before = query.created_before.map(|value| value.to_rfc3339());
536 let upper = query.upper_bound.to_rfc3339();
537 let after_time = query
538 .after
539 .as_ref()
540 .map(|value| value.created_at.to_rfc3339());
541 let after_id = query.after.as_ref().map(|value| value.id);
542 let limit = i64::from(query.limit.clamp(1, 500));
543 c.query_many("browse operation history", r#"
544 SELECT jsonb_build_object('scope', record->'scope', 'summary', record#>'{detail,summary}', 'result_dataset_id', record#>'{result,data,dataset,dataset_id}')::text
545 FROM public.operations
546 WHERE actor=CURRENT_USER::text COLLATE "default"
547 AND ($1::text IS NULL OR created_at > $1::text::timestamptz)
548 AND ($2::text IS NULL OR created_at < $2::text::timestamptz)
549 AND created_at <= $3::text::timestamptz
550 AND ($4::text IS NULL OR (created_at, operation_id) < ($4::text::timestamptz, $5::uuid))
551 ORDER BY created_at DESC, operation_id DESC
552 LIMIT $6
553 "#, &[&lower, &before, &upper, &after_time, &after_id, &limit])?
554 .into_iter().map(|row| serde_json::from_str(row.text(0)?).map_err(|_| recovery_error())).collect()
555}
556
557impl PgMetadataRepository<'_> {
558 pub fn operation_now(&self) -> Result<chrono::DateTime<chrono::Utc>, CoreError> {
560 self.with_session_connection("read operation clock", operation_now, operation_now)
561 }
562}
563fn operation_now<E: MetadataExecutor>(
564 c: &mut E,
565) -> Result<chrono::DateTime<chrono::Utc>, PgMetaError>
566where
567 E::Row: ReceiptRow,
568{
569 let row = c
570 .query_optional(
571 "read operation clock",
572 r#"SELECT to_char(public.operation_now() AT TIME ZONE 'UTC', 'YYYY-MM-DD"T"HH24:MI:SS.US"Z"')"#,
573 &[],
574 )?
575 .ok_or_else(recovery_error)?;
576 row.text(0)?.parse().map_err(|_| recovery_error())
577}