# Snowflake: Resonance Engine Objects

## Target Schema

**Schema**: `FANSIFTER_APP_REPORTING.DEV_MMACHADO` (dev). All objects prefixed with `RE_` for namespace isolation within the shared dev schema.

## Object Inventory

| Object | Type | Purpose |
|--------|------|---------|
| `RE_SPOTIFY_DATA_LANDING` | Table | Raw VARIANT rows from Snowpipe Streaming sink connector |
| `RE_SPOTIFY_DATA_STREAM_TOP_ARTISTS` | Stream | Append-only CDC stream — dedicated to top artists child task |
| `RE_SPOTIFY_DATA_STREAM_RECENTLY_PLAYED` | Stream | Append-only CDC stream — dedicated to recently played child task |
| `RE_SPOTIFY_DATA_STREAM_TOKENS` | Stream | Append-only CDC stream — dedicated to token extraction child task (Phase 5a) |
| `RE_FAN_TOP_ARTISTS` | Table | Flattened top artists per fan (up to 50 per collection) |
| `RE_FAN_RECENTLY_PLAYED` | Table | Flattened recently played tracks per fan (up to 50 per collection) |
| `RE_FAN_TOKENS` | Table | Refresh tokens per fan — source-of-truth for re-collection (Phase 5a) |
| `RE_TRANSFORM_SPOTIFY_DATA` | Task | Parent task — gates on any stream having data, runs every 5 min |
| `RE_TRANSFORM_TOP_ARTISTS` | Task | Child — flattens `endpoints.top_artists.items` via LATERAL FLATTEN |
| `RE_TRANSFORM_RECENTLY_PLAYED` | Task | Child — flattens `endpoints.recently_played.items` via LATERAL FLATTEN |
| `RE_EXTRACT_TOKENS` | Task | Child — MERGEs `refresh_token`/`new_refresh_token` into `RE_FAN_TOKENS` (Phase 5a) |
| `RE_RECOLLECTION_STAGE` | External Stage | S3 stage for recollection export files (Phase 5b) |
| `RE_RECOLLECTION_EXPORT` | Task | Scheduled — COPY INTO @stage to export active/stale fan tokens to S3 (Phase 5b) |
| `RE_TOKEN_MASK` | Masking Policy | Masks `RE_FAN_TOKENS.REFRESH_TOKEN` — returns `********` for non-privileged roles |
| `RE_LANDING_CONTENT_MASK` | Masking Policy | Masks `RE_SPOTIFY_DATA_LANDING.RECORD_CONTENT` — strips `refresh_token`/`new_refresh_token` keys via `OBJECT_DELETE` for non-privileged roles; analytics data remains visible |

## Masking Policy

Only `FANSIFTER_APP_REPORTING_DB_DEV_MMACHADO_SCHEMA_READWRITE` (sink connector + tasks) and `ACCOUNTADMIN` see plaintext refresh tokens. All other roles see masked values. This follows the org-standard infrastructure-layer encryption approach (see `token-encryption-strategy.md`). Application-layer field encryption (KMS encrypt/decrypt per token) was evaluated and rejected as inconsistent with org patterns.

## Design Notes

- Each child task has its own dedicated stream. Snowflake advances a stream's offset when the consuming DML commits, so parallel child tasks sharing a single stream would race — the second to commit sees an empty result set.
- **Warehouse**: `DEV_OWS_WH` (existing dev warehouse — no new warehouse creation needed).

## Snowflake Sink Connector (Fargate — Phase 4b)

Kafka Connect Snowflake sink connector running on Fargate. Consumes from MSK topic `resonance-engine.spotify-data` and writes to `RE_SPOTIFY_DATA_LANDING` via Snowpipe Streaming.

| Setting | Value |
|---------|-------|
| **Fargate service** | `kc-sfsink-resonance-engine` |
| **Docker image** | Shared org `kafka-connect-sfsink:latest` (account `086679231553`) |
| **MSK cluster** | `dev-managed-kafka-cdc-destination` (shared dev CDC cluster) |
| **Topics** | `resonance-engine.spotify-data` (6 partitions), `dlq.resonance-engine` (3 partitions) |
| **Snowflake user** | `DEV_KAFKA_CONNECT_RESONANCE_ENGINE` |
| **Snowflake role** | `FANSIFTER_APP_REPORTING_DB_DEV_MMACHADO_SCHEMA_READWRITE` |
| **Ingestion method** | Snowpipe Streaming |
| **Buffer** | 10,000 records / 5 MB / 180s flush |
| **Sizing** | 2 vCPU, 4 GB memory, 1 task (max 2) |

### Prerequisites (manual, before `terraform apply`)

1. Create Kafka topics via AKHQ: `resonance-engine.spotify-data` (6 partitions) and `dlq.resonance-engine` (3 partitions)
2. Create Snowflake service user `DEV_KAFKA_CONNECT_RESONANCE_ENGINE` with RSA key pair (via Snowflake admin)
3. After `terraform apply`: populate Secrets Manager secrets (`SNOWFLAKE_PRIVATE_KEY`, `SNOWFLAKE_PRIVATE_KEY_PASSPHRASE`)
4. Set `sink_bootstrap_servers` variable to dev MSK CDC bootstrap servers
5. Set `sink_desired_task_count = 1` to start the Fargate service
