package raw_events import ( "context" "encoding/json" "fmt" "time" "github.com/jackc/pgx/v5" "github.com/jackc/pgx/v5/pgtype" "github.com/jackc/pgx/v5/pgxpool" ) type ( Store struct { db *pgxpool.Pool } // StoreWithContext allows currying all Store methods by context.Context StoreWithContext struct { ctx context.Context db *Store } Event struct { Platform string StoreID string EventID string EventTimestamp time.Time Payload json.RawMessage } ) func NewStore(ctx context.Context) (*Store, error) { pool, err := newPool(ctx) if err != nil { return nil, err } return &Store{ db: pool, }, nil } func newPool(ctx context.Context) (*pgxpool.Pool, error) { pool, err := pgxpool.New(ctx, "postgres://app_client:app_password@localhost:5432/inventory_2?sslmode=disable") if err != nil { return nil, fmt.Errorf("failed to create database client: %w", err) } conn, err := pool.Acquire(ctx) if err != nil { return nil, fmt.Errorf("failed to create a database connection: %w", err) } conn.Release() return pool, nil } func (db *Store) WithContext(ctx context.Context) *StoreWithContext { return &StoreWithContext{ ctx: ctx, db: db, } } func (db *Store) Save(ctx context.Context, e *Event) error { _, err := db.db.Exec( ctx, ` INSERT INTO raw_store_events ( platform, store_id, event_timestamp, event_id, raw_payload ) VALUES ( @platform, @store_id, @event_timestamp, @event_id, @raw_payload ) `, pgx.NamedArgs{ "platform": e.Platform, "store_id": e.StoreID, "event_timestamp": pgtype.Timestamptz{ Time: e.EventTimestamp, Valid: !e.EventTimestamp.IsZero(), }, "event_id": e.EventID, "raw_payload": string(e.Payload), }, ) if err != nil { return fmt.Errorf("failed to perform query: %w", err) } return nil } func (db *StoreWithContext) Save(e *Event) error { return db.db.Save(db.ctx, e) } func (db *Store) LoadEventsForStore( ctx context.Context, platform string, storeID string, ) ([]Event, error) { rows, err := db.db.Query( ctx, ` SELECT platform as Platform, store_id as StoreID, event_timestamp as EventTimestamp, event_id as EventID, raw_payload as Payload FROM raw_store_events WHERE platform = @platform AND store_id = @store_id ORDER BY event_timestamp DESC LIMIT 100 `, pgx.NamedArgs{ "platform": platform, "store_id": storeID, }, ) if err != nil { return nil, fmt.Errorf("failed to perform query: %w", err) } evts, err := pgx.CollectRows(rows, pgx.RowToStructByNameLax[Event]) if err != nil { return nil, fmt.Errorf("failed to scan events: %w", err) } return evts, nil } func (db *StoreWithContext) LoadEventsForStore( platform string, storeID string, ) ([]Event, error) { return db.db.LoadEventsForStore(db.ctx, platform, storeID) }