Skip to main content

IntegrationEventRepository

Struct IntegrationEventRepository 

Source
pub struct IntegrationEventRepository(/* private fields */);
Expand description

Repository for IntegrationEvent entities.

All standard CRUD, soft-delete, pagination, and bulk methods are provided automatically via Deref to backbone_orm::GenericCrudRepository.

Implementations§

Source§

impl IntegrationEventRepository

Source

pub fn new(pool: PgPool) -> Self

Create a new repository instance.

Source§

impl IntegrationEventRepository

Hand-written IntegrationEvent SQL. Lives here (not in the write service) per the module’s 4-layer rule: services orchestrate and own the unit of work, repositories hold the SQL.

Source

pub async fn claim_event( &self, pool: &PgPool, e: &NewEventRow<'_>, ) -> Result<Option<Uuid>, Error>

Claim the (connector, business_key) dedup slot. Ok(None) = this business action was already received (a webhook retry, or a second notification for the same action) and the caller must re-read the original rather than re-map it. A NEW business action does not conflict.

Runs outside any transaction; it rides the request-dedicated connection when the composing service bound one — the decorator’s org fence then keeps the dedup idempotent within the tenant even off the request path — else runs plainly on the pool (a decorated deployment refuses it, fail closed; ADR-0029).

Source

pub async fn fetch_by_business_key( &self, pool: &PgPool, connector_id: Uuid, business_key: &str, ) -> Result<EventOutcomeRow, Error>

Re-read the original after a losing dedup claim — the row is known to exist. Same connection discipline as Self::claim_event.

Source

pub async fn mark_mapped( &self, conn: &mut PgConnection, event_id: Uuid, mapped_ref_type: &str, mapped_ref_id: Uuid, ) -> Result<(), Error>

Record a successful mapping to an internal action. State-guarded on received.

Takes the CALLER’S connection so this and the outbox stage commit as one unit. The caller has already bound the ambient org scope on it when one is present — don’t re-bind here.

Source

pub async fn mark_ignored( &self, conn: &mut PgConnection, event_id: Uuid, reason: &str, ) -> Result<(), Error>

Record an event the target intentionally ignored. State-guarded on received; same caller-owned-tx contract as Self::mark_mapped.

Source

pub async fn mark_failed( &self, conn: &mut PgConnection, event_id: Uuid, error_detail: &str, ) -> Result<(), Error>

Record a mapping rejection. State-guarded on received; same caller-owned-tx contract as Self::mark_mapped.

Source

pub async fn fetch_failed( &self, pool: &PgPool, connector_id: Uuid, ) -> Result<Vec<FailedEventRow>, Error>

A connector’s FAILED events — the operator’s failure report and the retry loop’s work list.

ID-only: the connector id alone identifies the set, so the read rides the request-dedicated connection when the composing service bound one — the decorator’s org fence then hides another tenant’s connector — else runs plainly on the pool (ADR-0029).

Source

pub async fn retry_mark_mapped( &self, conn: &mut PgConnection, event_id: Uuid, mapped_ref_type: &str, mapped_ref_id: Uuid, ) -> Result<u64, Error>

The RETRY path’s failed → mapped transition, clearing the stale error. Distinct from Self::mark_mapped: state-guarded on failed, not received. Returns the rows affected so the caller only stages the event (and counts it) when the transition really happened.

Takes the CALLER’S connection so the transition and the outbox stage commit as one unit — and so the caller can roll back when it loses. The caller has already bound the ambient org scope on it when one is present — don’t re-bind here.

Source

pub async fn retry_mark_ignored( &self, pool: &PgPool, event_id: Uuid, reason: &str, ) -> Result<(), Error>

The RETRY path’s failed → ignored transition (the target now says this event is a no-op). State-guarded on failed.

Runs outside any transaction; it rides the request-dedicated connection when the composing service bound one, else runs plainly on the pool (ADR-0029).

Source

pub async fn set_error_detail( &self, pool: &PgPool, event_id: Uuid, error_detail: &str, ) -> Result<(), Error>

Refresh a still-failing event’s error after a retry attempt — the status deliberately stays failed so the next retry picks it up again. Same connection discipline as Self::retry_mark_ignored.

Methods from Deref<Target = GenericCrudRepository<IntegrationEvent, SoftDelete>>§

Source

pub fn pool(&self) -> &Pool<Postgres>

Source

pub fn table_name(&self) -> &str

Source

pub async fn create(&self, entity: &T) -> Result<T, Error>
where T: Serialize + Send + Sync,

Insert a new entity row and return the created row.

Source

pub async fn bulk_create(&self, entities: &[T]) -> Result<Vec<T>, Error>
where T: Serialize + Send + Sync,

Insert multiple entity rows inside a single transaction.

Source

pub async fn run_filtered_query( &self, pagination: PaginationParams, base_condition: Option<&str>, filters: &HashMap<String, String>, column_types: &HashMap<String, String>, search_fields: &[&str], ) -> Result<PaginatedResult<T>, Error>
where T: Send + Sync,

Execute a filtered / paginated query against this entity’s table.

base_condition — when Some, it is inserted as the __base_condition filter key which the ORM injects verbatim into the WHERE clause. Use this to add soft-delete guards without touching the caller-supplied filters.

NOTE: prefer the *_scoped entry points for anything a client can reach — they compose the company fence into this condition. This raw form applies no fence and is for internal/job callers that have already established scope.

Source

pub async fn run_aggregate_query( &self, spec: &AggregateSpec, base_condition: Option<&str>, filters: &HashMap<String, String>, column_types: &HashMap<String, String>, search_fields: &[&str], ) -> Result<AggregateResult, Error>
where T: EntityRepoMeta + Send + Sync,

The aggregate twin of Self::run_filtered_query.

Takes the same base_condition so an aggregate is computed over exactly the row set the matching list would return — a total that counted soft-deleted rows, or another tenant’s, would be wrong in a way no caller could see.

Source

pub async fn find_by_text_field( &self, field: &str, value: &str, ) -> Result<Option<T>, Error>

Find an active entity by a unique text field.

Filters metadata->>'deleted_at' IS NULL automatically.

Source

pub async fn exists_by_text_field( &self, field: &str, value: &str, ) -> Result<bool, Error>

Check existence by a unique text field (active records only).

Source

pub async fn find_by_uuid_field( &self, field: &str, value: Uuid, ) -> Result<Option<T>, Error>

Find an active entity by a unique UUID field.

Source

pub async fn exists_by_uuid_field( &self, field: &str, value: Uuid, ) -> Result<bool, Error>

Check existence by a unique UUID field (active records only).

Source

pub async fn list_paginated_filtered( &self, pagination: PaginationParams, filters: Option<&HashMap<String, String>>, ) -> Result<PaginatedResult<T>, Error>
where T: EntityRepoMeta + Send + Sync,

Paginate active entities with filter and search support.

Requires T: EntityRepoMeta for column type hints and search fields.

Source

pub async fn aggregate_filtered( &self, spec: &AggregateSpec, filters: Option<&HashMap<String, String>>, ) -> Result<AggregateResult, Error>
where T: EntityRepoMeta + Send + Sync,

Group and reduce active entities, skipping the soft-deleted.

Source

pub async fn list_paginated_filtered_scoped( &self, pagination: PaginationParams, filters: Option<&HashMap<String, String>>, company: Option<Uuid>, ) -> Result<PaginatedResult<T>, Error>
where T: EntityRepoMeta + Send + Sync,

Paginate active entities, fenced to company.

The company-scoped counterpart of Self::list_paginated_filtered and the one any client-reachable read should use. Fails closed: a fenced entity with no company returns MissingCompanyScope rather than every company’s rows.

The fence is ANDed into the same __base_condition the soft-delete guard uses, so it lands in SQL — which is what keeps total honest. A post-filter in the handler would fetch limit rows and return the survivors, leaving callers unable to tell “end of data” from “filtered”, and could not fix COUNT at all.

Source

pub async fn list_deleted_filtered_scoped( &self, pagination: PaginationParams, filters: Option<&HashMap<String, String>>, company: Option<Uuid>, ) -> Result<PaginatedResult<T>, Error>
where T: EntityRepoMeta + Send + Sync,

Paginate soft-deleted entities, fenced to company.

/trash is the worse leak of the two: it serves rows whose owners were told the data is gone. Same fence, same fail-closed contract.

Source

pub async fn list_deleted_filtered( &self, pagination: PaginationParams, filters: Option<&HashMap<String, String>>, ) -> Result<PaginatedResult<T>, Error>
where T: EntityRepoMeta + Send + Sync,

Paginate soft-deleted entities with filter support.

Source

pub async fn find_by_id(&self, id: &str) -> Result<Option<T>, Error>

Find an active (non-deleted) entity by primary key.

Source

pub async fn find_all(&self) -> Result<Vec<T>, Error>

Return all active (non-deleted) entities.

Source

pub async fn update(&self, id: &str, entity: &T) -> Result<Option<T>, Error>
where T: Serialize + Send + Sync,

Full update — skips silently if the record is already soft-deleted.

Source

pub async fn delete(&self, id: &str) -> Result<bool, Error>

Soft-delete an entity (sets metadata.deleted_at).

Source

pub async fn count(&self) -> Result<u64, Error>

Count active (non-deleted) entities.

Source

pub async fn exists(&self, id: &str) -> Result<bool, Error>

Return true if an active entity with the given ID exists.

Source

pub async fn list_paginated( &self, pagination: PaginationParams, ) -> Result<PaginatedResult<T>, Error>

Paginate active entities (most-recent-first by ID).

Source

pub async fn soft_delete(&self, id: &str) -> Result<bool, Error>

Set metadata.deleted_at to NOW() (soft delete).

Source

pub async fn restore(&self, id: &str) -> Result<Option<T>, Error>

Remove deleted_at from metadata, restoring the entity.

Source

pub async fn list_deleted( &self, pagination: PaginationParams, ) -> Result<PaginatedResult<T>, Error>

Paginate soft-deleted entities (trash view).

Source

pub async fn empty_trash(&self) -> Result<u64, Error>

Permanently delete all soft-deleted rows (empty trash).

Source

pub async fn find_deleted_by_id(&self, id: &str) -> Result<Option<T>, Error>

Find a soft-deleted entity by primary key.

Source

pub async fn permanent_delete(&self, id: &str) -> Result<bool, Error>

Permanently delete a soft-deleted entity by primary key.

Source

pub async fn count_active(&self) -> Result<u64, Error>

Count active (non-deleted) entities.

Source

pub async fn count_deleted(&self) -> Result<u64, Error>

Count soft-deleted entities.

Source

pub async fn bulk_soft_delete(&self, ids: &[String]) -> Result<u64, Error>

Soft-delete many active rows atomically.

Source

pub async fn bulk_restore(&self, ids: &[String]) -> Result<Vec<T>, Error>

Restore many soft-deleted rows atomically, returning the restored rows.

Source

pub async fn bulk_permanent_delete(&self, ids: &[String]) -> Result<u64, Error>

Permanently delete many soft-deleted rows atomically.

Source

pub async fn restore_all(&self) -> Result<Vec<T>, Error>

Restore every soft-deleted row, returning the restored rows. A single UPDATE ... RETURNING is atomic, and the returned rows let the service layer emit a Restored event per entity.

Source

pub async fn bulk_update(&self, entities: &[T]) -> Result<Vec<T>, Error>
where T: Serialize + Send + Sync,

Update many active rows atomically. Every entity must reference an existing active (non-soft-deleted) row or the whole batch is rolled back.

Trait Implementations§

Source§

impl CrudRepository<IntegrationEvent> for IntegrationEventRepository

Hydrate ?include= relations: fetch rows from table by id list, as JSON (raw row_to_json). table comes from EntityRepoMeta::relations() (generator-emitted), never client input. Default: no expansion. The Postgres-backed generated repos override this via impl_crud_repository!.
Source§

fn create<'life0, 'async_trait>( &'life0 self, entity: IntegrationEvent, ) -> Pin<Box<dyn Future<Output = Result<IntegrationEvent, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Create a new entity
Source§

fn find_by_id<'life0, 'life1, 'async_trait>( &'life0 self, id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<IntegrationEvent>, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Find entity by ID (excluding soft-deleted)
Source§

fn find_by_id_including_deleted<'life0, 'life1, 'async_trait>( &'life0 self, id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<IntegrationEvent>, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Find entity by ID (including soft-deleted, for trash operations)
Source§

fn update<'life0, 'async_trait>( &'life0 self, entity: IntegrationEvent, ) -> Pin<Box<dyn Future<Output = Result<IntegrationEvent, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Update an existing entity
Source§

fn soft_delete<'life0, 'life1, 'async_trait>( &'life0 self, id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<bool, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Soft delete an entity
Source§

fn restore<'life0, 'life1, 'async_trait>( &'life0 self, id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<Option<IntegrationEvent>, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Restore a soft-deleted entity
Source§

fn hard_delete<'life0, 'life1, 'async_trait>( &'life0 self, id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<bool, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Permanently delete an entity
Source§

fn list<'life0, 'async_trait>( &'life0 self, page: u32, limit: u32, ) -> Pin<Box<dyn Future<Output = Result<(Vec<IntegrationEvent>, u64), RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

List entities with pagination (excluding soft-deleted)
Source§

fn list_filtered<'life0, 'async_trait>( &'life0 self, page: u32, limit: u32, filters: HashMap<String, String>, ) -> Pin<Box<dyn Future<Output = Result<(Vec<IntegrationEvent>, u64), RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

List entities with pagination and filters (excluding soft-deleted) Read more
Source§

fn list_filtered_with_info<'life0, 'async_trait>( &'life0 self, page: u32, limit: u32, filters: HashMap<String, String>, ) -> Pin<Box<dyn Future<Output = Result<(Vec<IntegrationEvent>, PaginationInfo), RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

list_filtered, carrying the pagination info (cursor positions included) so the HTTP layer can surface keyset paging. Default: the tuple form with the cursor fields absent; the generated Postgres repositories override this with the real keyset walk.
Source§

fn table_name(&self) -> Option<&str>

Group and reduce entities under the same filters as list_filtered. Read more
Source§

fn aggregate_filtered<'life0, 'life1, 'async_trait>( &'life0 self, spec: &'life1 AggregateSpec, filters: HashMap<String, String>, ) -> Pin<Box<dyn Future<Output = Result<AggregateResult, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Source§

fn list_deleted<'life0, 'async_trait>( &'life0 self, page: u32, limit: u32, ) -> Pin<Box<dyn Future<Output = Result<(Vec<IntegrationEvent>, u64), RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

List soft-deleted entities with pagination
Source§

fn count<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<u64, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Count all entities (excluding soft-deleted)
Source§

fn count_deleted<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<u64, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Count soft-deleted entities
Source§

fn bulk_create<'life0, 'async_trait>( &'life0 self, entities: Vec<IntegrationEvent>, ) -> Pin<Box<dyn Future<Output = Result<Vec<IntegrationEvent>, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Bulk create entities
Source§

fn empty_trash<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<u64, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Permanently delete all soft-deleted entities
Source§

fn bulk_soft_delete<'life0, 'life1, 'async_trait>( &'life0 self, ids: &'life1 [String], ) -> Pin<Box<dyn Future<Output = Result<u64, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Soft-delete many entities by id. Returns the number affected.
Source§

fn bulk_restore<'life0, 'life1, 'async_trait>( &'life0 self, ids: &'life1 [String], ) -> Pin<Box<dyn Future<Output = Result<Vec<IntegrationEvent>, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Restore many soft-deleted entities by id. Returns the restored entities.
Source§

fn bulk_hard_delete<'life0, 'life1, 'async_trait>( &'life0 self, ids: &'life1 [String], ) -> Pin<Box<dyn Future<Output = Result<u64, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Permanently delete many entities by id. Returns the number affected.
Source§

fn restore_all<'life0, 'async_trait>( &'life0 self, ) -> Pin<Box<dyn Future<Output = Result<Vec<IntegrationEvent>, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Restore every soft-deleted entity. Returns the restored entities.
Source§

fn bulk_update<'life0, 'async_trait>( &'life0 self, entities: Vec<IntegrationEvent>, ) -> Pin<Box<dyn Future<Output = Result<Vec<IntegrationEvent>, RepositoryError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Update many entities atomically. Returns the updated entities.
Source§

fn exists<'life0, 'life1, 'async_trait>( &'life0 self, id: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<bool, RepositoryError>> + Send + 'async_trait>>
where 'life0: 'async_trait, 'life1: 'async_trait, Self: 'async_trait,

Check if an entity exists by ID
Source§

impl Deref for IntegrationEventRepository

Source§

type Target = GenericCrudRepository<IntegrationEvent>

The resulting type after dereferencing.
Source§

fn deref(&self) -> &Self::Target

Dereferences the value.

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<P, T> Receiver for P
where P: Deref<Target = T> + ?Sized, T: ?Sized,

Source§

type Target = T

🔬This is a nightly-only experimental API. (arbitrary_self_types)
The target type on which the method may be called.
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more