All checks were successful
Co-authored-by: Jeffrey Smith <jasafpro@gmail.com> Co-committed-by: Jeffrey Smith <jasafpro@gmail.com>
288 lines
8.5 KiB
Go
288 lines
8.5 KiB
Go
package postgres
|
|
|
|
import (
|
|
"context"
|
|
"database/sql"
|
|
"encoding/json"
|
|
"fmt"
|
|
"time"
|
|
|
|
"switchboard-core/models"
|
|
"switchboard-core/store"
|
|
)
|
|
|
|
type TriggerStore struct{}
|
|
|
|
func NewTriggerStore() *TriggerStore {
|
|
return &TriggerStore{}
|
|
}
|
|
|
|
func (s *TriggerStore) Create(ctx context.Context, t *models.Trigger) error {
|
|
cfg, _ := json.Marshal(t.Config)
|
|
if len(cfg) == 0 {
|
|
cfg = []byte("{}")
|
|
}
|
|
_, err := DB.ExecContext(ctx,
|
|
`INSERT INTO triggers (id, package_id, type, enabled, slug, secret, event_pattern,
|
|
entry_point, config)
|
|
VALUES (gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7, $8)
|
|
RETURNING id`,
|
|
t.PackageID, t.Type, t.Enabled, nilIfEmpty(t.Slug), nilIfEmpty(t.Secret),
|
|
nilIfEmpty(t.EventPattern), t.EntryPoint, cfg)
|
|
return err
|
|
}
|
|
|
|
func (s *TriggerStore) GetByID(ctx context.Context, id string) (*models.Trigger, error) {
|
|
var t models.Trigger
|
|
var cfg []byte
|
|
err := DB.QueryRowContext(ctx,
|
|
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
|
|
COALESCE(event_pattern,''), entry_point, config, fire_count,
|
|
last_fire_at, COALESCE(last_error,''), created_at, updated_at
|
|
FROM triggers WHERE id = $1`, id).
|
|
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
|
|
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
|
|
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt)
|
|
if err == sql.ErrNoRows {
|
|
return nil, nil
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
t.Config = json.RawMessage(cfg)
|
|
return &t, nil
|
|
}
|
|
|
|
func (s *TriggerStore) Update(ctx context.Context, t *models.Trigger) error {
|
|
cfg, _ := json.Marshal(t.Config)
|
|
if len(cfg) == 0 {
|
|
cfg = []byte("{}")
|
|
}
|
|
_, err := DB.ExecContext(ctx,
|
|
`UPDATE triggers SET enabled = $1, slug = $2, secret = $3, event_pattern = $4,
|
|
entry_point = $5, config = $6, updated_at = NOW()
|
|
WHERE id = $7`,
|
|
t.Enabled, nilIfEmpty(t.Slug), nilIfEmpty(t.Secret),
|
|
nilIfEmpty(t.EventPattern), t.EntryPoint, cfg, t.ID)
|
|
return err
|
|
}
|
|
|
|
func (s *TriggerStore) Delete(ctx context.Context, id string) error {
|
|
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE id = $1`, id)
|
|
return err
|
|
}
|
|
|
|
func (s *TriggerStore) List(ctx context.Context, opts store.TriggerListOptions) ([]models.Trigger, int, error) {
|
|
where := "1=1"
|
|
args := []any{}
|
|
n := 0
|
|
|
|
if opts.PackageID != "" {
|
|
n++
|
|
where += fmt.Sprintf(" AND package_id = $%d", n)
|
|
args = append(args, opts.PackageID)
|
|
}
|
|
if opts.Type != "" {
|
|
n++
|
|
where += fmt.Sprintf(" AND type = $%d", n)
|
|
args = append(args, opts.Type)
|
|
}
|
|
if opts.Enabled != nil {
|
|
n++
|
|
where += fmt.Sprintf(" AND enabled = $%d", n)
|
|
args = append(args, *opts.Enabled)
|
|
}
|
|
|
|
// Count
|
|
var total int
|
|
countArgs := make([]any, len(args))
|
|
copy(countArgs, args)
|
|
err := DB.QueryRowContext(ctx, "SELECT COUNT(*) FROM triggers WHERE "+where, countArgs...).Scan(&total)
|
|
if err != nil {
|
|
return nil, 0, err
|
|
}
|
|
|
|
// Query
|
|
limit := opts.Limit
|
|
if limit <= 0 {
|
|
limit = 50
|
|
}
|
|
n++
|
|
query := fmt.Sprintf(
|
|
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
|
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
|
created_at, updated_at
|
|
FROM triggers WHERE %s ORDER BY created_at DESC LIMIT $%d OFFSET $%d`,
|
|
where, n, n+1)
|
|
args = append(args, limit, opts.Offset)
|
|
|
|
rows, err := DB.QueryContext(ctx, query, args...)
|
|
if err != nil {
|
|
return nil, 0, err
|
|
}
|
|
defer rows.Close()
|
|
|
|
var triggers []models.Trigger
|
|
for rows.Next() {
|
|
var t models.Trigger
|
|
var cfg []byte
|
|
if err := rows.Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug,
|
|
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
|
|
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt); err != nil {
|
|
return nil, 0, err
|
|
}
|
|
t.Config = json.RawMessage(cfg)
|
|
triggers = append(triggers, t)
|
|
}
|
|
if triggers == nil {
|
|
triggers = []models.Trigger{}
|
|
}
|
|
return triggers, total, nil
|
|
}
|
|
|
|
func (s *TriggerStore) ListByPackage(ctx context.Context, packageID string) ([]models.Trigger, error) {
|
|
rows, err := DB.QueryContext(ctx,
|
|
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
|
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
|
created_at, updated_at
|
|
FROM triggers WHERE package_id = $1 ORDER BY created_at`, packageID)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanTriggerRows(rows)
|
|
}
|
|
|
|
func (s *TriggerStore) ListEnabledByType(ctx context.Context, triggerType string) ([]models.Trigger, error) {
|
|
rows, err := DB.QueryContext(ctx,
|
|
`SELECT id, package_id, type, enabled, COALESCE(slug,''), COALESCE(event_pattern,''),
|
|
entry_point, config, fire_count, last_fire_at, COALESCE(last_error,''),
|
|
created_at, updated_at
|
|
FROM triggers WHERE type = $1 AND enabled = true ORDER BY created_at`, triggerType)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanTriggerRows(rows)
|
|
}
|
|
|
|
func (s *TriggerStore) SetEnabled(ctx context.Context, id string, enabled bool) error {
|
|
_, err := DB.ExecContext(ctx,
|
|
`UPDATE triggers SET enabled = $1, updated_at = NOW() WHERE id = $2`, enabled, id)
|
|
return err
|
|
}
|
|
|
|
func (s *TriggerStore) UpdateFireState(ctx context.Context, id string, lastFireAt time.Time, lastError string) error {
|
|
_, err := DB.ExecContext(ctx,
|
|
`UPDATE triggers SET last_fire_at = $1, last_error = $2, fire_count = fire_count + 1,
|
|
updated_at = NOW() WHERE id = $3`,
|
|
lastFireAt, nilIfEmpty(lastError), id)
|
|
return err
|
|
}
|
|
|
|
func (s *TriggerStore) GetWebhook(ctx context.Context, packageID, slug string) (*models.Trigger, error) {
|
|
var t models.Trigger
|
|
var cfg []byte
|
|
err := DB.QueryRowContext(ctx,
|
|
`SELECT id, package_id, type, enabled, COALESCE(slug,''), secret,
|
|
COALESCE(event_pattern,''), entry_point, config, fire_count,
|
|
last_fire_at, COALESCE(last_error,''), created_at, updated_at
|
|
FROM triggers
|
|
WHERE package_id = $1 AND slug = $2 AND type = 'webhook'`, packageID, slug).
|
|
Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug, &t.Secret,
|
|
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
|
|
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt)
|
|
if err == sql.ErrNoRows {
|
|
return nil, nil
|
|
}
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
t.Config = json.RawMessage(cfg)
|
|
return &t, nil
|
|
}
|
|
|
|
func (s *TriggerStore) DeleteForPackage(ctx context.Context, packageID string) error {
|
|
_, err := DB.ExecContext(ctx, `DELETE FROM triggers WHERE package_id = $1`, packageID)
|
|
return err
|
|
}
|
|
|
|
func (s *TriggerStore) LogExecution(ctx context.Context, log *models.TriggerLog) error {
|
|
_, err := DB.ExecContext(ctx,
|
|
`INSERT INTO trigger_logs (id, trigger_id, scheduled_task_id, fired_at, duration_ms, success, error, output)
|
|
VALUES (gen_random_uuid(), $1, $2, $3, $4, $5, $6, $7)`,
|
|
nilIfEmpty(log.TriggerID), nilIfEmpty(log.ScheduledTaskID),
|
|
log.FiredAt, log.DurationMs, log.Success, nilIfEmpty(log.Error), nilIfEmpty(log.Output))
|
|
return err
|
|
}
|
|
|
|
func (s *TriggerStore) ListLogs(ctx context.Context, triggerID string, limit int) ([]models.TriggerLog, error) {
|
|
if limit <= 0 {
|
|
limit = 50
|
|
}
|
|
rows, err := DB.QueryContext(ctx,
|
|
`SELECT id, COALESCE(trigger_id::text,''), COALESCE(scheduled_task_id::text,''),
|
|
fired_at, duration_ms, success, COALESCE(error,''), COALESCE(output,'')
|
|
FROM trigger_logs
|
|
WHERE trigger_id = $1
|
|
ORDER BY fired_at DESC LIMIT $2`, triggerID, limit)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer rows.Close()
|
|
return scanLogRows(rows)
|
|
}
|
|
|
|
func (s *TriggerStore) PruneLogs(ctx context.Context, before time.Time) (int64, error) {
|
|
result, err := DB.ExecContext(ctx,
|
|
`DELETE FROM trigger_logs WHERE fired_at < $1`, before)
|
|
if err != nil {
|
|
return 0, err
|
|
}
|
|
return result.RowsAffected()
|
|
}
|
|
|
|
// ── Helpers ──────────────────────────────────
|
|
|
|
func scanTriggerRows(rows *sql.Rows) ([]models.Trigger, error) {
|
|
var triggers []models.Trigger
|
|
for rows.Next() {
|
|
var t models.Trigger
|
|
var cfg []byte
|
|
if err := rows.Scan(&t.ID, &t.PackageID, &t.Type, &t.Enabled, &t.Slug,
|
|
&t.EventPattern, &t.EntryPoint, &cfg, &t.FireCount,
|
|
&t.LastFireAt, &t.LastError, &t.CreatedAt, &t.UpdatedAt); err != nil {
|
|
return nil, err
|
|
}
|
|
t.Config = json.RawMessage(cfg)
|
|
triggers = append(triggers, t)
|
|
}
|
|
if triggers == nil {
|
|
triggers = []models.Trigger{}
|
|
}
|
|
return triggers, nil
|
|
}
|
|
|
|
func scanLogRows(rows *sql.Rows) ([]models.TriggerLog, error) {
|
|
var logs []models.TriggerLog
|
|
for rows.Next() {
|
|
var l models.TriggerLog
|
|
if err := rows.Scan(&l.ID, &l.TriggerID, &l.ScheduledTaskID,
|
|
&l.FiredAt, &l.DurationMs, &l.Success, &l.Error, &l.Output); err != nil {
|
|
return nil, err
|
|
}
|
|
logs = append(logs, l)
|
|
}
|
|
if logs == nil {
|
|
logs = []models.TriggerLog{}
|
|
}
|
|
return logs, nil
|
|
}
|
|
|
|
func nilIfEmpty(s string) *string {
|
|
if s == "" {
|
|
return nil
|
|
}
|
|
return &s
|
|
}
|