Skip to main content

ahri_tre_pgmeta/
operation.rs

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    /// Recheck current Study authority before another governed stage. The
10    /// operation record is correlation evidence and never supplies the grant.
11    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    /// A fresh user-authenticated lifecycle connection. The application exposes
20    /// only inspection and cancellation; it shares no mutex with the producer.
21    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}
152/// Stored m0002 results predate canonical public lineage. Decode their private
153/// record using the retained Datastore scope; never guess it from a caller.
154fn 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
225/// Operation recovery is restricted to stopped private attempts under the same
226/// Dataset executor exclusion used by Lake cleanup. It can only fail unfinished
227/// records; committed completion/result evidence is never rewritten.
228impl 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    /// A bounded lifecycle pass. The only cascade is into expired key bindings;
276    /// Dataset ownership, private attempts, governed Assets and provenance are
277    /// independent tables and remain available to trusted reconciliation.
278    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        // Order after any in-flight final admission before inspecting its receipt.
345        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        // This capability cannot change identity, result, prior events, or claim success.
356        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        // A normal cleanup may have released the reservation before its user
370        // connection lost the terminal response. Keep evidence until this point.
371        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
383/// Serializes only admission for one actor-scoped key on a dedicated user
384/// connection. Dropping it after committed acceptance lets retries inspect the
385/// operation without waiting for Lake execution or occupying a worker permit.
386pub 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    /// An absent deadline protects active operations indefinitely. Terminal
413    /// bindings last through both the 24-hour minimum and operation retention.
414    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    /// Must join the reservation/operation transaction. An expired binding can
423    /// be replaced only after the new request has acquired output ownership.
424    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    /// Evaluated inside output acceptance after ownership has been acquired.
493    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    /// One authoritative lifecycle instant, also used by key and cleanup queries.
559    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}