追加した会話データ基盤

稼働中の Go + NATS + Dagster イベントパイプラインに、会話の分岐、評価、学習データ生成を扱う3つのアセットグループを設計・実装しました。LLM 会話を、任意の地点から再実行できるグラフとして記録するためです。

PostgreSQL は再実行に必要なリネージュ、DuckDB は評価・クラスタリング、Iceberg / Parquet は凍結済みデータセットを保存します。用途を分け、各層を独立に変更できる構成にしました。

実装レビューでは、グルーピング、メタデータ、turn_index、未実装の asset に4つの問題が見つかりました。修正方針と残課題も記録します。


前提

既存パイプラインの状態

agent-gateway は Go + Gin の OpenAI 互換プロキシです。NATS JetStream、Dagster、PostgreSQL + pgvector、DuckDB を使い、イベント処理、検索、分析を行います。

既存の pipeline_runtime では11個の Dagster アセットが稼働していました。NATS センサーが拾ったイベントを pipeline_event_inbox から llm_chat_turns、flow_nodes、flow_edges へ書き込みます。

不足していた機能

既存処理はリクエスト単位の正規化と蓄積までで、次の機能がありませんでした。

  • 会話全体を 分岐可能なグラフ として扱う仕組み
  • トラジェクトリの評価
  • トレーニング用データセットの自動生成

ローカル OSS LLM の fine-tuning、モデル評価、system prompt の最適化に使うデータ基盤として、これらを Dagster に追加します。


3つのアセットグループの設計

Conversation Lineage Assets

会話ツリーを記録するグループです。

アセット役割
seed_themes実験テーマの参照テーブル
conversation_sessionsセッションの管理
conversation_branch_nodes会話ツリーの全ノードを保持
taskification_resultsMCP やタスク化の中間結果
execution_run_records実行ファクト
response_branch_records最終レスポンスとトレース

各ノードは branch_kind(root・retry・tool_variant・sampling_variant・prompt_variant・taskify_variant)と variant_config を JSON で保存します。

依存グラフは次の直列構造です。

  seed_themes
  → conversation_sessions
    → conversation_branch_nodes
      → taskification_results
        → execution_run_records
          → response_branch_records
  

リプレイルール

  • taskify_rule_version が変わった場合は taskification_results 以降を再生成
  • model_snapshot_id が変わった場合は execution_run_records 以降を再生成
  • 特定ノードの条件変更はそのサブツリーのみ再生成

Evaluation and Distribution Labeling Assets

会話の処理経路を評価単位へ変換するグループです。

  trajectory_units          # ルートからターミナルノードまでのパスを正規化
  → trajectory_embeddings # ベクトル化
    → trajectory_clusters # クラスタリング
      → cluster_category_labels  # カテゴリラベル付与
        → human_sample_queue     # 人間レビューキュー
          → human_eval_labels    # 評価結果の蓄積
            → candidate_preference_pairs
              → preference_labels  # DPO ペア構築
  

カテゴリは tool_helped、tool_misuse、hallucination、good_recovery、overlong、underspecified、stable_high_quality などです。

設計原則

  • embedding はラベルではなく表現
  • クラスタは判断ではなく構造
  • モデルによるカテゴリ化と人間の選好は別レイヤーとして明確に分離する

Training Dataset Generation Assets

評価済みの処理経路から学習データを作るグループです。

  positive_topk_trajectories     # 層別 top-K 抽出
negative_bottomk_trajectories  # 層別 bottom-K 抽出
  → sft_messages_dataset       # SFT 用データセット
  → dpo_preference_dataset     # DPO 用データセット
  → suppression_dataset        # 抑制用データセット
  → unlearning_dataset         # アンラーニング用データセット
    → dataset_snapshot_registry  # 凍結済みスナップショットのメタデータ管理
  

negative type は hard_negative・tool_misuse・hallucination・noise_reject・suppression_candidate です。すべてを一律にアンラーニングへ使わず、種類ごとに分けます。


設計上の主要な意思決定

能動的オーケストレーターか受動的キャプチャか

seed_themes → conversation_sessions → branch_nodes が会話を駆動するのか、既存の会話を整理するのかを検討しました。

まず手動でテーマを決め、セッション ID を発行して対話を繰り返します。その ID で会話をマテリアライズし、任意のノードから再実行できる土台を作ります。会話を自動で進めるオーケストレーターは、その上に後から載せる設計です。

分岐処理の位置

既存の session / session_turn / session_lineage は実行ログとして残し、その上に実験用の asset を作ります。

親子関係を持つターンに、次の情報を追加します。

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

taskification_results は FlowLineageEvent と分け、prompt 分解と orchestration plan の中間表現として保存します。

ストレージ三層構造

レイヤーストレージ用途
Truth sourcePostgreSQLreplay 必須のリネージュデータ
分析ワークベンチDuckDB評価・クラスタリング向け派生データ
Immutable artifactIceberg / Parquet凍結済みデータセットスナップショット

リネージュの保存方式、評価処理、データセット形式をそれぞれ変更できるように分けています。

実装優先順位

実装は次の5段階で進めます。

  1. replay 可能なリネージュモデルの確立
  2. taskification と branch generation の導入
  3. trajectory 化
  4. embedding / clustering / labeling
  5. dataset generation と snapshot freeze

最初に OSS ローカル LLM の replay と branch node の正規化を固めます。その結果を評価とデータセット生成で使います。


実装計画

ディレクトリ構造

model-foundry/src/loftllc/ に pipeline_runtime と並列のモジュールを追加します。既存の domain / infra / presentation の構成に合わせます。

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

アセット実行パターンの違い

pipeline_runtime は NATS メッセージごとに1つの RunRequest を作ります。リネージュアセットは蓄積済みのデータをまとめて処理します。

  • NATS イベントはすでに pipeline_runtime で処理され llm_chat_turns や flow_nodes に蓄積済み
  • リネージュアセットはそこから読み取る 下流コンシューマ として動作
  • セッションの蓄積データをバッチ的にマテリアライズする方式

マテリアライズの単位は次のとおりです。

グループトリガー
Group 1(Conversation Lineage)手動または 60 秒間隔センサー
Group 2(Evaluation)Group 1 への依存による自動トリガー
Group 3(Dataset Generation)完全手動トリガー

パーティション戦略

DynamicPartitionsDefinition を使います。セッション ID は Go ゲートウェイが実行時に作るため、センサーが llm_chat_turns から発見し、動的パーティションへ追加して RunRequest を発行します。Group 3 はスナップショット ID を別の動的パーティションにします。

リプレイメカニズム

ノード X からの replay は次の処理です。

  1. ノード X の会話状態(凍結済みメッセージ履歴)を取得
  2. 1 つ以上の LLM にそのメッセージ履歴を送信
  3. 新しいレスポンスを response_branch_records として記録

infra リソース LLMClientResource から /v1/chat/completions へ HTTP リクエストを送ります。replay もゲートウェイを通り、通常と同じルーティング、ログ、NATS イベントを記録します。

PostgreSQL スキーマ設計

リネージュ用テーブルの定義です。

テーブル役割
seed_themes実験テーマの参照テーブル
conversation_sessionsllm_chat_turns から集約されたセッションビュー
conversation_branch_nodes会話内の分岐点(branch_node_id・parent_node_id・turn_index・message_history_json・branch_reason)
taskification_results分解されたサブプロンプト
execution_run_recordsLLM 実行ファクト(original または replay)
response_branch_records同一プロンプトに対する比較レコード(context_anchor_turn_id で分岐元を追跡)
lineage_watermark増分処理用のウォーターマーク

実装レビューで発見した設計問題

Phase 1 の会話リネージュ基盤を実装後、レビューで4つの問題が見つかりました。

問題1: response_branch_records のグルーピング破綻

replay ごとに新しい turn_id を付け、その ID でグループ分けしていました。このため各 replay が別グループになり、branch_index は実質0のままでした。同じ親から条件を変えた分岐を比較できず、応答の比較や preference pair の元データになりません。

修正方針:context_anchor_turn_id でグループ分けします。

問題2: original turn の execution_run_records メタデータ欠落

chat_history から取得するのは role / content / created_at だけでしたが、ビルダーは model / backend / assistant_message / token_usage / duration_ms も読む前提でした。元の実行の response_text と usage が空になります。

メタデータは llm_chat_turns(DuckDB)または pipeline_event_inbox(PostgreSQL)から取得する必要がありました。

問題3: replay branch の turn_index 衝突

replay branch に固定で turn_index = parent_turn_index + 1 を入れ、既存の次ターンと衝突していました。読み出しも ORDER BY turn_index ASC だけなので、兄弟分岐と直列ターンの順序が不安定です。任意ノードを replay root にする前提が崩れます。

修正方針:turn_index による直列モデルから、branch node table の親子ツリーへ変えます。

問題4: 仕様の中核 asset の未実装

conversation_sessions と conversation_branch_nodes のマテリアライズがなく、depth / branch_kind / context_snapshot_id / variant_config / node_status が保存されていませんでした。分岐グラフに必要な情報が不足していました。

修正後の残課題

1回目の修正後も、次の2点が残りました。

  • 通常マテリアライズで replay node を除外するか、build_execution_record_from_event() 側で original source のない turn を skip するかの選択
  • conversation_branch_nodes の親参照 2-pass 構築が必要

追加の修正方針です。

  • upsert_branch_node() の conflict update に topology 系カラム(parent_node_id・depth・branch_kind)も含める
  • original run でも status / error / tool_trace を pipeline_event から埋める

再実行から学習データへつなぐ

既存の session / session_turn / session_lineage の上に、実験用の5つのテーブルを作ります。任意ノードを context anchor にして条件を制御した replay を行い、応答分布を評価して学習用 dataset へつなぐ設計です。


結果

  • Conversation Lineage・Evaluation & Distribution Labeling・Training Dataset Generation の3アセットグループを設計し、合計約30個のアセットを追加
  • PostgreSQL / DuckDB / Iceberg の三層ストレージ戦略を確定
  • Phase 1(会話リネージュ基盤)の実装を完了
  • コードレビューで4つの重大な設計問題を発見・修正
  • 残課題(branch node の 2-pass 構築・context_anchor ベースのグルーピング・replay node の除外ロジック)を明確化
  • 実装優先順位を「replay 基盤 → taskification → trajectory 化 → evaluation → dataset generation」の5段階に確定

その後の方針転換

運用中に、次の方針へ変更しました。

テレメトリの追加と可視化パイプラインへの転換

Dagster のアセットグラフだけでは、会話の状態をリアルタイムに把握しにくい問題がありました。Vector → Prometheus → Grafana のテレメトリ経路を追加し、会話の状態、分岐、評価結果をダッシュボードで表示するようにしました。

Dagster アセットのフラット化

可視化をテレメトリへ移し、Dagster のアセットはフラットにする方針へ変えました。会話ツリーを深い依存グラフで表す必要がなくなったためです。

agent-gateway のオーケストレーターモデルへの再設計

agent-gateway も、単一ゲートウェイで全リクエストを処理する構成から、オーケストレーターモデルがワーカー LLM 群を操作する構成へ変えました。LLM 同士の対話フローをアセットにし、Grafana でリアルタイムに表示します。

この再設計は別の記事で記録する予定です。