Files

119 lines
2.1 KiB
Go

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
}