The added conversation data platform

I added three asset groups for conversation branching, evaluation, and training data generation to the working Go + NATS + Dagster event pipeline. They record LLM conversations as graphs that can be replayed from any node.

PostgreSQL stores the lineage needed for replay. DuckDB handles evaluation and clustering. Iceberg / Parquet holds frozen dataset snapshots. Each layer can change independently.

Code review found four problems in grouping, metadata, turn_index, and missing assets. The fixes and remaining work are recorded below.


Background

State of the Existing Pipeline

agent-gateway is an OpenAI-compatible Go + Gin proxy. It uses NATS JetStream and Dagster for events, PostgreSQL + pgvector for search, and DuckDB for analysis.

The existing pipeline_runtime group had 11 running assets. A NATS sensor feeds pipeline_event_inbox, then writes to llm_chat_turns, flow_nodes, and flow_edges.

What Was Missing

The pipeline normalized and stored individual requests, but lacked:

  • A mechanism to treat entire conversations as branchable graphs
  • Trajectory evaluation
  • Automated training dataset generation

I added these capabilities as a data platform for local OSS LLM fine-tuning, model evaluation, and system prompt optimization.


Designing the Three Asset Groups

Conversation Lineage Assets

This group records the conversation tree.

AssetRole
seed_themesReference table of experiment themes
conversation_sessionsSession management
conversation_branch_nodesHolds all nodes in the conversation tree
taskification_resultsIntermediate results from MCP and taskification
execution_run_recordsExecution facts
response_branch_recordsFinal responses and traces

Each node stores branch_kind (root, retry, tool_variant, sampling_variant, prompt_variant, taskify_variant) and variant_config as JSON.

The dependency graph is a serial chain:

  seed_themes
  → conversation_sessions
    → conversation_branch_nodes
      → taskification_results
        → execution_run_records
          → response_branch_records
  

Replay Rules

  • When taskify_rule_version changes: regenerate from taskification_results onward
  • When model_snapshot_id changes: regenerate from execution_run_records onward
  • When a specific node’s conditions change: regenerate only that subtree

Evaluation and Distribution Labeling Assets

This group converts conversation trajectories into evaluation units.

  trajectory_units          # Normalize paths from root to terminal node
  → trajectory_embeddings # Vectorize
    → trajectory_clusters # Cluster
      → cluster_category_labels  # Assign category labels
        → human_sample_queue     # Human review queue
          → human_eval_labels    # Accumulate evaluation results
            → candidate_preference_pairs
              → preference_labels  # Build DPO pairs
  

Categories include tool_helped, tool_misuse, hallucination, good_recovery, overlong, underspecified, and stable_high_quality.

Design Principles

  • Embeddings are representations, not labels
  • Clusters are structure, not judgments
  • Model-based categorization and human preference are explicitly separated into distinct layers

Training Dataset Generation Assets

This group builds training datasets from evaluated trajectories.

  positive_topk_trajectories     # Stratified top-K extraction
negative_bottomk_trajectories  # Stratified bottom-K extraction
  → sft_messages_dataset       # SFT dataset
  → dpo_preference_dataset     # DPO dataset
  → suppression_dataset        # Suppression dataset
  → unlearning_dataset         # Unlearning dataset
    → dataset_snapshot_registry  # Metadata for frozen snapshots
  

Negative types are hard_negative, tool_misuse, hallucination, noise_reject, and suppression_candidate. They stay separate rather than all being used for unlearning.


Key Design Decisions

Active Orchestrator vs. Passive Capture

I first decided whether seed_themes → conversation_sessions → branch_nodes would drive conversations or organize existing ones.

The first version starts with manually chosen themes and session IDs. After repeated interactions, it materializes conversations by session ID and supports replay from any node. An active orchestrator can be added later.

Where to Place Branching Logic

I kept session / session_turn / session_lineage as execution logs and built experimental assets above them.

The parent-child turn structure adds:

  • branch_kind
  • variant_config
  • node_status
  • context_snapshot_id
  • taskification replay point

taskification_results stays separate from FlowLineageEvent as an intermediate representation of prompt decomposition and orchestration plans.

Three-Tier Storage Structure

LayerStoragePurpose
Truth sourcePostgreSQLLineage data required for replay
Analytics workbenchDuckDBDerived data for evaluation and clustering
Immutable artifactIceberg / ParquetFrozen dataset snapshots

The split lets lineage storage, evaluation, and dataset formats change independently.

Implementation Priority

Implementation follows five stages:

  1. Establish a replayable lineage model
  2. Introduce taskification and branch generation
  3. Trajectory normalization
  4. Embedding / clustering / labeling
  5. Dataset generation and snapshot freeze

Replay with local OSS LLMs and branch-node normalization come first. Evaluation and datasets depend on that foundation.


Implementation Plan

Directory Structure

New modules go under model-foundry/src/loftllc/, alongside pipeline_runtime, following domain / infra / presentation layering:

  loftllc/
  domain/
    conversation_lineage.py
  infra/
    lineage_store.py
  presentation/
    dagster/
      assets/
        conversation_lineage.py
  

Difference in Asset Execution Patterns

pipeline_runtime creates one RunRequest per NATS message. Lineage assets batch the accumulated session data:

  • NATS events are already processed by pipeline_runtime and accumulated in llm_chat_turns and flow_nodes
  • Lineage assets operate as downstream consumers reading from those tables
  • Session-accumulated data is materialized in batch

Materialization units:

GroupTrigger
Group 1 (Conversation Lineage)Manual or 60-second interval sensor
Group 2 (Evaluation)Auto-triggered via dependency on Group 1
Group 3 (Dataset Generation)Fully manual

Partition Strategy

DynamicPartitionsDefinition handles session IDs created at runtime by the Go gateway. A sensor discovers IDs in llm_chat_turns, adds dynamic partitions, and issues a RunRequest. Group 3 uses separate snapshot-ID partitions.

Replay Mechanism

Replay from node X means:

  1. Retrieve node X’s conversation state (frozen message history)
  2. Send that message history to one or more LLMs
  3. Record new responses as response_branch_records

The new LLMClientResource sends HTTP requests to /v1/chat/completions. Replay uses the same gateway routing, logs, and NATS events as ordinary requests.

PostgreSQL Schema Design

The lineage tables are:

TableRole
seed_themesExperiment theme reference table
conversation_sessionsSession view aggregated from llm_chat_turns
conversation_branch_nodesBranch points in the conversation (branch_node_id, parent_node_id, turn_index, message_history_json, branch_reason)
taskification_resultsDecomposed sub-prompts
execution_run_recordsLLM execution facts (original or replay)
response_branch_recordsComparative records for the same prompt (branching origin tracked via context_anchor_turn_id)
lineage_watermarkWatermark for incremental processing

Design Problems Found in Code Review

Review after Phase 1 found four problems.

Problem 1: Grouping Breakdown in response_branch_records

Each replay received a new turn_id, then grouping used that ID. Replays therefore became separate groups and branch_index stayed effectively at 0. This prevented shared-parent comparisons and preference-pair construction.

Fix: group by context_anchor_turn_id.

Problem 2: Missing Metadata in original turn execution_run_records

chat_history supplied only role / content / created_at, but the builder expected model / backend / assistant_message / token_usage / duration_ms. Original runs ended up with empty response_text and usage.

The metadata needed to come from llm_chat_turns in DuckDB or pipeline_event_inbox in PostgreSQL.

Problem 3: turn_index Collision in Replay Branches

The fixed turn_index = parent_turn_index + 1 collided with the existing next turn. Reads used only ORDER BY turn_index ASC, leaving sibling branches and serial turns in unstable order. That broke replay from arbitrary nodes.

Fix: replace the serial turn_index model with a parent-child branch-node tree.

Problem 4: Unimplemented Core Assets

conversation_sessions and conversation_branch_nodes had no materialization. The required depth / branch_kind / context_snapshot_id / variant_config / node_status fields were therefore missing.

Remaining Issues After Initial Fix

Two issues remained after the first fixes:

  • Whether to exclude replay nodes during normal materialization, or to skip turns with no original source in build_execution_record_from_event()
  • conversation_branch_nodes requires 2-pass construction for parent references

Additional fixes:

  • Include topology columns (parent_node_id, depth, branch_kind) in upsert_branch_node() conflict updates
  • Populate status / error / tool_trace from pipeline_event for original runs as well

From replay to training data

Five experimental asset tables sit above session / session_turn / session_lineage. The goal is controlled replay from any context anchor, response-distribution evaluation, and training dataset generation.


Results

  • Designed three asset groups—Conversation Lineage, Evaluation & Distribution Labeling, Training Dataset Generation—adding approximately 30 assets in total
  • Finalized a three-tier storage strategy: PostgreSQL / DuckDB / Iceberg
  • Completed Phase 1 (conversation lineage foundation) implementation
  • Discovered and fixed four critical design problems through code review
  • Clarified remaining issues: branch node 2-pass construction, context_anchor-based grouping, replay node exclusion logic
  • Established implementation priority as five stages: replay foundation → taskification → trajectory normalization → evaluation → dataset generation

Subsequent Pivot

Operation led to the following changes.

Adding Telemetry and Pivoting to a Visualization Pipeline

The Dagster asset graph was hard to read in real time. I added Vector → Prometheus → Grafana telemetry to show conversation state, branches, and evaluation results in dashboards.

Flattening Dagster Assets

With visualization in telemetry, I moved to flat Dagster assets rather than a deep dependency graph representing the conversation tree.

Redesigning agent-gateway as an Orchestrator Model

I also changed agent-gateway from one gateway processing every request to an orchestrator model controlling worker LLMs. Their interaction flows become assets and are shown in Grafana in real time.

A separate article is planned for this redesign.