了解です。工程(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__.pydecikg/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.cypherdecikg/src/decikg/seeds/gold/B_sudden_drop.cypherdecikg/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.pydecikg/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.cypherdecikg/src/decikg/seeds/gold/A_deadlock.cypherdecikg/src/decikg/seeds/gold/B_sudden_drop.cypherdecikg/src/decikg/seeds/gold/C_hidden_key.cypherdecikg/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案件の中から確実に選ぶ形で固定できます。
コメントを残す