Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 15 additions & 15 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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: <json>} 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)

---

Expand Down Expand Up @@ -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`

Expand Down Expand Up @@ -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"**
Expand Down Expand Up @@ -1342,7 +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` - Schema with triggers
- `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
Expand Down Expand Up @@ -1437,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/.../V1__initial_schema.sql`
→ 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?"**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
127 changes: 88 additions & 39 deletions data-index/collectors/vector/mode1-postgresql/vector.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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: <json>}) ->
# PostgreSQL BEFORE INSERT triggers normalize in place
#
# Requires: Vector >= 0.46.0 (native `postgres` sink).
#
Expand Down Expand Up @@ -92,37 +93,89 @@ 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: 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: |
. = {
"tag": .flow_event.eventType,
"time": now(),
"data": .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: |
. = {
"tag": .flow_event.eventType,
"time": now(),
"data": .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
# WARNING: At 1000 events/sec, generates ~43GB/day of logs
Expand All @@ -138,20 +191,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
Expand All @@ -160,26 +209,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:
Expand Down Expand Up @@ -246,7 +295,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)
#
Expand All @@ -255,7 +304,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: <json>} rows directly into workflow_instances / task_instances
# * Exposes Prometheus metrics
#
# ============================================================================
Original file line number Diff line number Diff line change
Expand Up @@ -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: <json>} rows)
PostgreSQL Normalized Tables (BEFORE INSERT triggers normalize in place)
↓ (JPA/Hibernate)
GraphQL API
----
Expand All @@ -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
Expand Down Expand Up @@ -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:**
Expand All @@ -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:
Expand Down Expand Up @@ -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.
Expand Down
Loading
Loading