package store import ( "context" "encoding/json" "fmt" "slices" "atlas9.dev/c/core" "atlas9.dev/c/core/dbi" "atlas9.dev/c/core/envelope" "atlas9.dev/c/demo/lib/access" "atlas9.dev/c/demo/lib/webhooks" ) type SqliteWebhookStore struct { db dbi.DBI guard access.Guard } var _ webhooks.Store = (*SqliteWebhookStore)(nil) func NewSqliteWebhookStore(db dbi.DBI, guard access.Guard) *SqliteWebhookStore { return &SqliteWebhookStore{db: db, guard: guard} } // endpointColumns is the fixed column list, scanned by scanEndpoint. event_types // is a JSON array of strings, so it can't map straight onto Endpoint.EventTypes // via reflection — every read decodes it explicitly. const endpointColumns = `id, tenant, name, url, event_types, active, dek_id, secret_enc` // rowScanner is satisfied by both *sql.Row and *sql.Rows. type rowScanner interface { Scan(dest ...any) error } func scanEndpoint(sc rowScanner, out *webhooks.Endpoint) error { var types string err := sc.Scan(&out.ID, &out.Tenant, &out.Name, &out.URL, &types, &out.Active, &out.DekID, &out.SecretEnc) if err != nil { return err } if types != "" { if err := json.Unmarshal([]byte(types), &out.EventTypes); err != nil { return fmt.Errorf("decoding event_types: %w", err) } } return nil } func (s *SqliteWebhookStore) CreateEndpoint(ctx context.Context, e *webhooks.Endpoint) error { if err := s.guard.Check(ctx, webhooks.Cap_Webhooks_CreateEndpoint, e.Tenant, ""); err != nil { return err } // Refuse to persist anything that isn't sealed ciphertext, so a plaintext // signing secret can never reach the column (mirrors store_sso). if !envelope.IsSealed(e.SecretEnc) { return fmt.Errorf("webhooks: signing secret is not sealed") } types, err := json.Marshal(e.EventTypes) if err != nil { return err } _, err = s.db.Exec(ctx, ` INSERT INTO webhook_endpoints (`+endpointColumns+`) VALUES ($1, $2, $3, $4, $5, $6, $7, $8) `, e.ID, e.Tenant, e.Name, e.URL, string(types), e.Active, e.DekID, e.SecretEnc) return wrapQuotaErr(err) } func (s *SqliteWebhookStore) UpdateEndpoint(ctx context.Context, e *webhooks.Endpoint) error { if err := s.guard.Check(ctx, webhooks.Cap_Webhooks_UpdateEndpoint, e.Tenant, ""); err != nil { return err } if !envelope.IsSealed(e.SecretEnc) { return fmt.Errorf("webhooks: signing secret is not sealed") } types, err := json.Marshal(e.EventTypes) if err != nil { return err } res, err := s.db.Exec(ctx, ` UPDATE webhook_endpoints SET name = $1, url = $2, event_types = $3, active = $4, dek_id = $5, secret_enc = $6 WHERE id = $7 AND tenant = $8 `, e.Name, e.URL, string(types), e.Active, e.DekID, e.SecretEnc, e.ID, e.Tenant) if err != nil { return err } n, err := res.RowsAffected() if err != nil { return err } if n == 0 { return core.ErrNotFound } return nil } func (s *SqliteWebhookStore) GetEndpoint(ctx context.Context, tenant core.ID, id core.ID, out *webhooks.Endpoint) error { if err := s.guard.Check(ctx, webhooks.Cap_Webhooks_ReadEndpoint, tenant, ""); err != nil { return err } row := s.db.QueryRow(ctx, ` SELECT `+endpointColumns+` FROM webhook_endpoints WHERE id = $1 AND tenant = $2 `, id, tenant) return dbi.TranslateNotFound(scanEndpoint(row, out), core.ErrNotFound) } func (s *SqliteWebhookStore) ListEndpoints(ctx context.Context, tenant core.ID, page core.PageReq, out *core.Page[webhooks.Endpoint]) error { if err := s.guard.Check(ctx, webhooks.Cap_Webhooks_ReadEndpoint, tenant, ""); err != nil { return err } limit := page.Limit if limit <= 0 { limit = 100 } rows, err := s.db.Query(ctx, ` SELECT `+endpointColumns+` FROM webhook_endpoints WHERE tenant = $1 AND id > $2 ORDER BY id LIMIT $3 `, tenant, page.Cursor, limit) if err != nil { return err } defer rows.Close() for rows.Next() { var e webhooks.Endpoint if err := scanEndpoint(rows, &e); err != nil { return err } out.Items = append(out.Items, e) } if err := rows.Err(); err != nil { return err } if len(out.Items) == limit { out.Cursor = out.Items[limit-1].ID.String() } return nil } func (s *SqliteWebhookStore) DeleteEndpoint(ctx context.Context, tenant core.ID, id core.ID) error { if err := s.guard.Check(ctx, webhooks.Cap_Webhooks_DeleteEndpoint, tenant, ""); err != nil { return err } _, err := s.db.Exec(ctx, `DELETE FROM webhook_endpoints WHERE id = $1 AND tenant = $2`, id, tenant) return err } func (s *SqliteWebhookStore) ListActiveForEvent(ctx context.Context, tenant core.ID, eventType string) ([]webhooks.Endpoint, error) { if err := s.guard.Check(ctx, webhooks.Cap_Webhooks_ReadEndpoint, tenant, ""); err != nil { return nil, err } rows, err := s.db.Query(ctx, ` SELECT `+endpointColumns+` FROM webhook_endpoints WHERE tenant = $1 AND active = 1 `, tenant) if err != nil { return nil, err } defer rows.Close() var out []webhooks.Endpoint for rows.Next() { var e webhooks.Endpoint if err := scanEndpoint(rows, &e); err != nil { return nil, err } // Subscriptions are a small per-endpoint list, so match in memory // rather than depending on a SQL JSON extension. if slices.Contains(e.EventTypes, eventType) { out = append(out, e) } } return out, rows.Err() }