From 014a097fa37a520e489c97530f8f6f2db28241d9 Mon Sep 17 00:00:00 2001 From: Matheus Cruz Date: Mon, 14 Sep 2026 18:50:04 -0300 Subject: [PATCH 1/3] Remove raw tables, Vector insert directly to normalized tables instead Vector now inserts directly into workflow_instances/task_instances (raw_event JSONB column only), and the normalization triggers move onto those same tables via an INSERT ... ON CONFLICT DO UPDATE guarded by pg_trigger_depth() to prevent the nested self-insert from recursing. --- CLAUDE.md | 31 +-- .../collectors/VectorConfigValidationIT.java | 6 +- .../vector/mode1-postgresql/vector.yaml | 60 ++--- .../pages/architecture/postgresql-mode.adoc | 38 ++-- .../ROOT/pages/deployment/vector-config.adoc | 37 ++-- .../data-index-storage-migrations/README.md | 205 +++++++----------- .../V2__direct_normalized_inserts.sql | 197 +++++++++++++++++ .../scripts/e2e/verify-infrastructure.sh | 14 +- 8 files changed, 358 insertions(+), 230 deletions(-) create mode 100644 data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql diff --git a/CLAUDE.md b/CLAUDE.md index 18be33514..db84332b9 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -195,24 +195,22 @@ graphify --update ``` Quarkus Flow → stdout → /var/log/containers/*.log (JSON) - ↓ (Vector kubernetes_logs → postgres sink) - PostgreSQL raw tables (JSONB) - ↓ (BEFORE INSERT triggers) - PostgreSQL normalized tables + ↓ (Vector kubernetes_logs → postgres sink, {raw_event: } rows) + PostgreSQL normalized tables (BEFORE INSERT triggers normalize in place) ↓ (JPA/Hibernate) GraphQL API (SmallRye GraphQL) ``` **Key Components:** -- **Vector DaemonSet** - Tails container logs, writes raw events to PostgreSQL (`postgres` sink; config: `data-index/collectors/vector/mode1-postgresql/vector.yaml`). -- **PostgreSQL Triggers** - Normalize events immediately on INSERT +- **Vector DaemonSet** - Tails container logs, inserts directly into `workflow_instances`/`task_instances` (`postgres` sink, only the `raw_event` column set; config: `data-index/collectors/vector/mode1-postgresql/vector.yaml`). +- **PostgreSQL Triggers** - `BEFORE INSERT` on `workflow_instances`/`task_instances`; self-targeting `INSERT ... ON CONFLICT DO UPDATE`, guarded by `pg_trigger_depth()` to avoid recursion - **Data Index Service** - Quarkus app with GraphQL API - **JPA Entities** - Map to normalized tables (workflow_instances, task_instances) **NOT used in MODE 1:** - ❌ Event Processor service (removed in Phase 1) - ❌ Polling (triggers are immediate) -- ❌ Staging tables (raw tables → triggers → normalized tables) +- ❌ Raw staging tables (`workflow_events_raw`/`task_events_raw` — removed; Vector now inserts directly into the normalized tables) --- @@ -507,15 +505,16 @@ mvn quarkus:dev -Dquarkus.profile=elasticsearch ### 1. Trigger-Based Normalization (Not Polling) **DO:** -- ✅ Raw events stored in `workflow_events_raw` and `task_events_raw` (tag, time, data JSONB) -- ✅ Triggers extract fields from JSONB and UPSERT into normalized tables +- ✅ Vector inserts directly into `workflow_instances`/`task_instances`, with only `raw_event JSONB` set +- ✅ `BEFORE INSERT` triggers on these same tables extract fields and UPSERT in place - ✅ COALESCE handles out-of-order events - ✅ Real-time processing (< 1ms latency) **DON'T:** - ❌ Don't add Event Processor service (we removed it in Phase 1) - ❌ Don't use polling architecture -- ❌ Don't reference "staging tables" (we use raw tables + triggers) +- ❌ Don't reference `workflow_events_raw`/`task_events_raw` (removed) +- ❌ Don't drop the `pg_trigger_depth()` guard (self-targeting INSERT would recurse) **See:** `data-index/data-index-docs/modules/ROOT/pages/architecture/postgresql-mode.adoc` @@ -1281,9 +1280,10 @@ curl http://localhost:9200/_transform/workflow-instances-transform/_stats - Enable event tracing: `kubectl set env daemonset/vector -n logging DEBUG_EVENTS=true` - Confirm Vector can reach PostgreSQL (`POSTGRES_HOST`/`POSTGRES_PORT` env on the DaemonSet) -**"Raw tables populated but normalized tables empty"** -- Check triggers exist: `\d workflow_events_raw` in psql -- Check trigger functions: `\df normalize_workflow_event` +**"Rows inserted but fields not populated (namespace/status/etc. all NULL)"** +- Check triggers exist: `\d workflow_instances` / `\d task_instances` in psql (look for `Triggers:`) +- Check trigger functions: `\df normalize_workflow_instance`, `\df normalize_task_instance` +- Confirm the inserted row actually has `raw_event` set (a NULL `raw_event` is treated as a passthrough insert) - Check PostgreSQL logs for trigger errors **"GraphQL query returns empty taskExecutions"** @@ -1342,7 +1342,8 @@ curl http://localhost:9200/_transform/workflow-instances-transform/_stats - `data-index-storage-postgresql/src/main/java/.../entity/` - JPA entities - `data-index-storage-postgresql/src/main/java/.../mapper/` - MapStruct mappers - `data-index-storage-migrations/src/main/resources/db/migration/` - Flyway migrations -- `V1__initial_schema.sql` - Schema with triggers +- `V1__initial_schema.sql` - Original schema (raw staging tables, now removed) +- `V2__direct_normalized_inserts.sql` - Current triggers: self-targeting `BEFORE INSERT` on `workflow_instances`/`task_instances` **Code (MODE 2 - Elasticsearch):** - `data-index-storage-elasticsearch/src/main/java/.../` - Storage implementation @@ -1437,7 +1438,7 @@ curl http://localhost:9200/_transform/workflow-instances-transform/_stats → MODE 2: Write `@QuarkusTest` integration test with Elasticsearch profile + wait for transforms **"Where are the database triggers?"** -→ MODE 1: `data-index-storage-migrations/.../V1__initial_schema.sql` +→ MODE 1: `data-index-storage-migrations/.../V2__direct_normalized_inserts.sql` (defined directly on `workflow_instances`/`task_instances`) → MODE 2: Not applicable (uses Elasticsearch transforms) **"Where are the Elasticsearch transforms?"** diff --git a/data-index/collectors/src/test/java/org/kubesmarts/logic/apps/dataindex/collectors/VectorConfigValidationIT.java b/data-index/collectors/src/test/java/org/kubesmarts/logic/apps/dataindex/collectors/VectorConfigValidationIT.java index 0738dadbb..fb6a4c7a4 100644 --- a/data-index/collectors/src/test/java/org/kubesmarts/logic/apps/dataindex/collectors/VectorConfigValidationIT.java +++ b/data-index/collectors/src/test/java/org/kubesmarts/logic/apps/dataindex/collectors/VectorConfigValidationIT.java @@ -78,11 +78,11 @@ void mode1PostgreSQLConfigIsValid() throws Exception { .contains("kubernetes_logs:"); assertThat(configContent) - .as("Config should define one postgres sink per raw table") + .as("Config should define one postgres sink per normalized table") .contains("postgres_workflow:") .contains("postgres_task:") - .contains("table: workflow_events_raw") - .contains("table: task_events_raw"); + .contains("table: workflow_instances") + .contains("table: task_instances"); validateWithVectorContainer(configPath, MODE1_ENV); } diff --git a/data-index/collectors/vector/mode1-postgresql/vector.yaml b/data-index/collectors/vector/mode1-postgresql/vector.yaml index d088825f7..b0a4a4a58 100644 --- a/data-index/collectors/vector/mode1-postgresql/vector.yaml +++ b/data-index/collectors/vector/mode1-postgresql/vector.yaml @@ -7,15 +7,16 @@ # - Or our Maven tests: `mvn verify -pl collectors` # - Vector documentation: https://vector.dev/docs/reference/configuration/ # -# Purpose: Capture Quarkus Flow structured logs from container stdout -# PostgreSQL BEFORE INSERT triggers normalize raw events into -# structured tables (workflow_instances, task_instances). +# Purpose: Capture Quarkus Flow structured logs from container stdout and +# insert them directly into the normalized tables (workflow_instances, +# task_instances). BEFORE INSERT triggers on those same tables do +# the field extraction and idempotent merge. # # Pipeline: # Quarkus Flow -> stdout -> K8s /var/log/containers/ -> Vector kubernetes_logs -> # Parse JSON events -> Filter workflow events -> Route by event type -> -# PostgreSQL INSERT (workflow_events_raw / task_events_raw) -> -# PostgreSQL triggers -> Normalized tables +# PostgreSQL INSERT (workflow_instances / task_instances, {raw_event: }) -> +# PostgreSQL BEFORE INSERT triggers normalize in place # # Requires: Vector >= 0.46.0 (native `postgres` sink). # @@ -92,24 +93,15 @@ transforms: type: vrl source: starts_with(string(.flow_event.eventType) ?? "", "io.serverlessworkflow.task.") - # Build the row shape expected by the FluentBit-compatible raw tables: - # (tag TEXT, time TIMESTAMP WITH TIME ZONE, data JSONB) - # - # - tag: informational only (indexed) - the triggers never read it. - # - time: ingestion time - feeds only the created_at/updated_at audit columns. - # The triggers derive event ordering from `data->>'timestamp'`. - # - data: the untouched Quarkus Flow event. Triggers extract instanceId, - # status, startTime/endTime (epoch seconds), input/output, error, and - # taskPosition straight from this JSONB. + # Row inserted directly into workflow_instances/task_instances: just the + # untouched event as `raw_event` JSONB. The trigger extracts everything else. build_workflow_row: type: remap inputs: - route_by_type.workflow_events source: | . = { - "tag": .flow_event.eventType, - "time": now(), - "data": .flow_event + "raw_event": .flow_event } build_task_row: @@ -118,9 +110,7 @@ transforms: - route_by_type.task_events source: | . = { - "tag": .flow_event.eventType, - "time": now(), - "data": .flow_event + "raw_event": .flow_event } # Debug filter - only pass events if DEBUG_EVENTS=true @@ -138,20 +128,16 @@ transforms: to_bool(debug_enabled) ?? false # ============================================================================ -# SINKS: PostgreSQL Raw Tables +# SINKS: PostgreSQL Normalized Tables # ============================================================================ # -# Vector's `postgres` sink cannot template the table name per-event, so we use -# one sink per raw table. Each sink serializes the event to JSON and maps the -# top-level fields (tag, time, data) to the matching table columns via -# jsonb_populate_recordset - extra fields are ignored, and the `data` object -# maps cleanly to the JSONB column. +# One sink per table (Vector can't template the table name per-event). Each +# sink maps `raw_event` to the matching table's JSONB column via +# jsonb_populate_recordset; other columns are filled in by the trigger. # -# A batch is sent as a single multi-row INSERT. If one row makes a trigger -# raise an uncaught exception, the whole batch INSERT fails and Vector retries -# it. The workflow trigger is INSERT ... ON CONFLICT and the task trigger -# catches foreign_key_violation, so this is tolerable. If poison-pill rows -# appear, set `batch.max_events: 1` for strict row-by-row parity with FluentBit. +# A batch is one multi-row INSERT; if a trigger errors, the whole batch fails +# and retries. Triggers are idempotent, so retrying is safe. For poison-pill +# rows, set `batch.max_events: 1` for row-by-row parity with FluentBit. # # Buffer: the default in-memory buffer (matches the FluentBit MODE 1 setup this # replaces). Do NOT switch to `buffer.type: disk` here - in Vector 0.54 the @@ -160,26 +146,26 @@ transforms: # Vector restart re-reads recent lines and the idempotent triggers dedupe them. # https://vector.dev/docs/reference/configuration/sinks/postgres/ sinks: - # Workflow events -> workflow_events_raw + # Workflow events -> workflow_instances postgres_workflow: type: postgres inputs: - build_workflow_row endpoint: "postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@${POSTGRES_HOST}:${POSTGRES_PORT}/${POSTGRES_DB}?options=-c%20search_path%3D${POSTGRES_SCHEMA}" - table: workflow_events_raw + table: workflow_instances healthcheck: enabled: true batch: max_events: 100 timeout_secs: 1 - # Task events -> task_events_raw + # Task events -> task_instances postgres_task: type: postgres inputs: - build_task_row endpoint: "postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@${POSTGRES_HOST}:${POSTGRES_PORT}/${POSTGRES_DB}?options=-c%20search_path%3D${POSTGRES_SCHEMA}" - table: task_events_raw + table: task_instances healthcheck: enabled: true batch: @@ -246,7 +232,7 @@ api: # * WARNING: At 1000 events/sec, generates ~43GB/day of logs # # What stays the same (handled by PostgreSQL): -# - BEFORE INSERT triggers: normalize_workflow_event() / normalize_task_event() +# - BEFORE INSERT triggers: normalize_workflow_instance() / normalize_task_instance() # - Field-level idempotency (COALESCE) and out-of-order event handling # - Composite task identity (instance_id, task) # @@ -255,7 +241,7 @@ api: # * Parses CRI/Docker format automatically # * Filters workflow events (by eventType field) # * Routes by event type (workflow vs task) -# * Inserts raw events into workflow_events_raw / task_events_raw +# * Inserts {raw_event: } rows directly into workflow_instances / task_instances # * Exposes Prometheus metrics # # ============================================================================ diff --git a/data-index/data-index-docs/modules/ROOT/pages/architecture/postgresql-mode.adoc b/data-index/data-index-docs/modules/ROOT/pages/architecture/postgresql-mode.adoc index 9f2cb046d..f1400fff8 100644 --- a/data-index/data-index-docs/modules/ROOT/pages/architecture/postgresql-mode.adoc +++ b/data-index/data-index-docs/modules/ROOT/pages/architecture/postgresql-mode.adoc @@ -10,10 +10,8 @@ Data Index with PostgreSQL storage uses trigger-based normalization for real-tim Quarkus Flow App ↓ (structured logging → stdout) Vector DaemonSet - ↓ (kubernetes_logs source → postgres sink) -PostgreSQL Raw Tables (JSONB) - ↓ (BEFORE INSERT triggers) -PostgreSQL Normalized Tables + ↓ (kubernetes_logs source → postgres sink, {raw_event: } rows) +PostgreSQL Normalized Tables (BEFORE INSERT triggers normalize in place) ↓ (JPA/Hibernate) GraphQL API ---- @@ -23,9 +21,9 @@ GraphQL API 1. **Quarkus Flow emits events** - JSON to stdout 2. **Kubernetes captures logs** - `/var/log/containers/POD_NAME.log` 3. **Vector collects** - `kubernetes_logs` source tails container logs, filters events by `eventType`, routes workflow vs task -4. **INSERT to raw tables** - `postgres` sinks write `workflow_events_raw`, `task_events_raw` (JSONB `data` column) -5. **Triggers fire immediately** - Extract fields from JSONB -6. **UPSERT to normalized tables** - `workflow_instances`, `task_instances` +4. **INSERT directly into normalized tables** - `postgres` sinks write `workflow_instances`, `task_instances`, with only the `raw_event` JSONB column set +5. **Triggers fire immediately** - `BEFORE INSERT` triggers, defined on these same tables, extract fields from `raw_event` +6. **Merge in place** - the trigger self-targets an `INSERT ... ON CONFLICT DO UPDATE`, guarded by `pg_trigger_depth()`, and cancels the original pending insert 7. **GraphQL queries** - Via JPA entities == Key Characteristics @@ -58,8 +56,8 @@ GraphQL API Events are normalized in real-time as they arrive: -1. **Vector writes raw events** - Complete events stored as JSON for debugging -2. **Events normalized automatically** - Fields extracted and stored in optimized tables +1. **Vector inserts directly into the normalized tables** - each row carries only the event as JSON (`raw_event` column) +2. **Events normalized automatically** - the `BEFORE INSERT` trigger extracts fields from `raw_event` and merges them into the row before the insert completes 3. **Immediate availability** - Data ready for querying in less than 1ms **Benefits:** @@ -75,20 +73,18 @@ Events are normalized in real-time as they arrive: - Limited throughput vs. Elasticsearch mode (< 50K workflows/day) - Schema changes require database updates -== Raw Event Storage +== Raw Event Column -Raw events are preserved in their original JSON format: +There are no separate raw staging tables. Instead, `workflow_instances` and `task_instances` each +carry a `raw_event JSONB` column holding the **most recently applied event** for that row: -- **workflow_events_raw** - All workflow-related events -- **task_events_raw** - All task-related events - -**Benefits:** - -- **Debugging** - Original events preserved for troubleshooting -- **Replay** - Can reprocess if normalization logic changes -- **Audit** - Complete event history maintained +- **Debugging** - The latest original event is available for troubleshooting without a separate table join - **Flexibility** - Accepts any event structure without schema changes +**Trade-off vs. the previous raw-table design:** only the latest event per row is kept, not the +full event history - there's no replay of earlier events once a newer one has merged over them. +This was a deliberate simplification (fewer tables, no recursive-trigger risk). + == Normalized Tables Events are automatically normalized to optimized tables for querying: @@ -131,12 +127,12 @@ sinks: type: postgres inputs: [build_workflow_row] endpoint: "postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@${POSTGRES_HOST}:${POSTGRES_PORT}/${POSTGRES_DB}" - table: workflow_events_raw + table: workflow_instances postgres_task: type: postgres inputs: [build_task_row] endpoint: "postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@${POSTGRES_HOST}:${POSTGRES_PORT}/${POSTGRES_DB}" - table: task_events_raw + table: task_instances ---- See xref:deployment/vector-config.adoc[Vector Configuration] for the full pipeline. diff --git a/data-index/data-index-docs/modules/ROOT/pages/deployment/vector-config.adoc b/data-index/data-index-docs/modules/ROOT/pages/deployment/vector-config.adoc index 1920e5ae1..f2fe985cb 100644 --- a/data-index/data-index-docs/modules/ROOT/pages/deployment/vector-config.adoc +++ b/data-index/data-index-docs/modules/ROOT/pages/deployment/vector-config.adoc @@ -2,10 +2,9 @@ :page-aliases: Data Index uses https://vector.dev[Vector] as the log collector for both -MODE 1 (PostgreSQL) and MODE 2 (Elasticsearch). Vector tails Quarkus Flow -container logs, keeps only structured workflow/task events, and writes them to -the mode's raw storage. Normalization then happens in the storage layer -(PostgreSQL triggers or Elasticsearch Transforms). +MODE 1 (PostgreSQL) and MODE 2 (Elasticsearch). MODE 1 inserts directly into +`workflow_instances`/`task_instances`, normalized by `BEFORE INSERT` triggers. +MODE 2 writes to raw Elasticsearch indices, normalized by Transforms. == Sources of truth @@ -43,14 +42,15 @@ Quarkus Flow pod → stdout → /var/log/containers/*.log filter_workflow_events → keep eventType starting with "io.serverlessworkflow." route_by_type → workflow_events | task_events ↓ - MODE 1: build_{workflow,task}_row → postgres sinks → workflow_events_raw / task_events_raw + MODE 1: build_{workflow,task}_row → postgres sinks → workflow_instances / task_instances MODE 2: (enrich @timestamp) → elasticsearch sinks → workflow-events-* / task-events-* ---- == MODE 1 (PostgreSQL) The `postgres` sink requires **Vector >= 0.46.0**. Because it cannot template the -table name per event, MODE 1 uses one sink per raw table. +table name per event, MODE 1 uses one sink per normalized table, and inserts +directly into it — there is no separate raw staging table. [source,yaml] ---- @@ -59,20 +59,20 @@ transforms: type: remap inputs: [route_by_type.workflow_events] source: | - . = { "tag": .flow_event.eventType, "time": now(), "data": .flow_event } + . = { "raw_event": .flow_event } sinks: postgres_workflow: type: postgres inputs: [build_workflow_row] endpoint: "postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@${POSTGRES_HOST}:${POSTGRES_PORT}/${POSTGRES_DB}?options=-c%20search_path%3D${POSTGRES_SCHEMA}" - table: workflow_events_raw + table: workflow_instances batch: { max_events: 100, timeout_secs: 1 } postgres_task: type: postgres inputs: [build_task_row] endpoint: "postgres://${POSTGRES_USER}:${POSTGRES_PASSWORD}@${POSTGRES_HOST}:${POSTGRES_PORT}/${POSTGRES_DB}?options=-c%20search_path%3D${POSTGRES_SCHEMA}" - table: task_events_raw + table: task_instances batch: { max_events: 100, timeout_secs: 1 } ---- @@ -100,21 +100,16 @@ setup). Do not set `buffer.type: disk` on the `postgres` sink: in Vector 0.54 it is not drained until graceful shutdown, stranding low-volume events. The `kubernetes_logs` checkpoint plus idempotent triggers cover a restart. -Row shape maps directly to the raw table columns `(tag TEXT, time TIMESTAMPTZ, data JSONB)`: - -* `tag` — informational only (indexed); the triggers never read it. -* `time` — ingestion time; feeds only the `created_at`/`updated_at` audit columns. -* `data` — the untouched Quarkus Flow event. `normalize_workflow_event()` / - `normalize_task_event()` extract `instanceId`, `status`, `startTime`/`endTime` - (epoch seconds), `input`/`output`, `error`, and `taskPosition` from this JSONB. +Vector only ever sets `raw_event JSONB`; the `BEFORE INSERT` trigger extracts +everything else (`instanceId`, `status`, timestamps, `input`/`output`, `error`, +`taskPosition`) and self-targets `INSERT ... ON CONFLICT DO UPDATE`, guarded by +`pg_trigger_depth()` to avoid recursion. [NOTE] ==== -A Vector batch is one multi-row `INSERT`. If a row makes a trigger raise an -uncaught exception the whole batch fails and is retried. The workflow trigger is -`INSERT ... ON CONFLICT` and the task trigger catches `foreign_key_violation`, so -this is tolerable. If a poison-pill row appears, set `batch.max_events: 1` for -strict row-by-row parity with FluentBit. +A Vector batch is one multi-row `INSERT`; if a trigger errors, the whole batch +fails and retries. Triggers are idempotent, so retrying is safe. For poison-pill +rows, set `batch.max_events: 1` for row-by-row parity with FluentBit. ==== Required environment (set on the DaemonSet): diff --git a/data-index/data-index-storage/data-index-storage-migrations/README.md b/data-index/data-index-storage/data-index-storage-migrations/README.md index 04565dbf7..5704b3b1e 100644 --- a/data-index/data-index-storage/data-index-storage-migrations/README.md +++ b/data-index/data-index-storage/data-index-storage-migrations/README.md @@ -4,11 +4,10 @@ ## Purpose -This module contains Flyway migration scripts for Data Index PostgreSQL storage backend: +Flyway migration scripts for the Data Index PostgreSQL storage backend: -- **Raw Tables**: `workflow_events_raw`, `task_events_raw` (raw events from the Vector `postgres` sink) -- **Normalized Tables**: `workflow_instances`, `task_instances` (optimized for querying) -- **Real-time processing**: Events are normalized automatically as they arrive +- **Normalized Tables**: `workflow_instances`, `task_instances` (direct insert target for the MODE 1 Vector `postgres` sink) +- **Real-time processing**: `BEFORE INSERT` triggers on these same two tables normalize events as they arrive ## Architecture @@ -16,14 +15,13 @@ This module contains Flyway migration scripts for Data Index PostgreSQL storage Quarkus Flow (quarkus-flow 0.9.0+) ↓ (structured JSON events with epoch timestamps → stdout) Vector DaemonSet (postgres sink) - ├─→ workflow_events_raw (tag TEXT, time TIMESTAMP, data JSONB) - └─→ task_events_raw (tag TEXT, time TIMESTAMP, data JSONB) - ↓ (BEFORE INSERT triggers) + ├─→ workflow_instances (raw_event JSONB column only) + └─→ task_instances (raw_event JSONB column only) + ↓ (BEFORE INSERT triggers, self-targeting) PostgreSQL Triggers - ├─→ Extract fields from data JSONB - ├─→ Handle out-of-order events - ├─→ Upsert workflow_instances - └─→ Upsert task_instances + ├─→ Extract fields from raw_event JSONB + ├─→ Handle out-of-order events (COALESCE/GREATEST) + ├─→ UPDATE existing row + RETURN NULL, or fill in NEW + RETURN NEW ↓ (GraphQL queries) Data Index GraphQL API ``` @@ -32,90 +30,71 @@ Data Index GraphQL API ### V1__initial_schema.sql -Initial schema for Data Index v1.0.0 with trigger-based normalization: +Original schema for Data Index v1.0.0: creates `workflow_instances` and `task_instances` +(without `raw_event`), plus the now-removed raw staging tables and their triggers. -**Raw Staging Tables** (`tag TEXT, time TIMESTAMP WITH TIME ZONE, data JSONB`): -- `workflow_events_raw`: Stores tag (TEXT), time (TIMESTAMP), data (JSONB) -- `task_events_raw`: Stores tag (TEXT), time (TIMESTAMP), data (JSONB) +### V2__direct_normalized_inserts.sql -This 3-column shape originated with the FluentBit pgsql plugin's fixed schema and -is retained: the Vector `postgres` sink writes the same `{tag, time, data}` rows. +- Adds `raw_event JSONB` to `workflow_instances` and `task_instances`. +- Drops `workflow_events_raw`, `task_events_raw`, and their triggers/functions. +- Adds `normalize_workflow_instance()` / `normalize_task_instance()`, `BEFORE INSERT` on + `workflow_instances` / `task_instances` respectively — self-targeting `INSERT ... ON CONFLICT + DO UPDATE`, guarded by `pg_trigger_depth()` to avoid recursion. -**Normalized Tables** (Populated via triggers): -- `workflow_instances`: Extracted workflow state with individual columns -- `task_instances`: Extracted task state (with FK to workflow_instances) +## Field Mappings (raw_event JSONB → Normalized Tables) -**Trigger Functions**: -- `normalize_workflow_event()`: Extracts fields from JSONB and UPSERTs into workflow_instances -- `normalize_task_event()`: Extracts fields from JSONB and UPSERTs into task_instances +### Workflow Events (raw_event JSONB → workflow_instances) -**Advantages**: -- No Event Processor needed - triggers handle normalization in real-time -- Automatic out-of-order event handling via UPSERT with COALESCE -- Raw events preserved in staging tables for debugging/replay -- Simpler architecture - fewer moving parts +`normalize_workflow_instance()` extracts: -## Raw Table Schema - -The raw staging tables use a fixed 3-column shape: -- `tag TEXT` - event tag (the event's `eventType`, e.g. `io.serverlessworkflow.workflow.started.v1`) -- `time TIMESTAMP WITH TIME ZONE` - ingestion time (feeds only `created_at`/`updated_at`) -- `data JSONB` - the complete Quarkus Flow event as JSON - -The MODE 1 log collector (Vector's `postgres` sink; previously the FluentBit pgsql -plugin) writes rows in this shape, and the **triggers** normalize `data` into the -`workflow_instances` / `task_instances` tables. - -## Field Mappings (JSONB → Normalized Tables) - -### Workflow Events (data JSONB → workflow_instances) - -Trigger function extracts: - -| JSONB Path (data->>) | workflow_instances Column | Type | Conversion | -|----------------------|---------------------------|------|------------| +| JSONB Path (raw_event->>) | workflow_instances Column | Type | Conversion | +|----------------------------|----------------------------|------|------------| | `instanceId` | `id` | VARCHAR(255) | Direct | | `workflowNamespace` | `namespace` | VARCHAR(255) | Direct | | `workflowName` | `name` | VARCHAR(255) | Direct | | `workflowVersion` | `version` | VARCHAR(255) | Direct | | `status` | `status` | VARCHAR(50) | Direct | -| `startTime` | `start` | TIMESTAMP | `to_timestamp(::numeric)` | -| `endTime` | `end` | TIMESTAMP | `to_timestamp(::numeric)` | -| `lastUpdateTime` | `last_update` | TIMESTAMP | `to_timestamp(::numeric)` | -| `input` | `input` | JSONB | Direct (->) | -| `output` | `output` | JSONB | Direct (->) | +| `startTime` | `started_at` | TIMESTAMPTZ | `to_timestamp(::numeric)` | +| `endTime` | `ended_at` | TIMESTAMPTZ | `to_timestamp(::numeric)` | +| `lastUpdateTime` | `last_update` | TIMESTAMPTZ | `to_timestamp(::numeric)` | +| `timestamp` | `last_event_time` | TIMESTAMPTZ | `to_timestamp(::numeric)` (drives idempotency) | +| `input` | `input` | JSONB | Direct (`->`) | +| `output` | `output` | JSONB | Direct (`->`) | | `error->>'type'` | `error_type` | VARCHAR(255) | Nested | | `error->>'title'` | `error_title` | VARCHAR(255) | Nested | | `error->>'detail'` | `error_detail` | TEXT | Nested | | `error->>'status'` | `error_status` | INTEGER | Nested + cast | | `error->>'instance'` | `error_instance` | VARCHAR(255) | Nested | -**Auto-populated:** -- `created_at` (from trigger: NEW.time) -- `updated_at` (from trigger: NEW.time) +**Auto-populated:** `created_at`/`updated_at` (`now()`, set by the trigger). -### Task Events (data JSONB → task_instances) +### Task Events (raw_event JSONB → task_instances) -Trigger function extracts: +`normalize_task_instance()` extracts: -| JSONB Path (data->>) | task_instances Column | Type | Conversion | -|----------------------|-----------------------|------|------------| -| `instanceId` | `instance_id` | VARCHAR(255) | Direct (PK part 1) | -| `taskPosition` | `task_position` | VARCHAR(255) | Direct (PK part 2) | +| JSONB Path (raw_event->>) | task_instances Column | Type | Conversion | +|----------------------------|------------------------|------|------------| +| `instanceId` | `instance_id` | VARCHAR(255) | Direct (PK part 1, FK to workflow_instances) | +| `taskPosition` | `task` | VARCHAR(255) | Direct (PK part 2, JSON Pointer e.g. `/do/1/initialize`) | | `taskName` | `task_name` | VARCHAR(255) | Direct | | `status` | `status` | VARCHAR(50) | Direct | -| `startTime` | `start` | TIMESTAMP | `to_timestamp(::numeric)` | -| `endTime` | `end` | TIMESTAMP | `to_timestamp(::numeric)` | -| `input` | `input` | JSONB | Direct (->) | -| `output` | `output` | JSONB | Direct (->) | +| `startTime` | `started_at` | TIMESTAMPTZ | `to_timestamp(::numeric)` | +| `endTime` | `ended_at` | TIMESTAMPTZ | `to_timestamp(::numeric)` | +| `timestamp` | `last_event_time` | TIMESTAMPTZ | `to_timestamp(::numeric)` (drives idempotency) | +| `input` | `input` | JSONB | Direct (`->`) | +| `output` | `output` | JSONB | Direct (`->`) | +| `error->>'type'` | `error_type` | VARCHAR(255) | Nested | +| `error->>'title'` | `error_title` | VARCHAR(255) | Nested | +| `error->>'detail'` | `error_detail` | TEXT | Nested | +| `error->>'status'` | `error_status` | INTEGER | Nested + cast | +| `error->>'instance'` | `error_instance` | VARCHAR(255) | Nested | -**Auto-populated:** -- `created_at` (from trigger: NEW.time) -- `updated_at` (from trigger: NEW.time) +**Auto-populated:** `created_at`/`updated_at` (`now()`, set by the trigger). -**Note:** Task payloads are included when `quarkus.flow.structured-logging.include-task-payloads=true` +**Note:** no `task_execution_id` column — `TaskExecution.id` is derived as `instanceId + ":" + task`. -**Out-of-Order Handling:** The task trigger first ensures the workflow instance exists (creates placeholder if needed) before inserting the task. +**Out-of-Order Handling:** the task trigger backfills a placeholder `workflow_instances` row +(`raw_event` left `NULL`) if the parent doesn't exist yet, before processing the task event. ## Usage @@ -124,85 +103,57 @@ Trigger function extracts: Apply migrations manually: ```bash -# Connect to PostgreSQL -kubectl exec -n postgresql postgresql-0 -- env PGPASSWORD=dataindex123 \ - psql -U dataindex -d dataindex - -# Apply migration kubectl exec -n postgresql postgresql-0 -- env PGPASSWORD=dataindex123 \ psql -U dataindex -d dataindex -f /path/to/V1__initial_schema.sql +kubectl exec -n postgresql postgresql-0 -- env PGPASSWORD=dataindex123 \ + psql -U dataindex -d dataindex -f /path/to/V2__direct_normalized_inserts.sql ``` ### Kubernetes Operator -The Data Index operator will use Flyway to manage migrations automatically when users choose PostgreSQL storage. - -**Operator behavior:** -1. Detects PostgreSQL storage backend -2. Runs Flyway migrations on startup -3. Compares schema version with migration files -4. Applies pending migrations +The Data Index operator uses Flyway to manage migrations automatically for PostgreSQL storage. **Upgrade safety:** -- Migrations are idempotent (`CREATE TABLE IF NOT EXISTS`) -- Foreign keys use `ON DELETE CASCADE` for data integrity -- Indexes created with `IF NOT EXISTS` to prevent errors +- `V2__` is additive on top of an already-provisioned `V1__` schema, including existing rows. +- Foreign keys use `ON DELETE CASCADE`. +- Triggers are idempotent regardless of how many times the same event is (re)inserted. ## Trigger-Based Normalization -PostgreSQL triggers handle all normalization automatically - **no Event Processor needed!** +No Event Processor needed — PostgreSQL triggers handle normalization directly. ### How It Works -1. **Vector `postgres` sink INSERT** → `workflow_events_raw` or `task_events_raw` -2. **BEFORE INSERT trigger fires** → Extracts fields from JSONB `data` column -3. **UPSERT normalized table** → `workflow_instances` or `task_instances` -4. **Return NEW** → Raw event is also stored in staging table +1. Vector's `postgres` sink inserts into `workflow_instances`/`task_instances` with only `raw_event` set. +2. The `BEFORE INSERT` trigger fires on that same table (`pg_trigger_depth()` guard: a nested call from step 4 passes through unchanged). +3. Extracts fields from `raw_event` (`NULL` `raw_event` also passes through unchanged — placeholders, JPA-seeded test rows). +4. Runs its own `INSERT ... ON CONFLICT (id) DO UPDATE` against the same table (COALESCE/GREATEST rules). +5. `RETURN NULL`s to cancel the original pending insert — step 4 already applied the event. + +`ON CONFLICT` (not a manual `UPDATE`-then-`INSERT`) is required for atomicity: Vector's +`postgres_workflow`/`postgres_task` sinks write concurrently, and a check-then-act pattern is race-prone. ### Advantages -- **Real-time normalization**: No polling or batch delays -- **Out-of-order handling**: UPSERT with COALESCE handles events arriving in any order -- **Idempotent**: Same event can be inserted multiple times safely -- **Simpler architecture**: Fewer services to deploy and monitor -- **Debugging**: Raw events preserved in staging tables +- **Real-time**: no polling or batch delays +- **Out-of-order handling**: atomic `ON CONFLICT DO UPDATE` with COALESCE/GREATEST, safe under concurrency +- **Idempotent**: same event can be reinserted safely +- **Simpler**: two tables instead of four +- **Debugging**: most recent raw event kept per row via `raw_event` -### Example Workflow Event Processing +### Example ```sql --- Vector postgres sink inserts: -INSERT INTO workflow_events_raw (tag, time, data) -VALUES ('workflow.instance.started', '2026-04-23 22:04:49+00', '{"instanceId":"01KPY...","workflowName":"simple-set",...}'); - --- Trigger automatically executes: -INSERT INTO workflow_instances (id, namespace, name, version, status, start, ...) -VALUES ('01KPY...', 'org.acme', 'simple-set', '0.0.1', 'RUNNING', '2026-04-23 22:04:49+00', ...) -ON CONFLICT (id) DO UPDATE SET - status = EXCLUDED.status, - "end" = COALESCE(EXCLUDED."end", workflow_instances."end"), - ...; -``` - -### Retention Policy - -Raw staging tables can be cleaned up periodically: +-- Vector inserts: +INSERT INTO workflow_instances (raw_event) +VALUES ('{"instanceId":"01KPY...","workflowName":"simple-set","status":"RUNNING",...}'::jsonb); -```sql --- Delete raw events older than 7 days -DELETE FROM workflow_events_raw WHERE time < NOW() - INTERVAL '7 days'; -DELETE FROM task_events_raw WHERE time < NOW() - INTERVAL '7 days'; +-- Trigger extracts fields and upserts workflow_instances in place, then cancels this insert. ``` -This can be scheduled via PostgreSQL `pg_cron` extension or external cron job. - ## Notes -- **Raw table schema**: fixed `(tag TEXT, time TIMESTAMP, data JSONB)` shape; the Vector `postgres` sink maps event fields to these columns -- **Timestamps**: Epoch seconds from Quarkus Flow 0.9.0+ converted via `to_timestamp()` in triggers -- **JSONB Storage**: Complete events stored as JSONB in `data` column, triggers extract to columns -- **Raw Events**: Preserved in `*_raw` tables for debugging and potential replay -- **JSONB Performance**: Normalized columns indexed for fast GraphQL queries -- **Cascade Deletes**: Deleting a workflow instance cascades to its tasks -- **Idempotent**: All CREATE statements use `IF NOT EXISTS`, triggers use UPSERT with COALESCE -- **Out-of-Order**: Triggers handle events arriving in any order via UPSERT logic -- **No Event Processor**: PostgreSQL triggers replace the need for a separate Event Processor service +- `raw_event` holds only the most recent event per row (not full history). +- Timestamps: epoch seconds from Quarkus Flow, converted via `to_timestamp()`. +- Normalized columns are indexed for GraphQL queries. +- Deleting a workflow instance cascades to its tasks. \ No newline at end of file diff --git a/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql b/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql new file mode 100644 index 000000000..21f21e16c --- /dev/null +++ b/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql @@ -0,0 +1,197 @@ +-- ============================================================================ +-- Direct normalized-table inserts (replaces V1's raw staging tables) +-- ============================================================================ +-- +-- Vector now inserts directly into workflow_instances/task_instances as +-- {raw_event: } rows. Triggers self-target these same tables via +-- INSERT ... ON CONFLICT DO UPDATE, guarded by pg_trigger_depth() so the +-- nested self-insert doesn't recurse. raw_event holds only the latest event +-- per row (no full history). +-- ============================================================================ + +ALTER TABLE workflow_instances ADD COLUMN raw_event JSONB; +ALTER TABLE task_instances ADD COLUMN raw_event JSONB; + +DROP TRIGGER IF EXISTS normalize_workflow_events ON workflow_events_raw; +DROP TRIGGER IF EXISTS normalize_task_events ON task_events_raw; +DROP FUNCTION IF EXISTS normalize_workflow_event(); +DROP FUNCTION IF EXISTS normalize_task_event(); + +CREATE OR REPLACE FUNCTION normalize_workflow_instance() +RETURNS TRIGGER AS $$ +DECLARE + v_id VARCHAR(255); + v_namespace VARCHAR(255); + v_name VARCHAR(255); + v_version VARCHAR(255); + v_status VARCHAR(50); + v_started_at TIMESTAMP WITH TIME ZONE; + v_ended_at TIMESTAMP WITH TIME ZONE; + v_last_update TIMESTAMP WITH TIME ZONE; + v_input JSONB; + v_output JSONB; + v_error_type VARCHAR(255); + v_error_title VARCHAR(255); + v_error_detail TEXT; + v_error_status INTEGER; + v_error_instance VARCHAR(255); + v_event_timestamp TIMESTAMP WITH TIME ZONE; +BEGIN + -- Nested call from our own INSERT below; let it proceed as-is. + IF pg_trigger_depth() > 1 THEN + RETURN NEW; + END IF; + + -- Placeholder/JPA-inserted rows have no raw_event; pass through untouched. + IF NEW.raw_event IS NULL THEN + RETURN NEW; + END IF; + + v_id := NEW.raw_event->>'instanceId'; + v_namespace := NEW.raw_event->>'workflowNamespace'; + v_name := NEW.raw_event->>'workflowName'; + v_version := NEW.raw_event->>'workflowVersion'; + v_status := NEW.raw_event->>'status'; + v_started_at := to_timestamp((NEW.raw_event->>'startTime')::numeric); + v_ended_at := to_timestamp((NEW.raw_event->>'endTime')::numeric); + v_last_update := to_timestamp((NEW.raw_event->>'lastUpdateTime')::numeric); + v_input := NEW.raw_event->'input'; + v_output := NEW.raw_event->'output'; + v_error_type := NEW.raw_event->'error'->>'type'; + v_error_title := NEW.raw_event->'error'->>'title'; + v_error_detail := NEW.raw_event->'error'->>'detail'; + v_error_status := (NEW.raw_event->'error'->>'status')::integer; + v_error_instance := NEW.raw_event->'error'->>'instance'; + v_event_timestamp := to_timestamp((NEW.raw_event->>'timestamp')::numeric); + + INSERT INTO workflow_instances ( + id, namespace, name, version, status, started_at, ended_at, last_update, + input, output, error_type, error_title, error_detail, error_status, error_instance, + last_event_time, raw_event, created_at, updated_at + ) VALUES ( + v_id, v_namespace, v_name, v_version, v_status, v_started_at, v_ended_at, v_last_update, + v_input, v_output, v_error_type, v_error_title, v_error_detail, v_error_status, v_error_instance, + v_event_timestamp, NEW.raw_event, now(), now() + ) + ON CONFLICT (id) DO UPDATE SET + -- Status/timestamps: newer event wins. Immutable fields: first write + -- wins. Terminal fields: latest non-null wins, never cleared. + status = CASE + WHEN v_event_timestamp > workflow_instances.last_event_time + THEN v_status + ELSE workflow_instances.status + END, + namespace = COALESCE(workflow_instances.namespace, v_namespace), + name = COALESCE(workflow_instances.name, v_name), + version = COALESCE(workflow_instances.version, v_version), + started_at = COALESCE(workflow_instances.started_at, v_started_at), + input = COALESCE(workflow_instances.input, v_input), + ended_at = COALESCE(v_ended_at, workflow_instances.ended_at), + output = COALESCE(v_output, workflow_instances.output), + error_type = COALESCE(v_error_type, workflow_instances.error_type), + error_title = COALESCE(v_error_title, workflow_instances.error_title), + error_detail = COALESCE(v_error_detail, workflow_instances.error_detail), + error_status = COALESCE(v_error_status, workflow_instances.error_status), + error_instance = COALESCE(v_error_instance, workflow_instances.error_instance), + last_update = GREATEST(v_last_update, workflow_instances.last_update), + last_event_time = GREATEST(v_event_timestamp, workflow_instances.last_event_time), + raw_event = NEW.raw_event, + updated_at = now(); + + RETURN NULL; +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER normalize_workflow_instances + BEFORE INSERT ON workflow_instances + FOR EACH ROW + EXECUTE FUNCTION normalize_workflow_instance(); + +CREATE OR REPLACE FUNCTION normalize_task_instance() +RETURNS TRIGGER AS $$ +DECLARE + v_instance_id VARCHAR(255); + v_task VARCHAR(255); + v_task_name VARCHAR(255); + v_status VARCHAR(50); + v_started_at TIMESTAMP WITH TIME ZONE; + v_ended_at TIMESTAMP WITH TIME ZONE; + v_input JSONB; + v_output JSONB; + v_error_type VARCHAR(255); + v_error_title VARCHAR(255); + v_error_detail TEXT; + v_error_status INTEGER; + v_error_instance VARCHAR(255); + v_event_timestamp TIMESTAMP WITH TIME ZONE; +BEGIN + IF pg_trigger_depth() > 1 THEN + RETURN NEW; + END IF; + + IF NEW.raw_event IS NULL THEN + RETURN NEW; + END IF; + + v_instance_id := NEW.raw_event->>'instanceId'; + v_task := NEW.raw_event->>'taskPosition'; + v_task_name := NEW.raw_event->>'taskName'; + v_status := NEW.raw_event->>'status'; + v_started_at := to_timestamp((NEW.raw_event->>'startTime')::numeric); + v_ended_at := to_timestamp((NEW.raw_event->>'endTime')::numeric); + v_input := NEW.raw_event->'input'; + v_output := NEW.raw_event->'output'; + v_error_type := NEW.raw_event->'error'->>'type'; + v_error_title := NEW.raw_event->'error'->>'title'; + v_error_detail := NEW.raw_event->'error'->>'detail'; + v_error_status := (NEW.raw_event->'error'->>'status')::integer; + v_error_instance := NEW.raw_event->'error'->>'instance'; + v_event_timestamp := to_timestamp((NEW.raw_event->>'timestamp')::numeric); + + -- Backfill the parent workflow row if this task event arrived first. + IF NOT EXISTS (SELECT 1 FROM workflow_instances WHERE id = v_instance_id) THEN + INSERT INTO workflow_instances (id, created_at, updated_at, last_event_time) + VALUES (v_instance_id, now(), now(), v_event_timestamp) + ON CONFLICT (id) DO NOTHING; + END IF; + + INSERT INTO task_instances ( + instance_id, task, task_name, status, started_at, ended_at, + input, output, error_type, error_title, error_detail, error_status, error_instance, + last_event_time, raw_event, created_at, updated_at + ) VALUES ( + v_instance_id, v_task, v_task_name, v_status, v_started_at, v_ended_at, + v_input, v_output, v_error_type, v_error_title, v_error_detail, v_error_status, v_error_instance, + v_event_timestamp, NEW.raw_event, now(), now() + ) + ON CONFLICT (instance_id, task) DO UPDATE SET + task_name = COALESCE(task_instances.task_name, v_task_name), + status = CASE + WHEN v_event_timestamp >= task_instances.last_event_time + THEN v_status + ELSE task_instances.status + END, + started_at = COALESCE(task_instances.started_at, v_started_at), + ended_at = COALESCE(v_ended_at, task_instances.ended_at), + input = COALESCE(task_instances.input, v_input), + output = COALESCE(v_output, task_instances.output), + error_type = COALESCE(v_error_type, task_instances.error_type), + error_title = COALESCE(v_error_title, task_instances.error_title), + error_detail = COALESCE(v_error_detail, task_instances.error_detail), + error_status = COALESCE(v_error_status, task_instances.error_status), + error_instance = COALESCE(v_error_instance, task_instances.error_instance), + last_event_time = GREATEST(v_event_timestamp, task_instances.last_event_time), + raw_event = NEW.raw_event, + updated_at = now(); + + RETURN NULL; +END; +$$ LANGUAGE plpgsql; + +CREATE TRIGGER normalize_task_instances + BEFORE INSERT ON task_instances + FOR EACH ROW + EXECUTE FUNCTION normalize_task_instance(); + +DROP TABLE IF EXISTS workflow_events_raw; +DROP TABLE IF EXISTS task_events_raw; diff --git a/data-index/scripts/e2e/verify-infrastructure.sh b/data-index/scripts/e2e/verify-infrastructure.sh index 1419d127e..358e90891 100755 --- a/data-index/scripts/e2e/verify-infrastructure.sh +++ b/data-index/scripts/e2e/verify-infrastructure.sh @@ -61,27 +61,29 @@ verify_mode1_infrastructure() { log_success "PostgreSQL is accepting connections" # 3. Verify schema initialized (tables exist) + # Vector inserts directly into workflow_instances/task_instances (no raw staging tables + # since V2__direct_normalized_inserts.sql) log_info "Checking database schema..." SCHEMA=${POSTGRES_SCHEMA:-public} TABLES=$(kubectl exec -n postgresql postgresql-0 -- psql -U dataindex -d dataindex -t -c " SELECT COUNT(*) FROM information_schema.tables WHERE table_schema = '$SCHEMA' - AND table_name IN ('workflow_instances', 'task_instances', 'workflow_events_raw', 'task_events_raw') + AND table_name IN ('workflow_instances', 'task_instances') ") - if [ "$TABLES" -eq 4 ]; then - log_success "Database schema initialized (4 tables found)" + if [ "$TABLES" -eq 2 ]; then + log_success "Database schema initialized (2 tables found)" else - log_error "Database schema incomplete (expected 4 tables, found $TABLES)" + log_error "Database schema incomplete (expected 2 tables, found $TABLES)" return 1 fi - # 4. Verify triggers exist + # 4. Verify triggers exist (self-targeting BEFORE INSERT triggers on the normalized tables) log_info "Checking database triggers..." SCHEMA=${POSTGRES_SCHEMA:-public} TRIGGERS=$(kubectl exec -n postgresql postgresql-0 -- psql -U dataindex -d dataindex -t -c " SELECT COUNT(*) FROM information_schema.triggers WHERE event_object_schema = '$SCHEMA' - AND event_object_table IN ('workflow_events_raw', 'task_events_raw') + AND event_object_table IN ('workflow_instances', 'task_instances') ") if [ "$TRIGGERS" -ge 2 ]; then log_success "Database triggers created ($TRIGGERS triggers found)" From 6f3225e1a182fcba6e9da5f7a398a4ae360bff99 Mon Sep 17 00:00:00 2001 From: Matheus Cruz Date: Mon, 28 Sep 2026 17:18:53 -0300 Subject: [PATCH 2/3] Map event fields to typed columns in Vector, drop raw_event Vector extracts and types each field in VRL instead of shipping raw_event JSONB; triggers read NEW. directly. Adds from_event to distinguish Vector inserts from direct writers (JPA), fixes a stuck-status comparison, and always keeps from_event accurate on merge. --- .../vector/mode1-postgresql/vector.yaml | 75 ++++++- .../V2__direct_normalized_inserts.sql | 183 +++++++----------- 2 files changed, 136 insertions(+), 122 deletions(-) diff --git a/data-index/collectors/vector/mode1-postgresql/vector.yaml b/data-index/collectors/vector/mode1-postgresql/vector.yaml index b0a4a4a58..f6f4d965a 100644 --- a/data-index/collectors/vector/mode1-postgresql/vector.yaml +++ b/data-index/collectors/vector/mode1-postgresql/vector.yaml @@ -93,25 +93,88 @@ transforms: type: vrl source: starts_with(string(.flow_event.eventType) ?? "", "io.serverlessworkflow.task.") - # Row inserted directly into workflow_instances/task_instances: just the - # untouched event as `raw_event` JSONB. The trigger extracts everything else. + # Row inserted directly into workflow_instances: one key per column, + # matching the OW event field names to their normalized column names. + # Vector's postgres sink maps JSON keys to columns via + # jsonb_populate_recordset, so no sink config is needed beyond this shape - + # the trigger just merges (COALESCE/GREATEST), it no longer parses JSON. build_workflow_row: type: remap inputs: - route_by_type.workflow_events source: | - . = { - "raw_event": .flow_event + e = .flow_event + row = {} + row.from_event = true + row.id = e.instanceId + row.namespace = e.workflowNamespace + row.name = e.workflowName + row.version = e.workflowVersion + row.status = e.status + row.input = e.input + row.output = e.output + + if e.startTime != null { + row.started_at = from_unix_timestamp!(to_int!(e.startTime), unit: "seconds") + } + if e.endTime != null { + row.ended_at = from_unix_timestamp!(to_int!(e.endTime), unit: "seconds") + } + if e.lastUpdateTime != null { + row.last_update = from_unix_timestamp!(to_int!(e.lastUpdateTime), unit: "seconds") + } + if e.timestamp != null { + row.last_event_time = from_unix_timestamp!(to_int!(e.timestamp), unit: "seconds") } + err = e.error + if err != null { + row.error_type = err.type + row.error_title = err.title + row.error_detail = err.detail + row.error_status = err.status + row.error_instance = err.instance + } + + . = row + + # Row inserted directly into task_instances - same column-mapping approach + # as build_workflow_row above. build_task_row: type: remap inputs: - route_by_type.task_events source: | - . = { - "raw_event": .flow_event + e = .flow_event + row = {} + row.from_event = true + row.instance_id = e.instanceId + row.task = e.taskPosition + row.task_name = e.taskName + row.status = e.status + row.input = e.input + row.output = e.output + + if e.startTime != null { + row.started_at = from_unix_timestamp!(to_int!(e.startTime), unit: "seconds") + } + if e.endTime != null { + row.ended_at = from_unix_timestamp!(to_int!(e.endTime), unit: "seconds") } + if e.timestamp != null { + row.last_event_time = from_unix_timestamp!(to_int!(e.timestamp), unit: "seconds") + } + + err = e.error + if err != null { + row.error_type = err.type + row.error_title = err.title + row.error_detail = err.detail + row.error_status = err.status + row.error_instance = err.instance + } + + . = row # Debug filter - only pass events if DEBUG_EVENTS=true # Controlled by environment variable for troubleshooting diff --git a/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql b/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql index 21f21e16c..6f312b22e 100644 --- a/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql +++ b/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql @@ -2,100 +2,81 @@ -- Direct normalized-table inserts (replaces V1's raw staging tables) -- ============================================================================ -- --- Vector now inserts directly into workflow_instances/task_instances as --- {raw_event: } rows. Triggers self-target these same tables via --- INSERT ... ON CONFLICT DO UPDATE, guarded by pg_trigger_depth() so the --- nested self-insert doesn't recurse. raw_event holds only the latest event --- per row (no full history). +-- Vector now extracts and types each field in VRL (data-index/collectors/vector/ +-- mode1-postgresql/vector.yaml) and inserts workflow_instances/task_instances +-- directly, one named column per field - Vector's postgres sink maps JSON keys +-- to columns via jsonb_populate_recordset, so no custom sink logic is needed +-- beyond the row shape VRL builds. +-- +-- Triggers self-target these same tables via INSERT ... ON CONFLICT DO UPDATE, +-- guarded by pg_trigger_depth() so the nested self-insert doesn't recurse. They +-- read NEW. directly instead of parsing JSON. +-- +-- from_event distinguishes a Vector event (needs the merge/upsert below) from +-- a direct writer like a JPA test fixture (should insert plain, untouched). +-- This isn't optional: a BEFORE INSERT trigger that RETURN NULLs to cancel the +-- original statement makes Postgres report 0 rows affected for THAT statement, +-- even though the nested INSERT below succeeds - Hibernate's optimistic-lock +-- row-count check then fails on every single em.persist(), not just conflicts. +-- Entities never need to set this column; it defaults to false. -- ============================================================================ -ALTER TABLE workflow_instances ADD COLUMN raw_event JSONB; -ALTER TABLE task_instances ADD COLUMN raw_event JSONB; - DROP TRIGGER IF EXISTS normalize_workflow_events ON workflow_events_raw; DROP TRIGGER IF EXISTS normalize_task_events ON task_events_raw; DROP FUNCTION IF EXISTS normalize_workflow_event(); DROP FUNCTION IF EXISTS normalize_task_event(); +ALTER TABLE workflow_instances ADD COLUMN from_event BOOLEAN NOT NULL DEFAULT false; +ALTER TABLE task_instances ADD COLUMN from_event BOOLEAN NOT NULL DEFAULT false; + CREATE OR REPLACE FUNCTION normalize_workflow_instance() RETURNS TRIGGER AS $$ -DECLARE - v_id VARCHAR(255); - v_namespace VARCHAR(255); - v_name VARCHAR(255); - v_version VARCHAR(255); - v_status VARCHAR(50); - v_started_at TIMESTAMP WITH TIME ZONE; - v_ended_at TIMESTAMP WITH TIME ZONE; - v_last_update TIMESTAMP WITH TIME ZONE; - v_input JSONB; - v_output JSONB; - v_error_type VARCHAR(255); - v_error_title VARCHAR(255); - v_error_detail TEXT; - v_error_status INTEGER; - v_error_instance VARCHAR(255); - v_event_timestamp TIMESTAMP WITH TIME ZONE; BEGIN -- Nested call from our own INSERT below; let it proceed as-is. IF pg_trigger_depth() > 1 THEN RETURN NEW; END IF; - -- Placeholder/JPA-inserted rows have no raw_event; pass through untouched. - IF NEW.raw_event IS NULL THEN + -- Direct writer (JPA, etc.): insert as given, no merge. + IF NOT NEW.from_event THEN RETURN NEW; END IF; - v_id := NEW.raw_event->>'instanceId'; - v_namespace := NEW.raw_event->>'workflowNamespace'; - v_name := NEW.raw_event->>'workflowName'; - v_version := NEW.raw_event->>'workflowVersion'; - v_status := NEW.raw_event->>'status'; - v_started_at := to_timestamp((NEW.raw_event->>'startTime')::numeric); - v_ended_at := to_timestamp((NEW.raw_event->>'endTime')::numeric); - v_last_update := to_timestamp((NEW.raw_event->>'lastUpdateTime')::numeric); - v_input := NEW.raw_event->'input'; - v_output := NEW.raw_event->'output'; - v_error_type := NEW.raw_event->'error'->>'type'; - v_error_title := NEW.raw_event->'error'->>'title'; - v_error_detail := NEW.raw_event->'error'->>'detail'; - v_error_status := (NEW.raw_event->'error'->>'status')::integer; - v_error_instance := NEW.raw_event->'error'->>'instance'; - v_event_timestamp := to_timestamp((NEW.raw_event->>'timestamp')::numeric); - INSERT INTO workflow_instances ( id, namespace, name, version, status, started_at, ended_at, last_update, input, output, error_type, error_title, error_detail, error_status, error_instance, - last_event_time, raw_event, created_at, updated_at + last_event_time, from_event, created_at, updated_at ) VALUES ( - v_id, v_namespace, v_name, v_version, v_status, v_started_at, v_ended_at, v_last_update, - v_input, v_output, v_error_type, v_error_title, v_error_detail, v_error_status, v_error_instance, - v_event_timestamp, NEW.raw_event, now(), now() + NEW.id, NEW.namespace, NEW.name, NEW.version, NEW.status, NEW.started_at, NEW.ended_at, NEW.last_update, + NEW.input, NEW.output, NEW.error_type, NEW.error_title, NEW.error_detail, NEW.error_status, NEW.error_instance, + NEW.last_event_time, true, now(), now() ) ON CONFLICT (id) DO UPDATE SET -- Status/timestamps: newer event wins. Immutable fields: first write -- wins. Terminal fields: latest non-null wins, never cleared. status = CASE - WHEN v_event_timestamp > workflow_instances.last_event_time - THEN v_status + WHEN NEW.last_event_time >= workflow_instances.last_event_time + THEN NEW.status ELSE workflow_instances.status END, - namespace = COALESCE(workflow_instances.namespace, v_namespace), - name = COALESCE(workflow_instances.name, v_name), - version = COALESCE(workflow_instances.version, v_version), - started_at = COALESCE(workflow_instances.started_at, v_started_at), - input = COALESCE(workflow_instances.input, v_input), - ended_at = COALESCE(v_ended_at, workflow_instances.ended_at), - output = COALESCE(v_output, workflow_instances.output), - error_type = COALESCE(v_error_type, workflow_instances.error_type), - error_title = COALESCE(v_error_title, workflow_instances.error_title), - error_detail = COALESCE(v_error_detail, workflow_instances.error_detail), - error_status = COALESCE(v_error_status, workflow_instances.error_status), - error_instance = COALESCE(v_error_instance, workflow_instances.error_instance), - last_update = GREATEST(v_last_update, workflow_instances.last_update), - last_event_time = GREATEST(v_event_timestamp, workflow_instances.last_event_time), - raw_event = NEW.raw_event, + namespace = COALESCE(workflow_instances.namespace, NEW.namespace), + name = COALESCE(workflow_instances.name, NEW.name), + version = COALESCE(workflow_instances.version, NEW.version), + started_at = COALESCE(workflow_instances.started_at, NEW.started_at), + input = COALESCE(workflow_instances.input, NEW.input), + ended_at = COALESCE(NEW.ended_at, workflow_instances.ended_at), + output = COALESCE(NEW.output, workflow_instances.output), + error_type = COALESCE(NEW.error_type, workflow_instances.error_type), + error_title = COALESCE(NEW.error_title, workflow_instances.error_title), + error_detail = COALESCE(NEW.error_detail, workflow_instances.error_detail), + error_status = COALESCE(NEW.error_status, workflow_instances.error_status), + error_instance = COALESCE(NEW.error_instance, workflow_instances.error_instance), + last_update = GREATEST(NEW.last_update, workflow_instances.last_update), + last_event_time = GREATEST(NEW.last_event_time, workflow_instances.last_event_time), + -- The row may have been created first by the task backfill below + -- (from_event defaults false there); a real event landing on it now + -- means it's event-sourced from here on. + from_event = true, updated_at = now(); RETURN NULL; @@ -109,79 +90,49 @@ CREATE TRIGGER normalize_workflow_instances CREATE OR REPLACE FUNCTION normalize_task_instance() RETURNS TRIGGER AS $$ -DECLARE - v_instance_id VARCHAR(255); - v_task VARCHAR(255); - v_task_name VARCHAR(255); - v_status VARCHAR(50); - v_started_at TIMESTAMP WITH TIME ZONE; - v_ended_at TIMESTAMP WITH TIME ZONE; - v_input JSONB; - v_output JSONB; - v_error_type VARCHAR(255); - v_error_title VARCHAR(255); - v_error_detail TEXT; - v_error_status INTEGER; - v_error_instance VARCHAR(255); - v_event_timestamp TIMESTAMP WITH TIME ZONE; BEGIN IF pg_trigger_depth() > 1 THEN RETURN NEW; END IF; - IF NEW.raw_event IS NULL THEN + -- Direct writer (JPA, etc.): insert as given, no merge. + IF NOT NEW.from_event THEN RETURN NEW; END IF; - v_instance_id := NEW.raw_event->>'instanceId'; - v_task := NEW.raw_event->>'taskPosition'; - v_task_name := NEW.raw_event->>'taskName'; - v_status := NEW.raw_event->>'status'; - v_started_at := to_timestamp((NEW.raw_event->>'startTime')::numeric); - v_ended_at := to_timestamp((NEW.raw_event->>'endTime')::numeric); - v_input := NEW.raw_event->'input'; - v_output := NEW.raw_event->'output'; - v_error_type := NEW.raw_event->'error'->>'type'; - v_error_title := NEW.raw_event->'error'->>'title'; - v_error_detail := NEW.raw_event->'error'->>'detail'; - v_error_status := (NEW.raw_event->'error'->>'status')::integer; - v_error_instance := NEW.raw_event->'error'->>'instance'; - v_event_timestamp := to_timestamp((NEW.raw_event->>'timestamp')::numeric); - -- Backfill the parent workflow row if this task event arrived first. - IF NOT EXISTS (SELECT 1 FROM workflow_instances WHERE id = v_instance_id) THEN + IF NOT EXISTS (SELECT 1 FROM workflow_instances WHERE id = NEW.instance_id) THEN INSERT INTO workflow_instances (id, created_at, updated_at, last_event_time) - VALUES (v_instance_id, now(), now(), v_event_timestamp) + VALUES (NEW.instance_id, now(), now(), NEW.last_event_time) ON CONFLICT (id) DO NOTHING; END IF; INSERT INTO task_instances ( instance_id, task, task_name, status, started_at, ended_at, input, output, error_type, error_title, error_detail, error_status, error_instance, - last_event_time, raw_event, created_at, updated_at + last_event_time, from_event, created_at, updated_at ) VALUES ( - v_instance_id, v_task, v_task_name, v_status, v_started_at, v_ended_at, - v_input, v_output, v_error_type, v_error_title, v_error_detail, v_error_status, v_error_instance, - v_event_timestamp, NEW.raw_event, now(), now() + NEW.instance_id, NEW.task, NEW.task_name, NEW.status, NEW.started_at, NEW.ended_at, + NEW.input, NEW.output, NEW.error_type, NEW.error_title, NEW.error_detail, NEW.error_status, NEW.error_instance, + NEW.last_event_time, true, now(), now() ) ON CONFLICT (instance_id, task) DO UPDATE SET - task_name = COALESCE(task_instances.task_name, v_task_name), + task_name = COALESCE(task_instances.task_name, NEW.task_name), status = CASE - WHEN v_event_timestamp >= task_instances.last_event_time - THEN v_status + WHEN NEW.last_event_time >= task_instances.last_event_time + THEN NEW.status ELSE task_instances.status END, - started_at = COALESCE(task_instances.started_at, v_started_at), - ended_at = COALESCE(v_ended_at, task_instances.ended_at), - input = COALESCE(task_instances.input, v_input), - output = COALESCE(v_output, task_instances.output), - error_type = COALESCE(v_error_type, task_instances.error_type), - error_title = COALESCE(v_error_title, task_instances.error_title), - error_detail = COALESCE(v_error_detail, task_instances.error_detail), - error_status = COALESCE(v_error_status, task_instances.error_status), - error_instance = COALESCE(v_error_instance, task_instances.error_instance), - last_event_time = GREATEST(v_event_timestamp, task_instances.last_event_time), - raw_event = NEW.raw_event, + started_at = COALESCE(task_instances.started_at, NEW.started_at), + ended_at = COALESCE(NEW.ended_at, task_instances.ended_at), + input = COALESCE(task_instances.input, NEW.input), + output = COALESCE(NEW.output, task_instances.output), + error_type = COALESCE(NEW.error_type, task_instances.error_type), + error_title = COALESCE(NEW.error_title, task_instances.error_title), + error_detail = COALESCE(NEW.error_detail, task_instances.error_detail), + error_status = COALESCE(NEW.error_status, task_instances.error_status), + error_instance = COALESCE(NEW.error_instance, task_instances.error_instance), + last_event_time = GREATEST(NEW.last_event_time, task_instances.last_event_time), updated_at = now(); RETURN NULL; From c0decf734ec43c51996915fd36685c414befc7c1 Mon Sep 17 00:00:00 2001 From: Matheus Cruz Date: Mon, 28 Sep 2026 17:38:12 -0300 Subject: [PATCH 3/3] Squash V1+V2 migrations into a single initial schema --- CLAUDE.md | 5 +- .../data-index-storage-migrations/README.md | 127 +++-- .../db/migration/V1__initial_schema.sql | 438 +++++------------- .../V2__direct_normalized_inserts.sql | 148 ------ 4 files changed, 176 insertions(+), 542 deletions(-) delete mode 100644 data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql diff --git a/CLAUDE.md b/CLAUDE.md index db84332b9..06beb7483 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -1342,8 +1342,7 @@ curl http://localhost:9200/_transform/workflow-instances-transform/_stats - `data-index-storage-postgresql/src/main/java/.../entity/` - JPA entities - `data-index-storage-postgresql/src/main/java/.../mapper/` - MapStruct mappers - `data-index-storage-migrations/src/main/resources/db/migration/` - Flyway migrations -- `V1__initial_schema.sql` - Original schema (raw staging tables, now removed) -- `V2__direct_normalized_inserts.sql` - Current triggers: self-targeting `BEFORE INSERT` on `workflow_instances`/`task_instances` +- `V1__initial_schema.sql` - Normalized tables + self-targeting `BEFORE INSERT` triggers on `workflow_instances`/`task_instances` **Code (MODE 2 - Elasticsearch):** - `data-index-storage-elasticsearch/src/main/java/.../` - Storage implementation @@ -1438,7 +1437,7 @@ curl http://localhost:9200/_transform/workflow-instances-transform/_stats → MODE 2: Write `@QuarkusTest` integration test with Elasticsearch profile + wait for transforms **"Where are the database triggers?"** -→ MODE 1: `data-index-storage-migrations/.../V2__direct_normalized_inserts.sql` (defined directly on `workflow_instances`/`task_instances`) +→ MODE 1: `data-index-storage-migrations/.../V1__initial_schema.sql` (defined directly on `workflow_instances`/`task_instances`) → MODE 2: Not applicable (uses Elasticsearch transforms) **"Where are the Elasticsearch transforms?"** diff --git a/data-index/data-index-storage/data-index-storage-migrations/README.md b/data-index/data-index-storage/data-index-storage-migrations/README.md index 5704b3b1e..3458212d8 100644 --- a/data-index/data-index-storage/data-index-storage-migrations/README.md +++ b/data-index/data-index-storage/data-index-storage-migrations/README.md @@ -14,12 +14,12 @@ Flyway migration scripts for the Data Index PostgreSQL storage backend: ``` Quarkus Flow (quarkus-flow 0.9.0+) ↓ (structured JSON events with epoch timestamps → stdout) -Vector DaemonSet (postgres sink) - ├─→ workflow_instances (raw_event JSONB column only) - └─→ task_instances (raw_event JSONB column only) +Vector DaemonSet + ↓ (VRL extracts + types each field) + ├─→ workflow_instances (typed columns, from_event = true) + └─→ task_instances (typed columns, from_event = true) ↓ (BEFORE INSERT triggers, self-targeting) PostgreSQL Triggers - ├─→ Extract fields from raw_event JSONB ├─→ Handle out-of-order events (COALESCE/GREATEST) ├─→ UPDATE existing row + RETURN NULL, or fill in NEW + RETURN NEW ↓ (GraphQL queries) @@ -30,93 +30,87 @@ Data Index GraphQL API ### V1__initial_schema.sql -Original schema for Data Index v1.0.0: creates `workflow_instances` and `task_instances` -(without `raw_event`), plus the now-removed raw staging tables and their triggers. +Creates `workflow_instances` and `task_instances`, plus `normalize_workflow_instance()` / +`normalize_task_instance()` — `BEFORE INSERT` triggers on those same tables, self-targeting +`INSERT ... ON CONFLICT DO UPDATE`, guarded by `pg_trigger_depth()` to avoid recursion. -### V2__direct_normalized_inserts.sql +## Row Shape (Vector VRL → Normalized Tables) -- Adds `raw_event JSONB` to `workflow_instances` and `task_instances`. -- Drops `workflow_events_raw`, `task_events_raw`, and their triggers/functions. -- Adds `normalize_workflow_instance()` / `normalize_task_instance()`, `BEFORE INSERT` on - `workflow_instances` / `task_instances` respectively — self-targeting `INSERT ... ON CONFLICT - DO UPDATE`, guarded by `pg_trigger_depth()` to avoid recursion. +Vector's VRL transforms (`build_workflow_row`/`build_task_row` in +`data-index/collectors/vector/mode1-postgresql/vector.yaml`) extract and type each field +from the Quarkus Flow event before sending it — Vector's `postgres` sink maps JSON keys to +columns via `jsonb_populate_recordset`, so the row Vector sends already has one key per +column. The trigger reads `NEW.` directly; it does not parse JSON. -## Field Mappings (raw_event JSONB → Normalized Tables) +### Workflow Events → workflow_instances -### Workflow Events (raw_event JSONB → workflow_instances) - -`normalize_workflow_instance()` extracts: - -| JSONB Path (raw_event->>) | workflow_instances Column | Type | Conversion | -|----------------------------|----------------------------|------|------------| +| Event field | workflow_instances Column | Type | Conversion (in VRL) | +|-------------|----------------------------|------|----------------------| | `instanceId` | `id` | VARCHAR(255) | Direct | | `workflowNamespace` | `namespace` | VARCHAR(255) | Direct | | `workflowName` | `name` | VARCHAR(255) | Direct | | `workflowVersion` | `version` | VARCHAR(255) | Direct | | `status` | `status` | VARCHAR(50) | Direct | -| `startTime` | `started_at` | TIMESTAMPTZ | `to_timestamp(::numeric)` | -| `endTime` | `ended_at` | TIMESTAMPTZ | `to_timestamp(::numeric)` | -| `lastUpdateTime` | `last_update` | TIMESTAMPTZ | `to_timestamp(::numeric)` | -| `timestamp` | `last_event_time` | TIMESTAMPTZ | `to_timestamp(::numeric)` (drives idempotency) | -| `input` | `input` | JSONB | Direct (`->`) | -| `output` | `output` | JSONB | Direct (`->`) | -| `error->>'type'` | `error_type` | VARCHAR(255) | Nested | -| `error->>'title'` | `error_title` | VARCHAR(255) | Nested | -| `error->>'detail'` | `error_detail` | TEXT | Nested | -| `error->>'status'` | `error_status` | INTEGER | Nested + cast | -| `error->>'instance'` | `error_instance` | VARCHAR(255) | Nested | +| `startTime` | `started_at` | TIMESTAMPTZ | `from_unix_timestamp()` | +| `endTime` | `ended_at` | TIMESTAMPTZ | `from_unix_timestamp()` | +| `lastUpdateTime` | `last_update` | TIMESTAMPTZ | `from_unix_timestamp()` | +| `timestamp` | `last_event_time` | TIMESTAMPTZ | `from_unix_timestamp()` (drives idempotency) | +| `input` | `input` | JSONB | Direct | +| `output` | `output` | JSONB | Direct | +| `error.type` | `error_type` | VARCHAR(255) | Flattened | +| `error.title` | `error_title` | VARCHAR(255) | Flattened | +| `error.detail` | `error_detail` | TEXT | Flattened | +| `error.status` | `error_status` | INTEGER | Flattened | +| `error.instance` | `error_instance` | VARCHAR(255) | Flattened | +| *(set by VRL)* | `from_event` | BOOLEAN | Always `true` for Vector rows | **Auto-populated:** `created_at`/`updated_at` (`now()`, set by the trigger). -### Task Events (raw_event JSONB → task_instances) - -`normalize_task_instance()` extracts: +### Task Events → task_instances -| JSONB Path (raw_event->>) | task_instances Column | Type | Conversion | -|----------------------------|------------------------|------|------------| +| Event field | task_instances Column | Type | Conversion (in VRL) | +|-------------|-------------------------|------|-----------------------| | `instanceId` | `instance_id` | VARCHAR(255) | Direct (PK part 1, FK to workflow_instances) | | `taskPosition` | `task` | VARCHAR(255) | Direct (PK part 2, JSON Pointer e.g. `/do/1/initialize`) | | `taskName` | `task_name` | VARCHAR(255) | Direct | | `status` | `status` | VARCHAR(50) | Direct | -| `startTime` | `started_at` | TIMESTAMPTZ | `to_timestamp(::numeric)` | -| `endTime` | `ended_at` | TIMESTAMPTZ | `to_timestamp(::numeric)` | -| `timestamp` | `last_event_time` | TIMESTAMPTZ | `to_timestamp(::numeric)` (drives idempotency) | -| `input` | `input` | JSONB | Direct (`->`) | -| `output` | `output` | JSONB | Direct (`->`) | -| `error->>'type'` | `error_type` | VARCHAR(255) | Nested | -| `error->>'title'` | `error_title` | VARCHAR(255) | Nested | -| `error->>'detail'` | `error_detail` | TEXT | Nested | -| `error->>'status'` | `error_status` | INTEGER | Nested + cast | -| `error->>'instance'` | `error_instance` | VARCHAR(255) | Nested | +| `startTime` | `started_at` | TIMESTAMPTZ | `from_unix_timestamp()` | +| `endTime` | `ended_at` | TIMESTAMPTZ | `from_unix_timestamp()` | +| `timestamp` | `last_event_time` | TIMESTAMPTZ | `from_unix_timestamp()` (drives idempotency) | +| `input` | `input` | JSONB | Direct | +| `output` | `output` | JSONB | Direct | +| `error.type` | `error_type` | VARCHAR(255) | Flattened | +| `error.title` | `error_title` | VARCHAR(255) | Flattened | +| `error.detail` | `error_detail` | TEXT | Flattened | +| `error.status` | `error_status` | INTEGER | Flattened | +| `error.instance` | `error_instance` | VARCHAR(255) | Flattened | +| *(set by VRL)* | `from_event` | BOOLEAN | Always `true` for Vector rows | **Auto-populated:** `created_at`/`updated_at` (`now()`, set by the trigger). **Note:** no `task_execution_id` column — `TaskExecution.id` is derived as `instanceId + ":" + task`. **Out-of-Order Handling:** the task trigger backfills a placeholder `workflow_instances` row -(`raw_event` left `NULL`) if the parent doesn't exist yet, before processing the task event. +(`from_event` left at its default `false`) if the parent doesn't exist yet, before processing +the task event. ## Usage ### Local Development -Apply migrations manually: +Apply the migration manually: ```bash kubectl exec -n postgresql postgresql-0 -- env PGPASSWORD=dataindex123 \ psql -U dataindex -d dataindex -f /path/to/V1__initial_schema.sql -kubectl exec -n postgresql postgresql-0 -- env PGPASSWORD=dataindex123 \ - psql -U dataindex -d dataindex -f /path/to/V2__direct_normalized_inserts.sql ``` ### Kubernetes Operator The Data Index operator uses Flyway to manage migrations automatically for PostgreSQL storage. -**Upgrade safety:** -- `V2__` is additive on top of an already-provisioned `V1__` schema, including existing rows. -- Foreign keys use `ON DELETE CASCADE`. -- Triggers are idempotent regardless of how many times the same event is (re)inserted. +**Idempotency:** Triggers are idempotent regardless of how many times the same event is +(re)inserted, and foreign keys use `ON DELETE CASCADE`. ## Trigger-Based Normalization @@ -124,36 +118,39 @@ No Event Processor needed — PostgreSQL triggers handle normalization directly. ### How It Works -1. Vector's `postgres` sink inserts into `workflow_instances`/`task_instances` with only `raw_event` set. +1. Vector's VRL builds a row with typed columns (`from_event = true`) and the `postgres` sink inserts it into `workflow_instances`/`task_instances`. 2. The `BEFORE INSERT` trigger fires on that same table (`pg_trigger_depth()` guard: a nested call from step 4 passes through unchanged). -3. Extracts fields from `raw_event` (`NULL` `raw_event` also passes through unchanged — placeholders, JPA-seeded test rows). -4. Runs its own `INSERT ... ON CONFLICT (id) DO UPDATE` against the same table (COALESCE/GREATEST rules). +3. If `from_event` is `false` (a direct writer like a JPA test fixture, or the task trigger's placeholder backfill), the trigger passes the row through unchanged — no merge. +4. Otherwise it runs its own `INSERT ... ON CONFLICT (id) DO UPDATE` against the same table (COALESCE/GREATEST rules). 5. `RETURN NULL`s to cancel the original pending insert — step 4 already applied the event. `ON CONFLICT` (not a manual `UPDATE`-then-`INSERT`) is required for atomicity: Vector's `postgres_workflow`/`postgres_task` sinks write concurrently, and a check-then-act pattern is race-prone. +`from_event` exists because `RETURN NULL` makes Postgres report 0 rows affected for the +*original* statement even when the nested insert in step 4 succeeds — Hibernate's +optimistic-lock row-count check would otherwise fail on every `em.persist()`, not just on +actual conflicts. JPA entities never need to set it; it defaults to `false`. + ### Advantages - **Real-time**: no polling or batch delays - **Out-of-order handling**: atomic `ON CONFLICT DO UPDATE` with COALESCE/GREATEST, safe under concurrency - **Idempotent**: same event can be reinserted safely -- **Simpler**: two tables instead of four -- **Debugging**: most recent raw event kept per row via `raw_event` +- **Simpler**: two tables, one migration, no raw JSON parsing in PL/pgSQL ### Example ```sql --- Vector inserts: -INSERT INTO workflow_instances (raw_event) -VALUES ('{"instanceId":"01KPY...","workflowName":"simple-set","status":"RUNNING",...}'::jsonb); +-- Vector inserts (typed columns, from_event always true): +INSERT INTO workflow_instances (id, name, status, from_event, ...) +VALUES ('01KPY...', 'simple-set', 'RUNNING', true, ...); --- Trigger extracts fields and upserts workflow_instances in place, then cancels this insert. +-- Trigger merges into the existing row (or inserts fresh) and cancels this insert. ``` ## Notes -- `raw_event` holds only the most recent event per row (not full history). -- Timestamps: epoch seconds from Quarkus Flow, converted via `to_timestamp()`. +- Timestamps: epoch seconds from Quarkus Flow, converted to `TIMESTAMPTZ` by Vector's VRL (`from_unix_timestamp()`). - Normalized columns are indexed for GraphQL queries. -- Deleting a workflow instance cascades to its tasks. \ No newline at end of file +- Deleting a workflow instance cascades to its tasks. diff --git a/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V1__initial_schema.sql b/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V1__initial_schema.sql index 3108ea7ee..ec9e12ff5 100644 --- a/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V1__initial_schema.sql +++ b/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V1__initial_schema.sql @@ -1,40 +1,24 @@ -- ============================================================================ --- Data Index v1.1.0 - FluentBit + Trigger-based Normalization with Idempotency +-- Data Index - Normalized tables with self-targeting trigger normalization -- ============================================================================ -- --- This schema includes: --- - Raw staging tables (FluentBit pgsql plugin fixed schema) --- - Normalized tables with idempotency fields --- - Trigger functions with field-level idempotency logic --- - Out-of-order event handling --- - Event replay safety +-- Vector extracts and types each field in VRL (data-index/collectors/vector/ +-- mode1-postgresql/vector.yaml) and inserts workflow_instances/task_instances +-- directly, one named column per field - Vector's postgres sink maps JSON keys +-- to columns via jsonb_populate_recordset, so no custom sink logic is needed +-- beyond the row shape VRL builds. -- --- ============================================================================ - --- ============================================================================ --- RAW STAGING TABLES (FluentBit pgsql plugin fixed schema) --- ============================================================================ - -CREATE TABLE IF NOT EXISTS workflow_events_raw ( - tag TEXT, - time TIMESTAMP WITH TIME ZONE, - data JSONB -); - -CREATE INDEX IF NOT EXISTS idx_workflow_events_raw_time ON workflow_events_raw (time DESC); -CREATE INDEX IF NOT EXISTS idx_workflow_events_raw_tag ON workflow_events_raw (tag); - -CREATE TABLE IF NOT EXISTS task_events_raw ( - tag TEXT, - time TIMESTAMP WITH TIME ZONE, - data JSONB -); - -CREATE INDEX IF NOT EXISTS idx_task_events_raw_time ON task_events_raw (time DESC); -CREATE INDEX IF NOT EXISTS idx_task_events_raw_tag ON task_events_raw (tag); - --- ============================================================================ --- NORMALIZED TABLES (GraphQL API queries these) +-- Triggers self-target these same tables via INSERT ... ON CONFLICT DO UPDATE, +-- guarded by pg_trigger_depth() so the nested self-insert doesn't recurse. They +-- read NEW. directly instead of parsing JSON. +-- +-- from_event distinguishes a Vector event (needs the merge/upsert below) from +-- a direct writer like a JPA test fixture (should insert plain, untouched). +-- This isn't optional: a BEFORE INSERT trigger that RETURN NULLs to cancel the +-- original statement makes Postgres report 0 rows affected for THAT statement, +-- even though the nested INSERT below succeeds - Hibernate's optimistic-lock +-- row-count check then fails on every single em.persist(), not just conflicts. +-- Entities never need to set this column; it defaults to false. -- ============================================================================ CREATE TABLE IF NOT EXISTS workflow_instances ( @@ -54,6 +38,7 @@ CREATE TABLE IF NOT EXISTS workflow_instances ( error_detail TEXT, error_status INTEGER, error_instance VARCHAR(255), + from_event BOOLEAN NOT NULL DEFAULT false, created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW() ); @@ -78,6 +63,7 @@ CREATE TABLE IF NOT EXISTS task_instances ( error_detail TEXT, error_status INTEGER, error_instance VARCHAR(255), + from_event BOOLEAN NOT NULL DEFAULT false, created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), -- Composite primary key: (instance_id, task) uniquely identifies a task @@ -90,317 +76,117 @@ CREATE INDEX IF NOT EXISTS idx_task_instances_instance_id ON task_instances (ins CREATE INDEX IF NOT EXISTS idx_task_instances_status ON task_instances (status); CREATE INDEX IF NOT EXISTS idx_task_instances_last_event_time ON task_instances (last_event_time DESC); --- ============================================================================ --- TRIGGER FUNCTIONS (Extract from JSONB and normalize with idempotency) --- ============================================================================ - --- Function to normalize workflow events with field-level idempotency -CREATE OR REPLACE FUNCTION normalize_workflow_event() +CREATE OR REPLACE FUNCTION normalize_workflow_instance() RETURNS TRIGGER AS $$ -DECLARE - event_timestamp TIMESTAMP WITH TIME ZONE; BEGIN - -- Extract event timestamp from JSONB data - -- Quarkus Flow uses 'timestamp' field (epoch-seconds format) - event_timestamp := to_timestamp((NEW.data->>'timestamp')::numeric); + -- Nested call from our own INSERT below; let it proceed as-is. + IF pg_trigger_depth() > 1 THEN + RETURN NEW; + END IF; + + -- Direct writer (JPA, etc.): insert as given, no merge. + IF NOT NEW.from_event THEN + RETURN NEW; + END IF; - -- Upsert with field-level idempotency logic INSERT INTO workflow_instances ( - id, - namespace, - name, - version, - status, - started_at, - ended_at, - last_update, - input, - output, - error_type, - error_title, - error_detail, - error_status, - error_instance, - last_event_time, - created_at, - updated_at + id, namespace, name, version, status, started_at, ended_at, last_update, + input, output, error_type, error_title, error_detail, error_status, error_instance, + last_event_time, from_event, created_at, updated_at ) VALUES ( - NEW.data->>'instanceId', - NEW.data->>'workflowNamespace', - NEW.data->>'workflowName', - NEW.data->>'workflowVersion', - NEW.data->>'status', - to_timestamp((NEW.data->>'startTime')::numeric), - to_timestamp((NEW.data->>'endTime')::numeric), - to_timestamp((NEW.data->>'lastUpdateTime')::numeric), - NEW.data->'input', - NEW.data->'output', - NEW.data->'error'->>'type', - NEW.data->'error'->>'title', - NEW.data->'error'->>'detail', - (NEW.data->'error'->>'status')::integer, - NEW.data->'error'->>'instance', - event_timestamp, - NEW.time, - NEW.time + NEW.id, NEW.namespace, NEW.name, NEW.version, NEW.status, NEW.started_at, NEW.ended_at, NEW.last_update, + NEW.input, NEW.output, NEW.error_type, NEW.error_title, NEW.error_detail, NEW.error_status, NEW.error_instance, + NEW.last_event_time, true, now(), now() ) ON CONFLICT (id) DO UPDATE SET - -- Status: Use event timestamp to determine which status wins - -- If incoming event is newer, use its status; otherwise keep existing + -- Status/timestamps: newer event wins. Immutable fields: first write + -- wins. Terminal fields: latest non-null wins, never cleared. status = CASE - WHEN event_timestamp > workflow_instances.last_event_time - THEN EXCLUDED.status + WHEN NEW.last_event_time >= workflow_instances.last_event_time + THEN NEW.status ELSE workflow_instances.status END, - - -- Immutable fields: First event wins (never overwrite if already set) - -- These are set by workflow.started event and should never change - namespace = COALESCE(workflow_instances.namespace, EXCLUDED.namespace), - name = COALESCE(workflow_instances.name, EXCLUDED.name), - version = COALESCE(workflow_instances.version, EXCLUDED.version), - started_at = COALESCE(workflow_instances.started_at, EXCLUDED.started_at), - input = COALESCE(workflow_instances.input, EXCLUDED.input), - - -- Terminal fields: Preserve if already set (completion data) - -- Once a workflow completes/faults, these fields should not be cleared - ended_at = COALESCE(EXCLUDED.ended_at, workflow_instances.ended_at), - output = COALESCE(EXCLUDED.output, workflow_instances.output), - error_type = COALESCE(EXCLUDED.error_type, workflow_instances.error_type), - error_title = COALESCE(EXCLUDED.error_title, workflow_instances.error_title), - error_detail = COALESCE(EXCLUDED.error_detail, workflow_instances.error_detail), - error_status = COALESCE(EXCLUDED.error_status, workflow_instances.error_status), - error_instance = COALESCE(EXCLUDED.error_instance, workflow_instances.error_instance), - - -- last_update: Always take newer value - last_update = GREATEST( - COALESCE(EXCLUDED.last_update, workflow_instances.last_update), - COALESCE(workflow_instances.last_update, EXCLUDED.last_update) - ), - - -- Timestamp tracking: Keep latest event timestamp - last_event_time = GREATEST(event_timestamp, workflow_instances.last_event_time), - - -- Audit: Always update - updated_at = NEW.time; - - -- Return NEW to keep the raw event in staging table - RETURN NEW; + namespace = COALESCE(workflow_instances.namespace, NEW.namespace), + name = COALESCE(workflow_instances.name, NEW.name), + version = COALESCE(workflow_instances.version, NEW.version), + started_at = COALESCE(workflow_instances.started_at, NEW.started_at), + input = COALESCE(workflow_instances.input, NEW.input), + ended_at = COALESCE(NEW.ended_at, workflow_instances.ended_at), + output = COALESCE(NEW.output, workflow_instances.output), + error_type = COALESCE(NEW.error_type, workflow_instances.error_type), + error_title = COALESCE(NEW.error_title, workflow_instances.error_title), + error_detail = COALESCE(NEW.error_detail, workflow_instances.error_detail), + error_status = COALESCE(NEW.error_status, workflow_instances.error_status), + error_instance = COALESCE(NEW.error_instance, workflow_instances.error_instance), + last_update = GREATEST(NEW.last_update, workflow_instances.last_update), + last_event_time = GREATEST(NEW.last_event_time, workflow_instances.last_event_time), + -- The row may have been created first by the task backfill below + -- (from_event defaults false there); a real event landing on it now + -- means it's event-sourced from here on. + from_event = true, + updated_at = now(); + + RETURN NULL; END; $$ LANGUAGE plpgsql; --- Function to normalize task lifecycle events into task_instances. --- --- Each logical task execution is uniquely identified by (instance_id, task). --- Quarkus Flow generates different taskExecutionId values for each lifecycle event --- (task.started, task.completed, etc.), so we cannot use it as the primary key. --- Instead, we use the composite key (instance_id, task) which remains stable --- across all events for the same task execution. --- --- Use INSERT ... ON CONFLICT (instance_id, task) DO UPDATE instead of an --- UPDATE-first / INSERT-if-not-found pattern. The latter is race-prone under --- concurrent ingestion and can produce duplicate-key violations. --- --- Repeated task positions (e.g., loop iterations) remain separate because --- each iteration has a distinct task value. -CREATE OR REPLACE FUNCTION normalize_task_event() +CREATE TRIGGER normalize_workflow_instances + BEFORE INSERT ON workflow_instances + FOR EACH ROW + EXECUTE FUNCTION normalize_workflow_instance(); + +CREATE OR REPLACE FUNCTION normalize_task_instance() RETURNS TRIGGER AS $$ -DECLARE - event_timestamp TIMESTAMP WITH TIME ZONE; BEGIN - -- Extract event timestamp from JSONB data - event_timestamp := to_timestamp((NEW.data->>'timestamp')::numeric); - - -- Fast path: - -- Try to normalize the task event directly. In the common case, the workflow - -- instance already exists because workflow events have already created it. - BEGIN - INSERT INTO task_instances ( - instance_id, - task_name, - task, - status, - started_at, - ended_at, - input, - output, - error_type, - error_title, - error_detail, - error_status, - error_instance, - last_event_time, - created_at, - updated_at - ) VALUES ( - NEW.data->>'instanceId', - NEW.data->>'taskName', - NEW.data->>'taskPosition', - NEW.data->>'status', - CASE - WHEN NEW.data ? 'startTime' - THEN to_timestamp((NEW.data->>'startTime')::numeric) - ELSE NULL - END, - CASE - WHEN NEW.data ? 'endTime' - THEN to_timestamp((NEW.data->>'endTime')::numeric) - ELSE NULL - END, - NEW.data->'input', - NEW.data->'output', - NEW.data->'error'->>'type', - NEW.data->'error'->>'title', - NEW.data->'error'->>'detail', - CASE - WHEN NEW.data->'error' ? 'status' - THEN (NEW.data->'error'->>'status')::integer - ELSE NULL - END, - NEW.data->'error'->>'instance', - event_timestamp, - NEW.time, - NEW.time - ) - ON CONFLICT (instance_id, task) DO UPDATE SET - -- Keep first non-null values (immutable fields) - task_name = COALESCE(task_instances.task_name, EXCLUDED.task_name), - - -- Update status based on event timestamp (latest wins) - status = CASE - WHEN EXCLUDED.last_event_time >= task_instances.last_event_time - THEN EXCLUDED.status - ELSE task_instances.status - END, - - -- Keep first start time, update end time with latest non-null value - started_at = COALESCE(task_instances.started_at, EXCLUDED.started_at), - ended_at = COALESCE(EXCLUDED.ended_at, task_instances.ended_at), - - -- Keep first input, update output with latest non-null value - input = COALESCE(task_instances.input, EXCLUDED.input), - output = COALESCE(EXCLUDED.output, task_instances.output), - - -- Update error fields with latest non-null values - error_type = COALESCE(EXCLUDED.error_type, task_instances.error_type), - error_title = COALESCE(EXCLUDED.error_title, task_instances.error_title), - error_detail = COALESCE(EXCLUDED.error_detail, task_instances.error_detail), - error_status = COALESCE(EXCLUDED.error_status, task_instances.error_status), - error_instance = COALESCE(EXCLUDED.error_instance, task_instances.error_instance), - - -- Track latest event time and update timestamp - last_event_time = GREATEST(EXCLUDED.last_event_time, task_instances.last_event_time), - updated_at = EXCLUDED.updated_at; - - EXCEPTION WHEN foreign_key_violation THEN - -- Slow fallback: - -- A task event arrived before its workflow event. Create a minimal workflow - -- placeholder to satisfy fk_task_instance_workflow, then retry the task upsert. - INSERT INTO workflow_instances ( - id, - created_at, - updated_at, - last_event_time - ) - VALUES ( - NEW.data->>'instanceId', - NEW.time, - NEW.time, - event_timestamp - ) + IF pg_trigger_depth() > 1 THEN + RETURN NEW; + END IF; + + -- Direct writer (JPA, etc.): insert as given, no merge. + IF NOT NEW.from_event THEN + RETURN NEW; + END IF; + + -- Backfill the parent workflow row if this task event arrived first. + IF NOT EXISTS (SELECT 1 FROM workflow_instances WHERE id = NEW.instance_id) THEN + INSERT INTO workflow_instances (id, created_at, updated_at, last_event_time) + VALUES (NEW.instance_id, now(), now(), NEW.last_event_time) ON CONFLICT (id) DO NOTHING; + END IF; - INSERT INTO task_instances ( - instance_id, - task_name, - task, - status, - started_at, - ended_at, - input, - output, - error_type, - error_title, - error_detail, - error_status, - error_instance, - last_event_time, - created_at, - updated_at - ) VALUES ( - NEW.data->>'instanceId', - NEW.data->>'taskName', - NEW.data->>'taskPosition', - NEW.data->>'status', - CASE - WHEN NEW.data ? 'startTime' - THEN to_timestamp((NEW.data->>'startTime')::numeric) - ELSE NULL - END, - CASE - WHEN NEW.data ? 'endTime' - THEN to_timestamp((NEW.data->>'endTime')::numeric) - ELSE NULL - END, - NEW.data->'input', - NEW.data->'output', - NEW.data->'error'->>'type', - NEW.data->'error'->>'title', - NEW.data->'error'->>'detail', - CASE - WHEN NEW.data->'error' ? 'status' - THEN (NEW.data->'error'->>'status')::integer - ELSE NULL - END, - NEW.data->'error'->>'instance', - event_timestamp, - NEW.time, - NEW.time - ) - ON CONFLICT (instance_id, task) DO UPDATE SET - -- Keep first non-null values (immutable fields) - task_name = COALESCE(task_instances.task_name, EXCLUDED.task_name), - - -- Update status based on event timestamp (latest wins) - status = CASE - WHEN EXCLUDED.last_event_time >= task_instances.last_event_time - THEN EXCLUDED.status - ELSE task_instances.status - END, - - -- Keep first start time, update end time with latest non-null value - started_at = COALESCE(task_instances.started_at, EXCLUDED.started_at), - ended_at = COALESCE(EXCLUDED.ended_at, task_instances.ended_at), - - -- Keep first input, update output with latest non-null value - input = COALESCE(task_instances.input, EXCLUDED.input), - output = COALESCE(EXCLUDED.output, task_instances.output), - - -- Update error fields with latest non-null values - error_type = COALESCE(EXCLUDED.error_type, task_instances.error_type), - error_title = COALESCE(EXCLUDED.error_title, task_instances.error_title), - error_detail = COALESCE(EXCLUDED.error_detail, task_instances.error_detail), - error_status = COALESCE(EXCLUDED.error_status, task_instances.error_status), - error_instance = COALESCE(EXCLUDED.error_instance, task_instances.error_instance), - - -- Track latest event time and update timestamp - last_event_time = GREATEST(EXCLUDED.last_event_time, task_instances.last_event_time), - updated_at = EXCLUDED.updated_at; - END; - - RETURN NEW; + INSERT INTO task_instances ( + instance_id, task, task_name, status, started_at, ended_at, + input, output, error_type, error_title, error_detail, error_status, error_instance, + last_event_time, from_event, created_at, updated_at + ) VALUES ( + NEW.instance_id, NEW.task, NEW.task_name, NEW.status, NEW.started_at, NEW.ended_at, + NEW.input, NEW.output, NEW.error_type, NEW.error_title, NEW.error_detail, NEW.error_status, NEW.error_instance, + NEW.last_event_time, true, now(), now() + ) + ON CONFLICT (instance_id, task) DO UPDATE SET + task_name = COALESCE(task_instances.task_name, NEW.task_name), + status = CASE + WHEN NEW.last_event_time >= task_instances.last_event_time + THEN NEW.status + ELSE task_instances.status + END, + started_at = COALESCE(task_instances.started_at, NEW.started_at), + ended_at = COALESCE(NEW.ended_at, task_instances.ended_at), + input = COALESCE(task_instances.input, NEW.input), + output = COALESCE(NEW.output, task_instances.output), + error_type = COALESCE(NEW.error_type, task_instances.error_type), + error_title = COALESCE(NEW.error_title, task_instances.error_title), + error_detail = COALESCE(NEW.error_detail, task_instances.error_detail), + error_status = COALESCE(NEW.error_status, task_instances.error_status), + error_instance = COALESCE(NEW.error_instance, task_instances.error_instance), + last_event_time = GREATEST(NEW.last_event_time, task_instances.last_event_time), + updated_at = now(); + + RETURN NULL; END; $$ LANGUAGE plpgsql; --- ============================================================================ --- TRIGGERS (Auto-normalize on INSERT) --- ============================================================================ - -CREATE TRIGGER normalize_workflow_events - BEFORE INSERT ON workflow_events_raw - FOR EACH ROW - EXECUTE FUNCTION normalize_workflow_event(); - -CREATE TRIGGER normalize_task_events - BEFORE INSERT ON task_events_raw +CREATE TRIGGER normalize_task_instances + BEFORE INSERT ON task_instances FOR EACH ROW - EXECUTE FUNCTION normalize_task_event(); + EXECUTE FUNCTION normalize_task_instance(); diff --git a/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql b/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql deleted file mode 100644 index 6f312b22e..000000000 --- a/data-index/data-index-storage/data-index-storage-migrations/src/main/resources/db/migration/V2__direct_normalized_inserts.sql +++ /dev/null @@ -1,148 +0,0 @@ --- ============================================================================ --- Direct normalized-table inserts (replaces V1's raw staging tables) --- ============================================================================ --- --- Vector now extracts and types each field in VRL (data-index/collectors/vector/ --- mode1-postgresql/vector.yaml) and inserts workflow_instances/task_instances --- directly, one named column per field - Vector's postgres sink maps JSON keys --- to columns via jsonb_populate_recordset, so no custom sink logic is needed --- beyond the row shape VRL builds. --- --- Triggers self-target these same tables via INSERT ... ON CONFLICT DO UPDATE, --- guarded by pg_trigger_depth() so the nested self-insert doesn't recurse. They --- read NEW. directly instead of parsing JSON. --- --- from_event distinguishes a Vector event (needs the merge/upsert below) from --- a direct writer like a JPA test fixture (should insert plain, untouched). --- This isn't optional: a BEFORE INSERT trigger that RETURN NULLs to cancel the --- original statement makes Postgres report 0 rows affected for THAT statement, --- even though the nested INSERT below succeeds - Hibernate's optimistic-lock --- row-count check then fails on every single em.persist(), not just conflicts. --- Entities never need to set this column; it defaults to false. --- ============================================================================ - -DROP TRIGGER IF EXISTS normalize_workflow_events ON workflow_events_raw; -DROP TRIGGER IF EXISTS normalize_task_events ON task_events_raw; -DROP FUNCTION IF EXISTS normalize_workflow_event(); -DROP FUNCTION IF EXISTS normalize_task_event(); - -ALTER TABLE workflow_instances ADD COLUMN from_event BOOLEAN NOT NULL DEFAULT false; -ALTER TABLE task_instances ADD COLUMN from_event BOOLEAN NOT NULL DEFAULT false; - -CREATE OR REPLACE FUNCTION normalize_workflow_instance() -RETURNS TRIGGER AS $$ -BEGIN - -- Nested call from our own INSERT below; let it proceed as-is. - IF pg_trigger_depth() > 1 THEN - RETURN NEW; - END IF; - - -- Direct writer (JPA, etc.): insert as given, no merge. - IF NOT NEW.from_event THEN - RETURN NEW; - END IF; - - INSERT INTO workflow_instances ( - id, namespace, name, version, status, started_at, ended_at, last_update, - input, output, error_type, error_title, error_detail, error_status, error_instance, - last_event_time, from_event, created_at, updated_at - ) VALUES ( - NEW.id, NEW.namespace, NEW.name, NEW.version, NEW.status, NEW.started_at, NEW.ended_at, NEW.last_update, - NEW.input, NEW.output, NEW.error_type, NEW.error_title, NEW.error_detail, NEW.error_status, NEW.error_instance, - NEW.last_event_time, true, now(), now() - ) - ON CONFLICT (id) DO UPDATE SET - -- Status/timestamps: newer event wins. Immutable fields: first write - -- wins. Terminal fields: latest non-null wins, never cleared. - status = CASE - WHEN NEW.last_event_time >= workflow_instances.last_event_time - THEN NEW.status - ELSE workflow_instances.status - END, - namespace = COALESCE(workflow_instances.namespace, NEW.namespace), - name = COALESCE(workflow_instances.name, NEW.name), - version = COALESCE(workflow_instances.version, NEW.version), - started_at = COALESCE(workflow_instances.started_at, NEW.started_at), - input = COALESCE(workflow_instances.input, NEW.input), - ended_at = COALESCE(NEW.ended_at, workflow_instances.ended_at), - output = COALESCE(NEW.output, workflow_instances.output), - error_type = COALESCE(NEW.error_type, workflow_instances.error_type), - error_title = COALESCE(NEW.error_title, workflow_instances.error_title), - error_detail = COALESCE(NEW.error_detail, workflow_instances.error_detail), - error_status = COALESCE(NEW.error_status, workflow_instances.error_status), - error_instance = COALESCE(NEW.error_instance, workflow_instances.error_instance), - last_update = GREATEST(NEW.last_update, workflow_instances.last_update), - last_event_time = GREATEST(NEW.last_event_time, workflow_instances.last_event_time), - -- The row may have been created first by the task backfill below - -- (from_event defaults false there); a real event landing on it now - -- means it's event-sourced from here on. - from_event = true, - updated_at = now(); - - RETURN NULL; -END; -$$ LANGUAGE plpgsql; - -CREATE TRIGGER normalize_workflow_instances - BEFORE INSERT ON workflow_instances - FOR EACH ROW - EXECUTE FUNCTION normalize_workflow_instance(); - -CREATE OR REPLACE FUNCTION normalize_task_instance() -RETURNS TRIGGER AS $$ -BEGIN - IF pg_trigger_depth() > 1 THEN - RETURN NEW; - END IF; - - -- Direct writer (JPA, etc.): insert as given, no merge. - IF NOT NEW.from_event THEN - RETURN NEW; - END IF; - - -- Backfill the parent workflow row if this task event arrived first. - IF NOT EXISTS (SELECT 1 FROM workflow_instances WHERE id = NEW.instance_id) THEN - INSERT INTO workflow_instances (id, created_at, updated_at, last_event_time) - VALUES (NEW.instance_id, now(), now(), NEW.last_event_time) - ON CONFLICT (id) DO NOTHING; - END IF; - - INSERT INTO task_instances ( - instance_id, task, task_name, status, started_at, ended_at, - input, output, error_type, error_title, error_detail, error_status, error_instance, - last_event_time, from_event, created_at, updated_at - ) VALUES ( - NEW.instance_id, NEW.task, NEW.task_name, NEW.status, NEW.started_at, NEW.ended_at, - NEW.input, NEW.output, NEW.error_type, NEW.error_title, NEW.error_detail, NEW.error_status, NEW.error_instance, - NEW.last_event_time, true, now(), now() - ) - ON CONFLICT (instance_id, task) DO UPDATE SET - task_name = COALESCE(task_instances.task_name, NEW.task_name), - status = CASE - WHEN NEW.last_event_time >= task_instances.last_event_time - THEN NEW.status - ELSE task_instances.status - END, - started_at = COALESCE(task_instances.started_at, NEW.started_at), - ended_at = COALESCE(NEW.ended_at, task_instances.ended_at), - input = COALESCE(task_instances.input, NEW.input), - output = COALESCE(NEW.output, task_instances.output), - error_type = COALESCE(NEW.error_type, task_instances.error_type), - error_title = COALESCE(NEW.error_title, task_instances.error_title), - error_detail = COALESCE(NEW.error_detail, task_instances.error_detail), - error_status = COALESCE(NEW.error_status, task_instances.error_status), - error_instance = COALESCE(NEW.error_instance, task_instances.error_instance), - last_event_time = GREATEST(NEW.last_event_time, task_instances.last_event_time), - updated_at = now(); - - RETURN NULL; -END; -$$ LANGUAGE plpgsql; - -CREATE TRIGGER normalize_task_instances - BEFORE INSERT ON task_instances - FOR EACH ROW - EXECUTE FUNCTION normalize_task_instance(); - -DROP TABLE IF EXISTS workflow_events_raw; -DROP TABLE IF EXISTS task_events_raw;