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/slack" ) type SqliteSlackStore struct { db dbi.DBI guard access.Guard } var _ slack.Store = (*SqliteSlackStore)(nil) func NewSqliteSlackStore(db dbi.DBI, guard access.Guard) *SqliteSlackStore { return &SqliteSlackStore{db: db, guard: guard} } // slackEndpointColumns is the fixed column list, scanned by scanSlackEndpoint. // 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 slackEndpointColumns = `id, tenant, name, event_types, active, dek_id, url_enc` func scanSlackEndpoint(sc rowScanner, out *slack.Endpoint) error { var types string err := sc.Scan(&out.ID, &out.Tenant, &out.Name, &types, &out.Active, &out.DekID, &out.URLEnc) 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 *SqliteSlackStore) CreateEndpoint(ctx context.Context, e *slack.Endpoint) error { if err := s.guard.Check(ctx, slack.Cap_Slack_CreateEndpoint, e.Tenant, ""); err != nil { return err } // Refuse to persist anything that isn't sealed ciphertext, so a plaintext // webhook URL can never reach the column (mirrors store_webhooks). if !envelope.IsSealed(e.URLEnc) { return fmt.Errorf("slack: webhook url is not sealed") } types, err := json.Marshal(e.EventTypes) if err != nil { return err } _, err = s.db.Exec(ctx, ` INSERT INTO slack_endpoints (`+slackEndpointColumns+`) VALUES ($1, $2, $3, $4, $5, $6, $7) `, e.ID, e.Tenant, e.Name, string(types), e.Active, e.DekID, e.URLEnc) return wrapQuotaErr(err) } func (s *SqliteSlackStore) UpdateEndpoint(ctx context.Context, e *slack.Endpoint) error { if err := s.guard.Check(ctx, slack.Cap_Slack_UpdateEndpoint, e.Tenant, ""); err != nil { return err } if !envelope.IsSealed(e.URLEnc) { return fmt.Errorf("slack: webhook url is not sealed") } types, err := json.Marshal(e.EventTypes) if err != nil { return err } res, err := s.db.Exec(ctx, ` UPDATE slack_endpoints SET name = $1, event_types = $2, active = $3, dek_id = $4, url_enc = $5 WHERE id = $6 AND tenant = $7 `, e.Name, string(types), e.Active, e.DekID, e.URLEnc, 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 *SqliteSlackStore) GetEndpoint(ctx context.Context, tenant core.ID, id core.ID, out *slack.Endpoint) error { if err := s.guard.Check(ctx, slack.Cap_Slack_ReadEndpoint, tenant, ""); err != nil { return err } row := s.db.QueryRow(ctx, ` SELECT `+slackEndpointColumns+` FROM slack_endpoints WHERE id = $1 AND tenant = $2 `, id, tenant) return dbi.TranslateNotFound(scanSlackEndpoint(row, out), core.ErrNotFound) } func (s *SqliteSlackStore) ListEndpoints(ctx context.Context, tenant core.ID, page core.PageReq, out *core.Page[slack.Endpoint]) error { if err := s.guard.Check(ctx, slack.Cap_Slack_ReadEndpoint, tenant, ""); err != nil { return err } limit := page.Limit if limit <= 0 { limit = 100 } rows, err := s.db.Query(ctx, ` SELECT `+slackEndpointColumns+` FROM slack_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 slack.Endpoint if err := scanSlackEndpoint(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 *SqliteSlackStore) DeleteEndpoint(ctx context.Context, tenant core.ID, id core.ID) error { if err := s.guard.Check(ctx, slack.Cap_Slack_DeleteEndpoint, tenant, ""); err != nil { return err } _, err := s.db.Exec(ctx, `DELETE FROM slack_endpoints WHERE id = $1 AND tenant = $2`, id, tenant) return err } func (s *SqliteSlackStore) ListActiveForEvent(ctx context.Context, tenant core.ID, eventType string) ([]slack.Endpoint, error) { if err := s.guard.Check(ctx, slack.Cap_Slack_ReadEndpoint, tenant, ""); err != nil { return nil, err } rows, err := s.db.Query(ctx, ` SELECT `+slackEndpointColumns+` FROM slack_endpoints WHERE tenant = $1 AND active = 1 `, tenant) if err != nil { return nil, err } defer rows.Close() var out []slack.Endpoint for rows.Next() { var e slack.Endpoint if err := scanSlackEndpoint(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() }