Files
inventory-plus-plus/internal/domains/raw_events/events.go
T

135 lines
2.4 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"
)
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(db *pgxpool.Pool) *Store {
return &Store{
db: db,
}
}
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)
}