backbone-integrations 0.6.0

Integration registry: connectors, integration accounts and an idempotent inbound event lane, with one OAuth flow (HMAC-bound state, PKCE)
Documentation
//! Snapshot Store
//!
//! Generated by metaphor-schema. Do not edit manually.
//!
//! Provides snapshot persistence for aggregate optimization.

use anyhow::Result;
use async_trait::async_trait;
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use sqlx::PgPool;
use uuid::Uuid;

// ============================================================================
// SNAPSHOT
// ============================================================================

/// Snapshot of aggregate state at a point in time
#[derive(Debug, Clone, Serialize, Deserialize, sqlx::FromRow)]
pub struct Snapshot {
    pub aggregate_id: Uuid,
    pub aggregate_type: String,
    pub version: i64,
    pub state: serde_json::Value,
    pub created_at: DateTime<Utc>,
}

/// Trait for snapshot store implementations
#[async_trait]
pub trait SnapshotStore: Send + Sync {
    /// Save a snapshot
    async fn save(&self, snapshot: &Snapshot) -> Result<()>;

    /// Load the latest snapshot for an aggregate
    async fn load(&self, aggregate_id: Uuid) -> Result<Option<Snapshot>>;

    /// Delete snapshots for an aggregate
    async fn delete(&self, aggregate_id: Uuid) -> Result<()>;
}

// ============================================================================
// POSTGRESQL SNAPSHOT STORE
// ============================================================================

/// PostgreSQL-based snapshot store implementation
pub struct PostgresSnapshotStore {
    pool: PgPool,
    table_name: String,
}

impl PostgresSnapshotStore {
    /// Create a new PostgreSQL snapshot store
    pub fn new(pool: PgPool) -> Self {
        Self {
            pool,
            table_name: "integrations.aggregate_snapshots".to_string(),
        }
    }

    /// Create with custom table name
    pub fn with_table_name(pool: PgPool, table_name: impl Into<String>) -> Self {
        Self {
            pool,
            table_name: table_name.into(),
        }
    }
}

#[async_trait]
impl SnapshotStore for PostgresSnapshotStore {
    async fn save(&self, snapshot: &Snapshot) -> Result<()> {
        let query = format!(
            "INSERT INTO {} (aggregate_id, aggregate_type, version, state, created_at) 
             VALUES ($1, $2, $3, $4, $5) 
             ON CONFLICT (aggregate_id) DO UPDATE SET 
             version = EXCLUDED.version, state = EXCLUDED.state, created_at = EXCLUDED.created_at",
            self.table_name
        );

        sqlx::query(&query)
            .bind(snapshot.aggregate_id)
            .bind(&snapshot.aggregate_type)
            .bind(snapshot.version)
            .bind(&snapshot.state)
            .bind(snapshot.created_at)
            .execute(&self.pool)
            .await?;

        Ok(())
    }

    async fn load(&self, aggregate_id: Uuid) -> Result<Option<Snapshot>> {
        let query = format!(
            "SELECT * FROM {} WHERE aggregate_id = $1",
            self.table_name
        );

        let snapshot = sqlx::query_as::<_, Snapshot>(&query)
            .bind(aggregate_id)
            .fetch_optional(&self.pool)
            .await?;

        Ok(snapshot)
    }

    async fn delete(&self, aggregate_id: Uuid) -> Result<()> {
        let query = format!(
            "DELETE FROM {} WHERE aggregate_id = $1",
            self.table_name
        );

        sqlx::query(&query)
            .bind(aggregate_id)
            .execute(&self.pool)
            .await?;

        Ok(())
    }
}

// ============================================================================
// SNAPSHOT STRATEGY
// ============================================================================

/// Configuration for when to create snapshots
#[derive(Debug, Clone)]
pub struct SnapshotStrategy {
    /// Whether snapshots are enabled
    pub enabled: bool,
    /// Create snapshot every N events
    pub every_n_events: u32,
    /// Maximum age before creating new snapshot
    pub max_age_seconds: Option<u64>,
    /// Storage backend (e.g., "postgres", "redis")
    pub storage: Option<String>,
}

impl Default for SnapshotStrategy {
    fn default() -> Self {
        Self {
            enabled: true,
            every_n_events: 100,
            max_age_seconds: Some(86400), // 24 hours
            storage: None,
        }
    }
}

impl SnapshotStrategy {
    /// Create a new snapshot strategy
    pub fn new(enabled: bool, every_n_events: u32) -> Self {
        Self {
            enabled,
            every_n_events,
            max_age_seconds: None,
            storage: None,
        }
    }

    /// Set max age before new snapshot
    pub fn with_max_age(mut self, seconds: u64) -> Self {
        self.max_age_seconds = Some(seconds);
        self
    }

    /// Set storage backend
    pub fn with_storage(mut self, storage: impl Into<String>) -> Self {
        self.storage = Some(storage.into());
        self
    }

    /// Check if a snapshot should be created
    pub fn should_snapshot(&self, current_version: i64, last_snapshot_version: i64) -> bool {
        self.enabled && (current_version - last_snapshot_version) as u32 >= self.every_n_events
    }
}

// <<< CUSTOM SNAPSHOT STORE START >>>
// Add custom snapshot implementations here
// <<< CUSTOM SNAPSHOT STORE END >>>