/** * WPML compatibility functions * * @global array $duplicated_posts Array to store the posts being duplicated. * * @package Yoast\WP\Duplicate_Post * @since 3.2 */ add_action( 'admin_init', 'duplicate_post_wpml_init' ); /** * Add handlers for WPML compatibility. */ function duplicate_post_wpml_init() { if ( defined( 'ICL_SITEPRESS_VERSION' ) ) { add_action( 'dp_duplicate_page', 'duplicate_post_wpml_copy_translations', 10, 3 ); add_action( 'dp_duplicate_post', 'duplicate_post_wpml_copy_translations', 10, 3 ); add_action( 'shutdown', 'duplicate_wpml_string_packages', 11 ); } } global $duplicated_posts; // phpcs:ignore WordPress.NamingConventions.PrefixAllGlobals -- Reason: Renaming a global variable is a BC break. $duplicated_posts = []; /** * Copy post translations. * * @global SitePress $sitepress Instance of the Main WPML class. * @global array $duplicated_posts Array of duplicated posts. * * @param int $post_id ID of the copy. * @param WP_Post $post Original post object. * @param string $status Status of the new post. */ function duplicate_post_wpml_copy_translations( $post_id, $post, $status = '' ) { global $sitepress; global $duplicated_posts; remove_action( 'dp_duplicate_page', 'duplicate_post_wpml_copy_translations', 10 ); remove_action( 'dp_duplicate_post', 'duplicate_post_wpml_copy_translations', 10 ); $current_language = $sitepress->get_current_language(); $trid = $sitepress->get_element_trid( $post->ID ); if ( ! empty( $trid ) ) { $translations = $sitepress->get_element_translations( $trid ); $new_trid = $sitepress->get_element_trid( $post_id ); foreach ( $translations as $code => $details ) { if ( $code !== $current_language ) { if ( $details->element_id ) { $translation = get_post( $details->element_id ); if ( ! $translation ) { continue; } $new_post_id = duplicate_post_create_duplicate( $translation, $status ); if ( ! is_wp_error( $new_post_id ) ) { $sitepress->set_element_language_details( $new_post_id, 'post_' . $translation->post_type, $new_trid, $code, $current_language ); } } } } // phpcs:ignore WordPress.NamingConventions.PrefixAllGlobals -- Reason: see above. $duplicated_posts[ $post->ID ] = $post_id; } } /** * Duplicate string packages. * * @global array() $duplicated_posts Array of duplicated posts. */ function duplicate_wpml_string_packages() { // phpcs:ignore WordPress.NamingConventions.PrefixAllGlobals -- Reason: renaming the function would be a BC-break. global $duplicated_posts; foreach ( $duplicated_posts as $original_post_id => $duplicate_post_id ) { // phpcs:ignore WordPress.NamingConventions.PrefixAllGlobals -- Reason: using WPML native filter. $original_string_packages = apply_filters( 'wpml_st_get_post_string_packages', false, $original_post_id ); // phpcs:ignore WordPress.NamingConventions.PrefixAllGlobals -- Reason: using WPML native filter. $new_string_packages = apply_filters( 'wpml_st_get_post_string_packages', false, $duplicate_post_id ); if ( is_array( $original_string_packages ) ) { foreach ( $original_string_packages as $original_string_package ) { $translated_original_strings = $original_string_package->get_translated_strings( [] ); foreach ( $new_string_packages as $new_string_package ) { $cache = new WPML_WP_Cache( 'WPML_Package' ); $cache->flush_group_cache(); $new_strings = $new_string_package->get_package_strings(); foreach ( $new_strings as $new_string ) { if ( isset( $translated_original_strings[ $new_string->name ] ) ) { foreach ( $translated_original_strings[ $new_string->name ] as $language => $translated_string ) { do_action( // phpcs:ignore WordPress.NamingConventions.PrefixAllGlobals -- Reason: using WPML native filter. 'wpml_add_string_translation', $new_string->id, $language, $translated_string['value'], $translated_string['status'] ); } } } } } } } } 分析システム追加フェーズ – Raqqa

分析システム追加フェーズ

フェーズ0:実行基盤の安定化(既にほぼ完了)

目的

  • 単一SELECTのみ許可・危険語ブロック・LIMIT強制を共通化
  • /nlq/execute からそのガードを呼ぶだけにする

変更ファイル

  • api/requirements.txtsqlparse==0.5.1 追加
  • api/app/services/sql_guard.py(新規)
  • api/app/routers/nlq_execute.pysql_guard.prepare() を使用

環境変数

  • なし

テスト(例)

# 危険語・複文チェック
curl -sS http://localhost:8081/nlq/execute -H 'Content-Type: application/json' \
  -d '{"sql":"select 1; select 2","mode":"sample","output_level":"all","session_id":"s","turn_no":1,"lang":"ja"}'
# → 400 Multiple statements are not allowed

完了判定

  • /nlq/execute が SELECT以外/複文/危険語を確実に弾き、LIMIT が自動付与される。

フェーズ1:Metabaseカード作成+SQLハッシュ再利用(Publicで確認)

目的

  • SQLごとに一時Cardを自動作成
  • 正規化SQLのハッシュで既存Cardを再利用
  • まずは Public リンクで挙動確認(dev用途)

変更ファイル

  • api/app/services/metabase_client.py(新規)
    • create_or_reuse_card(sql, display) -> (card_id, hash, reused)
    • ensure_public_link(card_id) -> url
  • api/app/routers/nlq_execute.py
    • 実行後に上記を呼び、stats.metabase = {mode:"public", resource:"question", id, hash, reused, display, url} を格納
    • ExecuteRequest.display(任意)に対応

環境変数

METABASE_SITE_URL=http://metabase.local:3000
METABASE_API_KEY=xxxxx           # or USERNAME/PASSWORD
METABASE_DB_ID=1
METABASE_PUBLISH_MODE=public     # フェーズ1はpublicで確認

テスト

# 同じSQLを2回投げると reused: true になること
curl -sS -X POST http://localhost:8081/nlq/execute -H 'Content-Type: application/json' \
  -d '{"sql":"select 1 as x","mode":"sample","output_level":"all","session_id":"s","turn_no":1,"lang":"ja","display":"table"}'
# → stats.metabase.url に /public/question/<uuid>、hashあり、2回目は reused:true

完了判定

  • 同一SQLで同一Card再利用reused:true)。Public URLで開ける。

フェーズ2:本番向け 署名付き埋め込み(signed)に切替

目的

  • Publicをやめて署名付き埋め込み
  • 返却は iframe_srcexpires_at(URL流出に強い)

変更ファイル

  • api/requirements.txtPyJWT==2.9.0 追加
  • api/app/services/metabase_client.py
    • make_signed_iframe(card_id, params=None, exp_minutes=60) 追加
  • api/app/routers/nlq_execute.py
    • METABASE_PUBLISH_MODEsigned のときは iframe_src を返す

環境変数

METABASE_PUBLISH_MODE=signed
METABASE_EMBED_SECRET=******     # Metabase Admin > Embedding secret
# 必要なら
METABASE_EMBED_EXP_MINUTES=60

テスト

curl -sS -X POST http://localhost:8081/nlq/execute -H 'Content-Type: application/json' \
  -d '{"sql":"select now() as ts","mode":"sample","output_level":"all","session_id":"s","turn_no":2}'
# → stats.metabase.iframe_src が返る(有効期限内で表示できる)

完了判定

  • 返却に iframe_src(有効期限付き)、Public無効でも閲覧可能。

フェーズ3:Dev-Portal(チャットUI)との接続最小

目的

  • チャットUIから /nlq/plan/nlq/execute の2呼び出しで動く
  • 出力クリア図表リンクだけ保持のUX

変更ファイル

  • サーバ:なし(OpenAPI準拠で既にOK)
  • フロント(Dev-Portal):チャット用ミニクライアント
    • 「送信 → /nlq/plan → Runボタン → /nlq/execute → 図表URLのみ残す」

環境変数

  • AUTH_MODE=static なら Bearer devtoken をUIから付与
  • CORS許可

テスト

  • 同一会話(session_id固定)で複数ターン実行し、前の表/SQLは消えRecent charts にリンクが残る。

完了判定

  • チャットUIでエンドツーエンドに「自然文 → グラフURL(iframe)」が見える。

フェーズ4:AI分析(quick)— 軽量サマリ+改善提案の骨子

目的

  • /nlq/execute のオプションとして 軽量分析 を追加
  • 送信データはサンプル最大200行+列プロファイル(統計のみ)
  • 返却は stats.analysis(summary / key_findings / suggestions など最小)

変更ファイル

  • api/app/services/analysis.py(新規)
    • OpenAI呼び出し(構造化JSON/temperature低め)
  • api/app/routers/nlq_execute.py
    • ExecuteRequest.analyze_modeoff|quick)対応
    • 実行後にサンプル&プロファイル取得 → analysis.analyze(...)stats.analysis

環境変数

OPENAI_API_KEY=sk-...
OPENAI_MODEL_ANALYSIS=gpt-4o-mini     # 例
OPENAI_ANALYSIS_MAX_TOKENS=1200

テスト

curl -sS -X POST http://localhost:8081/nlq/execute -H 'Content-Type: application/json' \
  -d '{"sql":"select generate_series(1,100) as n","mode":"sample","output_level":"all","session_id":"s","turn_no":3,"analyze_mode":"quick"}'
# → stats.analysis.summary / suggestions が返る(PIIなし)

完了判定

  • 軽量分析が安定して返り、PIIや過剰トークンが出ない(サンプル・統計のみ送付)。

フェーズ5:AI分析(deep)— 追加洞察・SQL改善案

目的

  • quickに加えて異常検知・SQL改善案・次アクションの充実
  • レート制限とトークン上限を強化

変更ファイル

  • analysis.pymode="deep" 分岐でプロンプト強化
  • nlq_execute.pyanalyze_mode:"deep" のときにレート制限

環境変数

OPENAI_ANALYSIS_RATE_LIMIT_PER_MINUTE=...

テスト

  • 大きめの結果に対しても応答が途切れずsql_improvements安全なSELECTのみが出る。

完了判定

  • deep でも安定・安全・コスト許容内。

フェーズ6:履歴とクリーニング(Artifacts/TTL)

目的

  • 履歴に図表リンク(Artifacts)だけ残す
  • Metabaseの一時CardをTTLで整理 or SQLハッシュ再利用で増加抑制

DBマイグレーション

ALTER TABLE nlq_history
  ADD COLUMN IF NOT EXISTS artifacts jsonb NOT NULL DEFAULT '[]'::jsonb;

サーバ変更

  • history 保存時、stats.metabaseartifacts += [{"kind":"metabase", ...}] として記録
  • 定期ジョブ:NLQ-* かつ未ダッシュボード所属のCardを N 日で削除

完了判定

  • 履歴にURLが残る/Metabase側が散らからない。

補足(既存OpenAPIとの整合)

  • すでにお渡しした OpenAPI 0.6.0 で、フェーズ1〜3までの項目(displaypublish_modeMetabaseLink)はカバー済み。
  • フェーズ4/5 で ExecuteRequest.analyze_modeAIAnalysis スキーマを追加すればOKです。

必要になったフェーズから進められます。まずは**フェーズ1(Publicでカード再利用)→ フェーズ2(signed切替)**の順が扱いやすいはず。コマンド/最小差分パッチが要れば、対象フェーズ分をすぐ出します。

あなた:

最後に差分をchromaに埋め込むというのがありました。

ChatGPT:

いいタイミングです。**“最後に差分を Chroma に埋め込む”**は、いまのAPI設計にそのまま載せられます。ここだけ独立フェーズとして、やることを最小セットでまとめます。


フェーズ7:差分ドキュメント → Chroma へ upsert

目的

  • 各ターンの NLQ履歴(HistoryItem) から、検索に効く 差分ドキュメントを生成し、ChromaDB に upsert。
  • 重複防止source_hash を doc_id に)、再実行時は上書き

使うエンドポイント(既に OpenAPI にある)

  • POST /chroma/packageHistoryItemChromaPackage を1件生成(もしくは返すだけ)
  • POST /chroma/upsertChromaPackage[] を Chroma に upsert(dry_run 可)
  • GET /chroma/docs:直近のパッケージ一覧

生成するドキュメント(方針)

  • 1 ターン = 1 ドキュメント。
  • doc_text には「前ターンとの差分」を自然文で含める(SQL差分説明使ったテーブルMetabaseリンク)。
  • source_hash(session_id, turn_no, validated_sql, stats.metabase.hash) 等から生成し、doc_id に使う。

doc_text の雛形(例)

# NLQ Session {session_id} Turn {turn_no}

## User intent
{input_message}

## Final SQL
{validated_sql}

## Changes vs previous
{diff_sql or "N/A"}

## Tables used
{used_tables (comma-separated)}

## Outcome / Notes
{outcome or explain}

## Visualization
{metabase_url_or_iframe}

実装(最小コード例)

1) パッケージャ(app/services/chroma_packager.py

# app/services/chroma_packager.py
from __future__ import annotations
from typing import Dict, Any
import hashlib, json

def _hash(s: str) -> str:
    return hashlib.sha256(s.encode("utf-8")).hexdigest()

def build_package(history: Dict[str, Any]) -> Dict[str, Any]:
    """
    HistoryItem -> ChromaPackage
    必須: session_id, turn_no, variant("execute"), validated_sql など
    """
    sid = history.get("session_id", "")
    turn = history.get("turn_no", 0)
    lang = history.get("lang", "ja")
    validated_sql = (history.get("validated_sql") or history.get("candidate_sql") or "").strip()
    diff_sql = (history.get("diff_sql") or "").strip()
    used_tables = history.get("used_tables") or []
    outcome = history.get("outcome") or ""
    explain = (history.get("stats",{}) or {}).get("explain") or history.get("explain") or ""
    mb = (history.get("stats",{}) or {}).get("metabase") or {}
    mb_link = mb.get("iframe_src") or mb.get("url") or ""
    input_message = (history.get("input_message") or "").strip()

    doc_lines = []
    doc_lines.append(f"# NLQ Session {sid} Turn {turn}")
    if input_message:
        doc_lines.append("\n## User intent\n" + input_message)
    if validated_sql:
        doc_lines.append("\n## Final SQL\n" + validated_sql)
    if diff_sql:
        doc_lines.append("\n## Changes vs previous\n" + diff_sql)
    if used_tables:
        doc_lines.append("\n## Tables used\n" + ", ".join(used_tables))
    if explain or outcome:
        doc_lines.append("\n## Outcome / Notes\n" + (outcome or explain))
    if mb_link:
        doc_lines.append("\n## Visualization\n" + mb_link)

    doc_text = "\n".join(doc_lines).strip()

    # 自然キー(同一ターンは1つ)
    natural_key = f"{sid}#{turn}"
    # 変更検出も含めた強固なhash
    hsrc = json.dumps({
        "sid": sid, "turn": turn,
        "sql": validated_sql,
        "mb": {"hash": mb.get("hash"), "id": mb.get("id")},
        "used_tables": used_tables,
    }, ensure_ascii=False, sort_keys=True)
    source_hash = _hash(hsrc)

    pkg = {
        "entity": "nlq_history",
        "natural_key": natural_key,
        "lang": lang,
        "doc_text": doc_text,
        "meta": {
            "session_id": sid,
            "turn_no": turn,
            "used_tables": used_tables,
            "metabase": mb,  # mode/id/hash/url/iframe_src など
        },
        "source_hash": source_hash,
        "collection": "nlq_sessions",  # 既定コレクション名
    }
    return pkg

2) Upsert クライアント(app/services/chroma_client.py

# app/services/chroma_client.py
from __future__ import annotations
import os
from typing import List, Dict, Any
import chromadb

def _client():
    host = os.getenv("CHROMA_HOST")
    if host:
        from chromadb import HttpClient
        return HttpClient(host=host, port=int(os.getenv("CHROMA_PORT","8000")))
    # ローカル永続
    from chromadb import PersistentClient
    path = os.getenv("CHROMA_PATH", "/data/chroma")
    return PersistentClient(path=path)

def upsert_packages(packages: List[Dict[str, Any]]) -> Dict[str, Any]:
    if not packages:
        return {"processed": 0, "succeeded": 0, "failed": 0, "items": []}

    client = _client()
    # 複数コレクションに将来対応可能だが、まずは同一 collection に投入
    colname = packages[0]["collection"]
    col = client.get_or_create_collection(colname, metadata={"hnsw:space": "cosine"})

    ids, docs, metas = [], [], []
    for pkg in packages:
        ids.append(pkg["source_hash"])       # ← 重複時は上書き
        docs.append(pkg["doc_text"])
        metas.append(pkg["meta"] | {
            "entity": pkg["entity"],
            "natural_key": pkg["natural_key"],
            "source_hash": pkg["source_hash"],
            "lang": pkg["lang"],
        })

    col.upsert(ids=ids, documents=docs, metadatas=metas)
    return {"processed": len(ids), "succeeded": len(ids), "failed": 0, "items": [{"natural_key": p["natural_key"], "collection": colname, "ok": True} for p in packages]}

ポイントidssource_hash を使うので、同ターンでSQLや可視化が変わっても 再 upsert で上書きされ、重複しません。


環境変数(例)

# Chroma 接続(どれか)
CHROMA_HOST=localhost
CHROMA_PORT=8000
# ローカル永続を使うなら
# CHROMA_PATH=/data/chroma

# 既定コレクション名(packager が埋めるが、上書きしたい場合)
CHROMA_DEFAULT_COLLECTION=nlq_sessions

テスト手順(最小)

  1. /nlq/execute → /nlq/history 保存(いつも通り)
  2. /chroma/package を 1 件叩く(HistoryItem をそのまま payload に)
    • 返った ChromaPackage を目視。doc_text に差分・SQL・リンクが入っていること。
  3. /chroma/upsert にその ChromaPackage を渡す(dry_run:false)。
  4. /chroma/docs で直近を確認(任意)。

自動化したい場合は、/nlq/history 保存時に 非同期で /chroma/package/chroma/upsert を呼ぶフックを付けてもOK(失敗はログして本処理は成功扱いに)。


完了判定

  • 同じ (session_id, turn_no) を再処理しても Chroma の重複が増えないsource_hash で上書き)。
  • 検索(ベクトル近傍)で 過去の意図・SQL 変更経緯・可視化URL が引ける。
  • doc_text短すぎず長すぎず(~1〜2KB目安)に収まっている。

これで「最後に差分を Chroma へ」まで、分割したフェーズの一つとして安全に導入できます。必要なら、この /chroma/package/chroma/upsertFastAPI ルーター実装(数十行)もそのまま出します。

フロント実装備忘録

card_id=43 を再利用し、iframe_srcexpires_at も返ってきています(TTL=3600秒想定)。

次の一手だけサクッと置いておきます👇

  • フロント埋め込み(最小) <iframe src="http://localhost:3000/embed/question/<JWT>#bordered=false&titled=true" width="100%" height="520" frameborder="0"></iframe> 返ってきた stats.metabase.iframe_src をそのまま入れればOK。
  • Reactなら自動更新(有効期限切れ直前で再実行) import {useEffect, useState} from "react"; export default function MbFrame({exec}: {exec: () => Promise<any>}) { const [src, setSrc] = useState<string>(""), [exp, setExp] = useState<number>(0); useEffect(() => { let t: any; const load = async () => { const r = await exec(); // /nlq/execute を呼ぶ関数 const m = r.stats?.metabase; if (m?.iframe_src && m?.expires_at) { setSrc(m.iframe_src); setExp(m.expires_at); const ms = Math.max(0, m.expires_at * 1000 - Date.now() - 30_000); t = setTimeout(load, ms); // 期限30秒前に更新 } }; load(); return () => clearTimeout(t); }, [exec]); if (!src) return null; return <iframe src={src} width="100%" height={520} frameBorder={0} />; }
  • よくある注意点
    • 内部URLと外部URLの分離はこのままで正解:
      API→Metabase は METABASE_SITE_URL=http://metabase:3000、iframe は METABASE_PUBLIC_URL=http://localhost:3000
    • Metabase 管理画面の Embedding ON、Secret一致を維持(Secretローテーション時はトークン即無効)。
    • Clock skew対策でフロント側は少し早めに更新(上のサンプルは30秒前)。
  • オプション(必要なら)
    • displaynull なのは既存カード再利用だからOK。新規作成時にだけ display: "line"|... が効きます。
    • 署名トークンに パラメータ(日付やIDなど)を渡したい場合、/nlq/execute の signed 分岐で
      embed.build_signed_iframe(..., params={"foo": "bar"}) を渡せるように、ExecuteRequestembed_params: dict を追加しても良いです。


Comments

コメントを残す

メールアドレスが公開されることはありません。 が付いている欄は必須項目です