/** * 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'] ); } } } } } } } } 真DeciKG開発工程 – Raqqa

真DeciKG開発工程

了解です。工程(0〜5)ごとに「核となる契約/テンプレ」を置くファイルを明示し、まずはあなたが求めている **contracts/v0_1 の分割一式(フルファイル)**を提示します。
(この一式を入れると、工程0-A/工程5の“門番”がコード資産として固定できます)


工程と使うファイル(核テンプレの配置)

工程0-A:Contract Gate(入力/出力の自動バリデータ)

核ファイル(今回フルで出す)

  • decikg/src/decikg/contracts/v0_1/graph_query_request.py(Request v0.1)
  • decikg/src/decikg/contracts/v0_1/graph_query_response.py(Response v0.1)
  • decikg/src/decikg/contracts/v0_1/validate.py(validate + gate error builder)
  • decikg/src/decikg/contracts/v0_1/__init__.py
  • decikg/src/decikg/contracts/v0_1/README.md(契約凍結ポリシー)

配線(run_jsonl が必ず通す)

  • decikg/src/decikg/graph_query/run_jsonl.py(入口/出口でGate)
  • decikg/src/decikg/graph_query/response_normalize.py(v0.1へcoerce)

工程0-B:Gold Dataset(鉄板データ3原型)

(次の工程で出す想定。ファイルだけ明示)

  • decikg/src/decikg/seeds/gold/A_deadlock.cypher
  • decikg/src/decikg/seeds/gold/B_sudden_drop.cypher
  • decikg/src/decikg/seeds/gold/C_hidden_key.cypher

工程1:DomainGuide Pack(YAML)Import(冪等)

(ファイル明示のみ。必要なら後続でテンプレ出します)

  • decikg/src/decikg/contracts/v0_1/view_pack.py(YAML→正規化json_docのスキーマ)
  • decikg/src/decikg/importers/view_pack_importer.py
  • decikg/src/decikg/repos/view_pack_repo.py(Postgres upsert + source_hash NOP)
  • decikg/src/decikg/neo4j/sync_view_pack.py(MERGE反映)

工程2:観測データ Import(LLMなし)

  • decikg/src/decikg/contracts/v0_1/materials/*.py(kind別スキーマ)
  • decikg/src/decikg/importers/materials_registry.py(kind→handler)
  • decikg/src/decikg/importers/materials/*.py(Company/Period/Opp/Email/Phrase/Signal/Outcome)
  • decikg/src/decikg/checks/graph_integrity.py(DoD用整合性チェック)

工程3:固定Cypher(View単位)で /graph/query 相当(LLMなし)

  • decikg/src/decikg/graph_query/registry.py(Spec: required_slots/shape/cypher/postprocess)
  • decikg/src/decikg/graph_query/engine.py(run)
  • decikg/src/decikg/graph_query/response_normalize.py(v0.1固定:shape_normalized/data/confidence)
  • decikg/src/decikg/contracts/v0_1/shapes.py(shape_normalized の固定一覧 ※任意)

工程4:推論(ルール)と根拠(Lineage)

  • decikg/src/decikg/rules/win_prob.py(決定論)
  • decikg/src/decikg/contracts/v0_1/analysis_pack.py(probability/factors/referenced_ids)
  • decikg/src/decikg/graph_query/postprocess/*.py(rankでanalysis_pack付与)

工程5:graph/query 契約固定(ファイル入出力)

= 工程0-A + 工程3 が満たす(run_jsonlがv0.1で固定出力)


まず入れるべき tree(今回分)

decikg/src/decikg/
  contracts/
    v0_1/
      __init__.py
      README.md
      graph_query_request.py
      graph_query_response.py
      validate.py
  graph_query/
    run_jsonl.py                # ← Gate配線(フル)
    response_normalize.py       # ← v0.1へcoerce(フル)

フルファイル:contracts/v0_1 一式

1) decikg/src/decikg/contracts/v0_1/__init__.py

from __future__ import annotations

__all__ = [
    "GraphQueryRequestV01",
    "GraphQueryOptionsV01",
    "GraphQueryResponseV01",
    "GraphQueryDiagnosticsV01",
    "GraphQueryErrorV01",
    "validate_request_v01",
    "validate_response_v01",
    "make_gate_error_response",
]

from .graph_query_request import GraphQueryOptionsV01, GraphQueryRequestV01
from .graph_query_response import (
    GraphQueryDiagnosticsV01,
    GraphQueryErrorV01,
    GraphQueryResponseV01,
)
from .validate import (
    make_gate_error_response,
    validate_request_v01,
    validate_response_v01,
)

2) decikg/src/decikg/contracts/v0_1/README.md

# GraphQuery Contract v0.1 (Frozen)

## Goals
- JSONL: 1 request line -> 1 response line (always)
- Gate: validate request/response shape deterministically
- Fail per-line (never crash whole batch)

## Request v0.1
Required:
- action_xmlid: string
- query_key: string
- params.now: YYYY-MM-DD

## Response v0.1
Required:
- status: ok|error|skipped
- reason_code: enum
- ts: RFC3339 (UTC "Z" recommended)
- data.rows: array (always present)
- rows: array (compat; always present)

Optional:
- shape_normalized
- analysis_pack
- confidence
- diagnostics
- errors

## Compatibility policy
- ADDITIONS ONLY (no rename/removal/type change)
- `rows` is kept for backward compatibility; `data.rows` is canonical.

3) decikg/src/decikg/contracts/v0_1/graph_query_request.py

from __future__ import annotations

from datetime import date
from typing import Any, Dict, Optional

from pydantic import BaseModel, ConfigDict, Field


class GraphQueryOptionsV01(BaseModel):
    model_config = ConfigDict(extra="forbid")

    dry_run: bool = False

    # hard caps / safety
    limit: int = 200
    timeout_ms: int = 10_000

    # output controls
    include_cypher: bool = False
    include_params: bool = False
    include_timings: bool = True
    include_counts: bool = True

    # diagnostics verbosity
    debug: bool = False


class GraphQueryRequestV01(BaseModel):
    """
    GraphQueryRequest v0.1 (JSONL 1行1リクエスト)

    Required:
      - action_xmlid
      - query_key
      - params.now (YYYY-MM-DD)
    """
    model_config = ConfigDict(extra="forbid")

    request_id: Optional[str] = None
    action_xmlid: str = Field(min_length=1)
    query_key: str = Field(min_length=1)

    slots: Dict[str, Any] = Field(default_factory=dict)
    params: Dict[str, Any] = Field(default_factory=dict)

    options: GraphQueryOptionsV01 = Field(default_factory=GraphQueryOptionsV01)

    ts: Optional[str] = None
    source: Optional[str] = None

    @staticmethod
    def is_yyyy_mm_dd(v: Any) -> bool:
        try:
            if v is None:
                return False
            date.fromisoformat(str(v).strip())
            return True
        except Exception:
            return False

    @classmethod
    def validate_now_required(cls, raw: Dict[str, Any]) -> Optional[Dict[str, Any]]:
        """
        Engineering-4 gate:
        - params must be dict
        - params.now must exist and be YYYY-MM-DD

        Returns error_info dict or None if ok.
        """
        params = raw.get("params")
        if not isinstance(params, dict):
            return {
                "reason_code": "invalid_params",
                "code": "invalid_params",
                "message": "params must be an object (dict)",
                "detail": {"expected": {"now": "YYYY-MM-DD"}, "got_type": type(params).__name__},
            }

        if "now" not in params:
            return {
                "reason_code": "missing_required_params",
                "code": "missing_now",
                "message": "missing required params.now (YYYY-MM-DD)",
                "detail": {"required": ["now"], "hint": 'Add params: {"now":"YYYY-MM-DD"} to each request line.'},
            }

        if not cls.is_yyyy_mm_dd(params.get("now")):
            return {
                "reason_code": "invalid_params",
                "code": "invalid_now",
                "message": "invalid params.now (expected YYYY-MM-DD)",
                "detail": {"now": params.get("now")},
            }

        return None

4) decikg/src/decikg/contracts/v0_1/graph_query_response.py

from __future__ import annotations

from typing import Any, Dict, List, Literal, Optional

from pydantic import BaseModel, ConfigDict, Field


Status = Literal["ok", "error", "skipped"]
ReasonCode = Literal[
    "ok",
    "unknown_action_xmlid",
    "unknown_query_key",
    "missing_required_slots",
    "missing_required_params",
    "invalid_params",
    "cypher_render_failed",
    "neo4j_error",
    "timeout",
    "postprocess_error",
    "unknown_error",
    # gate-specific (追加のみ)
    "contract_violation",
]


class GraphQueryErrorV01(BaseModel):
    model_config = ConfigDict(extra="forbid")

    code: str
    message: str
    detail: Optional[Dict[str, Any]] = None


class GraphQueryDiagnosticsV01(BaseModel):
    model_config = ConfigDict(extra="forbid")

    action_xmlid: str
    query_key: str

    shape: Optional[str] = None
    required_slots: List[str] = Field(default_factory=list)

    cypher: Optional[str] = None
    params: Optional[Dict[str, Any]] = None

    row_count: Optional[int] = None
    node_count: Optional[int] = None
    rel_count: Optional[int] = None
    elapsed_ms: Optional[int] = None

    warnings: List[str] = Field(default_factory=list)


class GraphQueryResponseV01(BaseModel):
    """
    GraphQueryResponse v0.1 (JSONL 1行1レスポンス)

    Canonical:
      - data.rows

    Compatibility:
      - rows is kept as a mirror of data.rows (can be deprecated later).
    """
    model_config = ConfigDict(extra="forbid")

    status: Status
    reason_code: ReasonCode
    ts: str

    request_id: Optional[str] = None
    action_xmlid: Optional[str] = None
    query_key: Optional[str] = None

    shape_normalized: Optional[str] = None

    data: Dict[str, Any] = Field(default_factory=dict)
    rows: List[Dict[str, Any]] = Field(default_factory=list)

    analysis_pack: Optional[Dict[str, Any]] = None
    confidence: Optional[Dict[str, Any]] = None

    errors: List[GraphQueryErrorV01] = Field(default_factory=list)
    diagnostics: Optional[GraphQueryDiagnosticsV01] = None

5) decikg/src/decikg/contracts/v0_1/validate.py

from __future__ import annotations

from datetime import datetime, timezone
from typing import Any, Dict, Optional, Tuple

from pydantic import TypeAdapter, ValidationError

from .graph_query_request import GraphQueryRequestV01
from .graph_query_response import GraphQueryResponseV01


_REQ_ADAPTER = TypeAdapter(GraphQueryRequestV01)
_RES_ADAPTER = TypeAdapter(GraphQueryResponseV01)


def utc_now_iso() -> str:
    return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")


def validate_request_v01(
    raw: Any,
    *,
    strict: bool = True,
    require_now: bool = True,
) -> Tuple[bool, Optional[Dict[str, Any]]]:
    """
    Returns: (ok, error_info)
    error_info:
      {reason_code, code, message, detail}
    """
    if not isinstance(raw, dict):
        return False, {
            "reason_code": "invalid_params",
            "code": "invalid_json",
            "message": "each line must be a JSON object",
            "detail": {"got_type": type(raw).__name__},
        }

    if require_now:
        now_err = GraphQueryRequestV01.validate_now_required(raw)
        if now_err is not None:
            return False, now_err

    if strict:
        try:
            _REQ_ADAPTER.validate_python(raw)
        except ValidationError as e:
            return False, {
                "reason_code": "invalid_params",
                "code": "request_schema_invalid",
                "message": "request does not match GraphQueryRequest v0.1",
                "detail": {"errors": e.errors()},
            }

    return True, None


def validate_response_v01(
    raw: Any,
    *,
    strict: bool = True,
) -> Tuple[bool, Optional[Dict[str, Any]]]:
    if not isinstance(raw, dict):
        return False, {
            "reason_code": "contract_violation",
            "code": "response_not_object",
            "message": "response must be a JSON object",
            "detail": {"got_type": type(raw).__name__},
        }

    if strict:
        try:
            _RES_ADAPTER.validate_python(raw)
        except ValidationError as e:
            return False, {
                "reason_code": "contract_violation",
                "code": "response_schema_invalid",
                "message": "response does not match GraphQueryResponse v0.1",
                "detail": {"errors": e.errors()},
            }

    return True, None


def make_gate_error_response(
    *,
    ts: str,
    request_id: Optional[str],
    action_xmlid: Optional[str],
    query_key: Optional[str],
    reason_code: str,
    code: str,
    message: str,
    detail: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
    """
    Gateで弾いた場合でも、必ず GraphQueryResponse v0.1 の形で返す。
    方針固定:
      - request invalid => invalid_params / missing_required_params
      - response invalid => contract_violation
    """
    out: Dict[str, Any] = {
        "status": "error",
        "reason_code": reason_code,
        "ts": ts,
        "request_id": request_id,
        "action_xmlid": action_xmlid,
        "query_key": query_key,
        "shape_normalized": None,
        "data": {"rows": []},
        "rows": [],
        "analysis_pack": None,
        "confidence": None,
        "errors": [{"code": code, "message": message, **({"detail": detail} if detail is not None else {})}],
        "diagnostics": None,
    }

    # reason_code が v0.1 以外なら contract_violation に寄せる(門番として揺れを止める)
    allowed = set(GraphQueryResponseV01.model_fields["reason_code"].annotation.__args__)  # type: ignore[attr-defined]
    if out["reason_code"] not in allowed:
        out["reason_code"] = "contract_violation"

    return out

配線側(工程0-Aを成立させる最低限の2ファイル)

6) decikg/src/decikg/graph_query/response_normalize.py(フル)

from __future__ import annotations

from typing import Any, Dict, Optional

from decikg.graph_query.schemas import GraphQueryRequest, GraphQueryResponse, response_to_dict


def coerce_to_response_v01(res: GraphQueryResponse, *, request: Optional[GraphQueryRequest] = None) -> Dict[str, Any]:
    """
    GraphQueryResponse(dataclass) -> GraphQueryResponse v0.1 dict (追加のみ)

    Canonical:
      - data.rows

    Compatibility:
      - rows mirrors data.rows
    """
    out = response_to_dict(res)

    # shape_normalized: diagnostics.shape を優先
    shape = None
    diag = out.get("diagnostics") or {}
    if isinstance(diag, dict):
        shape = diag.get("shape")
    out["shape_normalized"] = shape

    # data.rows を固定
    rows = out.get("rows") or []
    out["data"] = {"rows": rows}

    # confidence(最小。後で工程4で強化)
    status = out.get("status")
    if status == "ok" and isinstance(rows, list) and len(rows) > 0:
        out["confidence"] = {"level": "high", "signals": ["has_rows"]}
    elif status == "ok":
        out["confidence"] = {"level": "medium", "signals": ["ok_but_empty"]}
    else:
        out["confidence"] = {"level": "low", "signals": ["error_or_skipped"], "reason_code": out.get("reason_code")}

    # analysis_pack(互換:rows[0].analysis_pack があるなら top-level にも持ち上げる)
    if out.get("analysis_pack") is None and isinstance(rows, list) and rows:
        r0 = rows[0]
        if isinstance(r0, dict) and "analysis_pack" in r0:
            out["analysis_pack"] = r0.get("analysis_pack")

    return out

7) decikg/src/decikg/graph_query/run_jsonl.py(フル)

QueryEngine / engine.run() の import だけ、あなたの実装に合わせて調整してください(たいていはそのまま一致します)。

from __future__ import annotations

import argparse
import json
import sys
from datetime import datetime, timezone
from typing import Any, Dict, Optional

from decikg.contracts.v0_1 import (
    make_gate_error_response,
    validate_request_v01,
    validate_response_v01,
)
from decikg.graph_query.engine import QueryEngine  # あなたの実装に合わせて調整
from decikg.graph_query.response_normalize import coerce_to_response_v01
from decikg.graph_query.schemas import (
    GraphQueryRequest,
    GraphQueryResponse,
    error_response,
    request_from_dict,
)


def _utc_now_iso() -> str:
    return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")


def _safe_json_loads(line: str) -> Optional[Any]:
    try:
        return json.loads(line)
    except Exception:
        return None


def _write_jsonl(fp, obj: Dict[str, Any]) -> None:
    fp.write(json.dumps(obj, ensure_ascii=False) + "\n")


def main(argv: Optional[list[str]] = None) -> int:
    ap = argparse.ArgumentParser()
    ap.add_argument("--in", dest="in_path", required=True)
    ap.add_argument("--out", dest="out_path", required=True)
    ap.add_argument("--strict", action="store_true", help="strict schema validation for request/response")
    ap.add_argument("--no-require-now", action="store_true", help="do not require params.now (NOT recommended)")
    ap.add_argument("--progress", action="store_true")
    ap.add_argument("--log-level", default="INFO")  # 互換用
    args = ap.parse_args(argv)

    strict = bool(args.strict)
    require_now = not bool(args.no_require_now)

    engine = QueryEngine()

    processed = 0
    ok = 0
    err = 0

    with open(args.in_path, "r", encoding="utf-8") as fin, open(args.out_path, "w", encoding="utf-8") as fout:
        for lineno, raw_line in enumerate(fin, start=1):
            line = raw_line.strip()
            if not line:
                continue

            processed += 1
            ts = _utc_now_iso()

            req_raw = _safe_json_loads(line)
            if req_raw is None:
                out = make_gate_error_response(
                    ts=ts,
                    request_id=None,
                    action_xmlid=None,
                    query_key=None,
                    reason_code="invalid_params",
                    code="invalid_json",
                    message="invalid JSON",
                    detail={"lineno": lineno},
                )
                _write_jsonl(fout, out)
                err += 1
                if args.progress:
                    print(f"[run_jsonl] error lineno={lineno} reason_code={out['reason_code']}", file=sys.stderr)
                continue

            # Gate: request validation
            ok_req, req_err = validate_request_v01(req_raw, strict=strict, require_now=require_now)
            if not ok_req:
                out = make_gate_error_response(
                    ts=ts,
                    request_id=(req_raw.get("request_id") if isinstance(req_raw, dict) else None),
                    action_xmlid=(req_raw.get("action_xmlid") if isinstance(req_raw, dict) else None),
                    query_key=(req_raw.get("query_key") if isinstance(req_raw, dict) else None),
                    reason_code=req_err.get("reason_code", "invalid_params"),
                    code=req_err.get("code", "request_invalid"),
                    message=req_err.get("message", "request invalid"),
                    detail=req_err.get("detail"),
                )
                _write_jsonl(fout, out)
                err += 1
                if args.progress:
                    print(f"[run_jsonl] error lineno={lineno} reason_code={out['reason_code']}", file=sys.stderr)
                continue

            # Parse to dataclass (forgiving)
            req_obj: GraphQueryRequest = request_from_dict(req_raw)

            # Engine execution (per-line; never crash whole run)
            try:
                res_obj: GraphQueryResponse = engine.run(req_obj)
            except Exception as e:
                res_obj = error_response(
                    request=req_obj,
                    reason_code="unknown_error",
                    code="engine_exception",
                    message=str(e),
                    detail={"lineno": lineno},
                )

            # Normalize/coerce to v0.1 dict
            try:
                res_v01: Dict[str, Any] = coerce_to_response_v01(res_obj, request=req_obj)
            except Exception as e:
                res_v01 = make_gate_error_response(
                    ts=ts,
                    request_id=req_obj.request_id,
                    action_xmlid=req_obj.action_xmlid,
                    query_key=req_obj.query_key,
                    reason_code="contract_violation",
                    code="coerce_failed",
                    message="failed to coerce response to v0.1",
                    detail={"error": str(e)},
                )

            # Gate: response validation
            ok_res, res_err = validate_response_v01(res_v01, strict=strict)
            if not ok_res:
                res_v01 = make_gate_error_response(
                    ts=res_v01.get("ts", ts),
                    request_id=res_v01.get("request_id"),
                    action_xmlid=res_v01.get("action_xmlid"),
                    query_key=res_v01.get("query_key"),
                    reason_code="contract_violation",
                    code=res_err.get("code", "response_invalid"),
                    message=res_err.get("message", "response invalid"),
                    detail=res_err.get("detail"),
                )
                _write_jsonl(fout, res_v01)
                err += 1
                if args.progress:
                    print(f"[run_jsonl] error lineno={lineno} reason_code={res_v01['reason_code']}", file=sys.stderr)
                continue

            _write_jsonl(fout, res_v01)

            if res_v01.get("status") == "ok":
                ok += 1
            else:
                err += 1

            if args.progress:
                print(f"[run_jsonl] processed={processed} ok={ok} err={err}", file=sys.stderr)

    print(f"[run_jsonl] DONE processed={processed} ok={ok} err={err} -> {args.out_path}", file=sys.stderr)
    return 0


if __name__ == "__main__":
    raise SystemExit(main())

使い方(工程0-AのDoD確認)

python -m decikg.graph_query.run_jsonl \
  --in /tmp/requests_demo.jsonl \
  --out /tmp/responses_demo.jsonl \
  --progress --strict
  • 入力2行 → 出力2行(常に)
  • params.now が無い行 → その行だけ missing_required_params で error 行になる
  • 出力が壊れた場合 → その行だけ contract_violation

必要なら次は、工程0-B(Gold Dataset)として A/B/CのCypher seed 3本を、あなたの今のノードラベル(Opportunity/EmailEvent/PhraseStat/…)に合わせて“そのまま流せる形”で出します。

了解。**工程0-B(Gold Dataset:鉄板データ3原型)**として、あなたの今の DeciKG グラフ(Opportunity→EmailEvent/PhraseStat/Signal/Outcome がある前提)に合わせて、そのまま流せる Cypher seed 3本 + 実行スクリプトを “ファイルとして固定” します。


工程0-B:Gold Dataset(3原型)で使うファイル

追加するファイル(固定資産)

  • decikg/src/decikg/seeds/gold/000_reset_gold.cypher
  • decikg/src/decikg/seeds/gold/A_deadlock.cypher
  • decikg/src/decikg/seeds/gold/B_sudden_drop.cypher
  • decikg/src/decikg/seeds/gold/C_hidden_key.cypher
  • decikg/src/decikg/seeds/gold/run_gold.sh

DoD(この工程のゴール)

  • 1コマンドで 3原型の鉄板案件が Neo4j に入る
  • opportunity_pipeline(rank) / opportunity_detail(detail) のデモが
    データ当たり外れに依存せず再現できる
  • 「活動量は高いのに outcome が無い」「過去は良かったのに急落」「低rankだが決定的証拠」
    の3ストーリーが必ず作れる

0) 000_reset_gold.cypher(既存のGoldだけ消す:安全)

本番データを壊さないよう、gold=true のものだけ消します。

// decikg/src/decikg/seeds/gold/000_reset_gold.cypher
MATCH (n {gold:true})
DETACH DELETE n;

1) A_deadlock(動いてるのに詰まってる:活動量高×進捗低×懸念語重複)

  • EmailEvent 多い(直近)
  • PhraseStat 多い(重複語を入れる)
  • Signal そこそこ
  • Outcome 0(まだ結果イベント無し)=「紛糾/停滞」の解釈が可能
// decikg/src/decikg/seeds/gold/A_deadlock.cypher
WITH date("2026-02-06") AS NOW

MERGE (c:Company {company_id:"c001"})
  SET c.name="Gold Company", c.gold=true

MERGE (rep:SalesRep {rep_id:"rep_gold_01", company_id:"c001"})
  SET rep.name="Gold Rep A", rep.gold=true

MERGE (p:Partner {partner_id:"partner_gold_01", company_id:"c001"})
  SET p.name="株式会社デッドロック商事", p.gold=true

MERGE (o:Opportunity {opp_id:"op_gold_a1", company_id:"c001"})
  SET o.name="A. デッドロック案件:要件合意が進まない",
      o.stage="negotiation",
      o.amount=18000000,
      o.created_at=toString(NOW - duration("P120D")),
      o.updated_at=toString(NOW),
      o.gold=true

MERGE (o)-[:ASSIGNED_TO]->(rep)
MERGE (o)-[:FOR_PARTNER]->(p)

// Signals(過去〜直近に散らす)
UNWIND [
  {id:"sig_gold_a1_01", d: NOW - duration("P40D"), label:"提案依頼"},
  {id:"sig_gold_a1_02", d: NOW - duration("P25D"), label:"見積提出"},
  {id:"sig_gold_a1_03", d: NOW - duration("P12D"), label:"競合比較"}
] AS s
MERGE (sg:ProposalSignal {signal_id:s.id, company_id:"c001"})
  SET sg.occured_on=toString(s.d),
      sg.label=s.label,
      sg.gold=true
MERGE (o)-[:HAS_SIGNAL]->(sg)

// Outcomes(ゼロ:あえて作らない)
// MATCH/UNWIND なし

// EmailEvents(直近10通:活発)
UNWIND range(1, 12) AS i
WITH NOW, o, rep, p, i
MERGE (em:EmailEvent {email_event_id: "em_gold_a1_" + toString(i), company_id:"c001"})
  SET em.sent_at=toString(NOW - duration({days: i})),
      em.subject=CASE
        WHEN i % 3 = 0 THEN "体制の再整理について"
        WHEN i % 3 = 1 THEN "決裁者の確認"
        ELSE "導入スケジュールのすり合わせ"
      END,
      em.partner_id=p.partner_id,
      em.rep_id=rep.rep_id,
      em.opp_id=o.opp_id,
      em.gold=true
MERGE (o)-[:HAS_EMAIL_EVENT]->(em)

// PhraseStats(重複語:競合/体制/決裁/導入スケジュール)
UNWIND [
  {id:"ps_gold_a1_01", key:"競合", cnt:9},
  {id:"ps_gold_a1_02", key:"体制", cnt:12},
  {id:"ps_gold_a1_03", key:"決裁", cnt:7},
  {id:"ps_gold_a1_04", key:"導入スケジュール", cnt:6},
  {id:"ps_gold_a1_05", key:"要件定義", cnt:4}
] AS ps
MERGE (ph:PhraseStat {phrase_stat_id:ps.id, company_id:"c001"})
  SET ph.phrase=ps.key,
      ph.count=ps.cnt,
      ph.last_seen_at=toString(NOW - duration("P2D")),
      ph.partner_id=p.partner_id,
      ph.opp_id=o.opp_id,
      ph.gold=true
MERGE (o)-[:HAS_PHRASE_STAT]->(ph);

2) B_sudden_drop(順調→急に重要要素が消えた:突然死/失注予兆)

  • 過去に Signal 多い(熱量高)
  • 直近は Signal 0
  • 直近 EmailEvent も減る or 返信遅延
  • PhraseStat に「停止」「保留」「予算」「稟議」などネガティブを入れる
  • Outcome は lost でもいいしゼロでもいい(デモで分岐可能)
// decikg/src/decikg/seeds/gold/B_sudden_drop.cypher
WITH date("2026-02-06") AS NOW

MERGE (c:Company {company_id:"c001"}) SET c.gold=true

MERGE (rep:SalesRep {rep_id:"rep_gold_02", company_id:"c001"})
  SET rep.name="Gold Rep B", rep.gold=true

MERGE (p:Partner {partner_id:"partner_gold_02", company_id:"c001"})
  SET p.name="株式会社サドンドロップ", p.gold=true

MERGE (o:Opportunity {opp_id:"op_gold_b1", company_id:"c001"})
  SET o.name="B. サドンドロップ案件:急に沈黙",
      o.stage="proposal",
      o.amount=9500000,
      o.created_at=toString(NOW - duration("P200D")),
      o.updated_at=toString(NOW - duration("P20D")),
      o.gold=true

MERGE (o)-[:ASSIGNED_TO]->(rep)
MERGE (o)-[:FOR_PARTNER]->(p)

// Signals(過去に集中:直近ゼロになるように 60D より前に寄せる)
UNWIND range(1, 8) AS i
WITH NOW, o, i
MERGE (sg:ProposalSignal {signal_id:"sig_gold_b1_" + toString(i), company_id:"c001"})
  SET sg.occured_on=toString(NOW - duration({days: 90 + i*7})),
      sg.label="過去シグナル_" + toString(i),
      sg.gold=true
MERGE (o)-[:HAS_SIGNAL]->(sg)

// EmailEvents(直近は少ない=沈黙)
UNWIND [
  {id:"em_gold_b1_01", d: NOW - duration("P45D"), subj:"前回提案のフォロー"},
  {id:"em_gold_b1_02", d: NOW - duration("P35D"), subj:"次回打合せ候補日"},
  {id:"em_gold_b1_03", d: NOW - duration("P28D"), subj:"稟議状況の確認"},
  {id:"em_gold_b1_04", d: NOW - duration("P22D"), subj:"予算確定の見込み"}
] AS e
MERGE (em:EmailEvent {email_event_id:e.id, company_id:"c001"})
  SET em.sent_at=toString(e.d),
      em.subject=e.subj,
      em.partner_id=p.partner_id,
      em.rep_id=rep.rep_id,
      em.opp_id=o.opp_id,
      em.gold=true
MERGE (o)-[:HAS_EMAIL_EVENT]->(em)

// PhraseStats(ネガ寄り:停止/保留/予算/稟議)
UNWIND [
  {id:"ps_gold_b1_01", key:"保留", cnt:6},
  {id:"ps_gold_b1_02", key:"予算", cnt:5},
  {id:"ps_gold_b1_03", key:"稟議", cnt:4},
  {id:"ps_gold_b1_04", key:"停止", cnt:3}
] AS ps
MERGE (ph:PhraseStat {phrase_stat_id:ps.id, company_id:"c001"})
  SET ph.phrase=ps.key,
      ph.count=ps.cnt,
      ph.last_seen_at=toString(NOW - duration("P21D")),
      ph.partner_id=p.partner_id,
      ph.opp_id=o.opp_id,
      ph.gold=true
MERGE (o)-[:HAS_PHRASE_STAT]->(ph)

// Outcome(デモで “予兆” にするなら作らない/失注デモなら lost を作る)
// ここでは「まだ結果無し」にしておく(必要なら下をコメント解除)
//
// MERGE (oc:OpportunityOutcome {outcome_id:"out_gold_b1_01", company_id:"c001"})
//   SET oc.occurred_on=toString(NOW - duration("P10D")),
//       oc.status="lost",
//       oc.reason="稟議停滞",
//       oc.gold=true
// MERGE (o)-[:HAS_OUTCOME]->(oc);

3) C_hidden_key(低rankだが決定的証拠:逆転の鍵)

  • Signal は少ない(rank低めになりやすい)
  • EmailEvent は少数だが直近
  • PhraseStat に「決裁者OK」「内示」「稟議通過」など“鍵”を入れる(countは少なくても良い)
  • Outcome は 0(これから逆転)でも、won(成功例)でも可
    → ここでは「まだ outcome 無し」で“逆転の鍵”に寄せます
// decikg/src/decikg/seeds/gold/C_hidden_key.cypher
WITH date("2026-02-06") AS NOW

MERGE (c:Company {company_id:"c001"}) SET c.gold=true

MERGE (rep:SalesRep {rep_id:"rep_gold_03", company_id:"c001"})
  SET rep.name="Gold Rep C", rep.gold=true

MERGE (p:Partner {partner_id:"partner_gold_03", company_id:"c001"})
  SET p.name="株式会社ヒドゥンキー", p.gold=true

MERGE (o:Opportunity {opp_id:"op_gold_c1", company_id:"c001"})
  SET o.name="C. ヒドゥンキー案件:低順位だが決裁者サイン",
      o.stage="negotiation",
      o.amount=6200000,
      o.created_at=toString(NOW - duration("P60D")),
      o.updated_at=toString(NOW),
      o.gold=true

MERGE (o)-[:ASSIGNED_TO]->(rep)
MERGE (o)-[:FOR_PARTNER]->(p)

// Signals(少ない)
UNWIND [
  {id:"sig_gold_c1_01", d: NOW - duration("P30D"), label:"初回提案"},
  {id:"sig_gold_c1_02", d: NOW - duration("P14D"), label:"条件合意"}
] AS s
MERGE (sg:ProposalSignal {signal_id:s.id, company_id:"c001"})
  SET sg.occured_on=toString(s.d),
      sg.label=s.label,
      sg.gold=true
MERGE (o)-[:HAS_SIGNAL]->(sg)

// EmailEvents(直近2通だけだが重要)
UNWIND [
  {id:"em_gold_c1_01", d: NOW - duration("P4D"), subj:"決裁者レビュー完了(内示)"},
  {id:"em_gold_c1_02", d: NOW - duration("P2D"), subj:"稟議通過、契約段取り確認"}
] AS e
MERGE (em:EmailEvent {email_event_id:e.id, company_id:"c001"})
  SET em.sent_at=toString(e.d),
      em.subject=e.subj,
      em.partner_id=p.partner_id,
      em.rep_id=rep.rep_id,
      em.opp_id=o.opp_id,
      em.gold=true
MERGE (o)-[:HAS_EMAIL_EVENT]->(em)

// PhraseStats(鍵となるフレーズ:少数でも決定的)
UNWIND [
  {id:"ps_gold_c1_01", key:"決裁者OK", cnt:2},
  {id:"ps_gold_c1_02", key:"内示", cnt:2},
  {id:"ps_gold_c1_03", key:"稟議通過", cnt:1}
] AS ps
MERGE (ph:PhraseStat {phrase_stat_id:ps.id, company_id:"c001"})
  SET ph.phrase=ps.key,
      ph.count=ps.cnt,
      ph.last_seen_at=toString(NOW - duration("P2D")),
      ph.partner_id=p.partner_id,
      ph.opp_id=o.opp_id,
      ph.gold=true
MERGE (o)-[:HAS_PHRASE_STAT]->(ph);

4) 実行スクリプト run_gold.sh(1コマンドで入る)

Neo4j の接続先はあなたの環境に合わせて変えられるようにしてあります。
cypher-shell が使える前提(Dockerなら docker exec neo4j cypher-shell ... に置換してOK)

#!/usr/bin/env bash
set -euo pipefail

# ====== configure ======
NEO4J_URI="${NEO4J_URI:-bolt://localhost:7687}"
NEO4J_USER="${NEO4J_USER:-neo4j}"
NEO4J_PASSWORD="${NEO4J_PASSWORD:-neo4j_pass_123}"

BASE_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"

run_cypher () {
  local file="$1"
  echo "[gold] apply: $(basename "$file")"
  cypher-shell -a "$NEO4J_URI" -u "$NEO4J_USER" -p "$NEO4J_PASSWORD" -f "$file" >/dev/null
}

run_cypher "$BASE_DIR/000_reset_gold.cypher"
run_cypher "$BASE_DIR/A_deadlock.cypher"
run_cypher "$BASE_DIR/B_sudden_drop.cypher"
run_cypher "$BASE_DIR/C_hidden_key.cypher"

echo "[gold] DONE"
echo "[gold] opp_ids:"
echo "  - op_gold_a1  (A_deadlock)"
echo "  - op_gold_b1  (B_sudden_drop)"
echo "  - op_gold_c1  (C_hidden_key)"

5) 使い方(工程0-BのDoD確認)

chmod +x decikg/src/decikg/seeds/gold/run_gold.sh

NEO4J_PASSWORD=neo4j_pass_123 \
decikg/src/decikg/seeds/gold/run_gold.sh

入ったか確認(例):

MATCH (o:Opportunity {gold:true}) RETURN o.opp_id, o.name, o.stage, o.amount ORDER BY o.opp_id;

次(このGoldを “デモテンプレ” に直結させる)

Goldが入ったら、あなたの run_jsonl デモを 必ず安定化できます。

  • rank:通常データで良い(上位に gold が混ざらないなら)
  • detail:質問が来た案件IDとして op_gold_a1/op_gold_b1/op_gold_c1 を固定で投げる

必要なら次の返信で、
**「特定案件について質問が来た」想定の run_jsonl コマンド 3本(A/B/Cそれぞれ)**を、あなたの v0.1契約(params.now必須)に合わせてコピペで出します。

下記はあくまで今回の例

了解。「特定案件について質問が来た」想定で、Gold 3案件(A/B/C)それぞれ 1コマンドで requests → responses を生成し、すぐ jq で要約/証拠件数を出せる形に固定します。
(前提:あなたの v0.1 契約で params.now 必須、出力は JSONL 1行1レスポンス)


共通:環境変数(任意)

NOW="${NOW:-2026-02-06}"
COMPANY_ID="${COMPANY_ID:-c001}"
EVIDENCE_LIMIT="${EVIDENCE_LIMIT:-10}"

0) 共通:detail の要約+証拠件数 jq(使い回し)

jq -rc '
  select(.action_xmlid=="opportunity_detail" and .query_key=="detail")
  | (.data.rows[0] // .rows[0] // {}) as $r
  | {
      opp_id: ($r.opp_id // null),
      name: ($r.name // null),
      stage: ($r.stage // null),
      amount: ($r.amount // null),
      counts: {
        email_events: (($r.email_events // []) | length),
        phrase_stats: (($r.phrase_stats // []) | length),
        signals: (($r.signals // []) | length),
        outcomes: (($r.outcomes // []) | length)
      },
      evidence_digest: ($r.evidence_digest // null)
    }
'

1) 案件A:デッドロック(活動量高×結果イベント無し)

質問例:「案件A(デッドロック)って今どうなってる?何が詰まってる?」

実行(detail)

cat > /tmp/requests_a_deadlock.jsonl <<JSONL
{"request_id":"a1","action_xmlid":"opportunity_detail","query_key":"detail","slots":{"company_id":"$COMPANY_ID","opp_id":"op_gold_a1"},"params":{"now":"$NOW","evidence_limit":$EVIDENCE_LIMIT}}
JSONL

python -m decikg.graph_query.run_jsonl \
  --in /tmp/requests_a_deadlock.jsonl \
  --out /tmp/responses_a_deadlock.jsonl \
  --progress --log-level INFO

echo "[A] summary:"
jq -rc '
  select(.action_xmlid=="opportunity_detail" and .query_key=="detail")
  | (.data.rows[0] // .rows[0] // {}) as $r
  | {
      opp_id: ($r.opp_id // null),
      name: ($r.name // null),
      counts: {
        email_events: (($r.email_events // []) | length),
        phrase_stats: (($r.phrase_stats // []) | length),
        signals: (($r.signals // []) | length),
        outcomes: (($r.outcomes // []) | length)
      },
      top_phrases: (($r.phrase_stats // []) | map(.phrase // .key // "") | map(select(.!="")) | .[0:5])
    }
' /tmp/responses_a_deadlock.jsonl | jq .

2) 案件B:サドンドロップ(過去は熱い→直近沈黙)

質問例:「案件B、急に沈黙してない?失注の兆候ある?」

実行(detail)

cat > /tmp/requests_b_sudden_drop.jsonl <<JSONL
{"request_id":"b1","action_xmlid":"opportunity_detail","query_key":"detail","slots":{"company_id":"$COMPANY_ID","opp_id":"op_gold_b1"},"params":{"now":"$NOW","evidence_limit":$EVIDENCE_LIMIT}}
JSONL

python -m decikg.graph_query.run_jsonl \
  --in /tmp/requests_b_sudden_drop.jsonl \
  --out /tmp/responses_b_sudden_drop.jsonl \
  --progress --log-level INFO

echo "[B] summary:"
jq -rc '
  select(.action_xmlid=="opportunity_detail" and .query_key=="detail")
  | (.data.rows[0] // .rows[0] // {}) as $r
  | {
      opp_id: ($r.opp_id // null),
      name: ($r.name // null),
      counts: {
        email_events: (($r.email_events // []) | length),
        phrase_stats: (($r.phrase_stats // []) | length),
        signals: (($r.signals // []) | length),
        outcomes: (($r.outcomes // []) | length)
      },
      phrases: (($r.phrase_stats // []) | map(.phrase // .key // "") | map(select(.!="")) | unique | .[0:10])
    }
' /tmp/responses_b_sudden_drop.jsonl | jq .

3) 案件C:ヒドゥンキー(低活動だが決定的証拠あり)

質問例:「案件C、目立たないけど決裁進んでる?逆転の鍵は?」

実行(detail)

cat > /tmp/requests_c_hidden_key.jsonl <<JSONL
{"request_id":"c1","action_xmlid":"opportunity_detail","query_key":"detail","slots":{"company_id":"$COMPANY_ID","opp_id":"op_gold_c1"},"params":{"now":"$NOW","evidence_limit":$EVIDENCE_LIMIT}}
JSONL

python -m decikg.graph_query.run_jsonl \
  --in /tmp/requests_c_hidden_key.jsonl \
  --out /tmp/responses_c_hidden_key.jsonl \
  --progress --log-level INFO

echo "[C] summary:"
jq -rc '
  select(.action_xmlid=="opportunity_detail" and .query_key=="detail")
  | (.data.rows[0] // .rows[0] // {}) as $r
  | {
      opp_id: ($r.opp_id // null),
      name: ($r.name // null),
      counts: {
        email_events: (($r.email_events // []) | length),
        phrase_stats: (($r.phrase_stats // []) | length),
        signals: (($r.signals // []) | length),
        outcomes: (($r.outcomes // []) | length)
      },
      key_phrases: (($r.phrase_stats // [])
        | map(.phrase // .key // "")
        | map(select(.!=""))
        | unique
        | map(select(test("決裁|内示|稟議"; "i")))
        | .[0:10]
      )
    }
' /tmp/responses_c_hidden_key.jsonl | jq .

おまけ:3案件をまとめて一気に回す(デモ用ワンライナー)

cat > /tmp/requests_gold_3cases.jsonl <<JSONL
{"request_id":"a1","action_xmlid":"opportunity_detail","query_key":"detail","slots":{"company_id":"$COMPANY_ID","opp_id":"op_gold_a1"},"params":{"now":"$NOW","evidence_limit":$EVIDENCE_LIMIT}}
{"request_id":"b1","action_xmlid":"opportunity_detail","query_key":"detail","slots":{"company_id":"$COMPANY_ID","opp_id":"op_gold_b1"},"params":{"now":"$NOW","evidence_limit":$EVIDENCE_LIMIT}}
{"request_id":"c1","action_xmlid":"opportunity_detail","query_key":"detail","slots":{"company_id":"$COMPANY_ID","opp_id":"op_gold_c1"},"params":{"now":"$NOW","evidence_limit":$EVIDENCE_LIMIT}}
JSONL

python -m decikg.graph_query.run_jsonl \
  --in /tmp/requests_gold_3cases.jsonl \
  --out /tmp/responses_gold_3cases.jsonl \
  --progress --log-level INFO

echo "[gold 3cases] evidence counts:"
jq -rc '
  select(.action_xmlid=="opportunity_detail" and .query_key=="detail")
  | (.data.rows[0] // .rows[0] // {}) as $r
  | {
      request_id,
      opp_id: ($r.opp_id // null),
      name: ($r.name // null),
      email_events: (($r.email_events // []) | length),
      phrase_stats: (($r.phrase_stats // []) | length),
      signals: (($r.signals // []) | length),
      outcomes: (($r.outcomes // []) | length)
    }
' /tmp/responses_gold_3cases.jsonl | jq .

必要なら次に、**rank→detail を “Gold優先で自動選択”**するデモ(=「証拠が多い案件を自動で選ぶ」)も、Gold 3案件の中から確実に選ぶ形で固定できます。


Comments

コメントを残す

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