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" "ruben/inventory2/logging" ) type ( Store struct { log *logging.Logger db *pgxpool.Pool } Event struct { Platform string StoreID string EventID string EventTimestamp time.Time Payload json.RawMessage } ) func NewStore(logger *logging.Logger, db *pgxpool.Pool) *Store { return &Store{ log: logger, db: db, } } func (db *Store) WithContext(ctx context.Context) *StoreWithContext { return NewStoreWithContext(ctx, 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 *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 }