Managing conversation branches and evaluation data in Dagster
Conversation replay, evaluation, and training datasets with Go, NATS, and Dagster: an AI development data platform and four issues found in review.
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.
| Asset | Role |
|---|---|
seed_themes | Reference table of experiment themes |
conversation_sessions | Session management |
conversation_branch_nodes | Holds all nodes in the conversation tree |
taskification_results | Intermediate results from MCP and taskification |
execution_run_records | Execution facts |
response_branch_records | Final 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_versionchanges: regenerate fromtaskification_resultsonward - When
model_snapshot_idchanges: regenerate fromexecution_run_recordsonward - 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_kindvariant_confignode_statuscontext_snapshot_idtaskification replay point
taskification_results stays separate from FlowLineageEvent as an intermediate representation of prompt decomposition and orchestration plans.
Three-Tier Storage Structure
| Layer | Storage | Purpose |
|---|---|---|
| Truth source | PostgreSQL | Lineage data required for replay |
| Analytics workbench | DuckDB | Derived data for evaluation and clustering |
| Immutable artifact | Iceberg / Parquet | Frozen dataset snapshots |
The split lets lineage storage, evaluation, and dataset formats change independently.
Implementation Priority
Implementation follows five stages:
- Establish a replayable lineage model
- Introduce taskification and branch generation
- Trajectory normalization
- Embedding / clustering / labeling
- 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_runtimeand accumulated inllm_chat_turnsandflow_nodes - Lineage assets operate as downstream consumers reading from those tables
- Session-accumulated data is materialized in batch
Materialization units:
| Group | Trigger |
|---|---|
| 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:
- Retrieve node X’s conversation state (frozen message history)
- Send that message history to one or more LLMs
- 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:
| Table | Role |
|---|---|
seed_themes | Experiment theme reference table |
conversation_sessions | Session view aggregated from llm_chat_turns |
conversation_branch_nodes | Branch points in the conversation (branch_node_id, parent_node_id, turn_index, message_history_json, branch_reason) |
taskification_results | Decomposed sub-prompts |
execution_run_records | LLM execution facts (original or replay) |
response_branch_records | Comparative records for the same prompt (branching origin tracked via context_anchor_turn_id) |
lineage_watermark | Watermark 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_nodesrequires 2-pass construction for parent references
Additional fixes:
- Include topology columns (
parent_node_id,depth,branch_kind) inupsert_branch_node()conflict updates - Populate
status / error / tool_tracefrompipeline_eventfor 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.
