| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106 |
- import json
- import re
- from decimal import Decimal
- from pathlib import Path
- from content_agent.integrations.database_runtime import (
- ContentSupplyDbConfig,
- DatabaseRuntimeStore,
- )
- class FakeCursor:
- def __init__(self, connection):
- self.connection = connection
- self._one = None
- self._all = []
- def __enter__(self):
- return self
- def __exit__(self, *_args):
- return None
- def execute(self, sql, params=None):
- self.connection.statements.append((sql, list(params or [])))
- if sql.startswith("SELECT COUNT(*)"):
- self._one = {"cnt": 0}
- self._all = []
- elif sql.startswith("SELECT"):
- self._one = self.connection.select_one_result
- self._all = self.connection.select_all_result
- def fetchone(self):
- return self._one
- def fetchall(self):
- return self._all
- class FakeConnection:
- def __init__(self):
- self.statements = []
- self.commit_count = 0
- self.select_one_result = None
- self.select_all_result = []
- def __enter__(self):
- return self
- def __exit__(self, *_args):
- return None
- def cursor(self):
- return FakeCursor(self)
- def commit(self):
- self.commit_count += 1
- def test_content_supply_db_config_reads_env_file(tmp_path):
- env_file = tmp_path / ".env"
- env_file.write_text(
- "\n".join(
- [
- "CONTENT_SUPPLY_DB_HOST=127.0.0.1",
- "CONTENT_SUPPLY_DB_PORT=3307",
- "CONTENT_SUPPLY_DB_NAME=content-deconstruction-supply",
- "CONTENT_SUPPLY_DB_USER=content_rw",
- "CONTENT_SUPPLY_DB_" + "PASS" + "WORD=dummy_password",
- ]
- ),
- encoding="utf-8",
- )
- config = ContentSupplyDbConfig.from_env(env_file=env_file)
- assert config.host == "127.0.0.1"
- assert config.port == 3307
- assert config.database == "content-deconstruction-supply"
- assert config.user == "content_rw"
- def test_content_supply_db_config_requires_all_project_db_keys(tmp_path):
- env_file = tmp_path / ".env"
- env_file.write_text(
- "\n".join(
- [
- "CONTENT_SUPPLY_DB_HOST=127.0.0.1",
- "CONTENT_SUPPLY_DB_PORT=3307",
- "CONTENT_SUPPLY_DB_NAME=content-deconstruction-supply",
- "CONTENT_SUPPLY_DB_USER=content_rw",
- ]
- ),
- encoding="utf-8",
- )
- try:
- ContentSupplyDbConfig.from_env(env_file=env_file)
- except ValueError as exc:
- assert "CONTENT_SUPPLY_DB_PASSWORD" in str(exc)
- else:
- raise AssertionError("expected missing db env key to fail")
- def test_database_runtime_writes_source_context_with_db_schema_version():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.write_json(
- "run_001",
- "source_context.json",
- {
- "schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "demand_content_id": "123",
- "ext_data": {
- "evidence_pack": {
- "pattern_source_system": "pg_pattern_v2",
- "source_kind": "pattern_itemset",
- "source_post_id": "post_001",
- "pattern_execution_id": 581,
- "mining_config_id": 2081,
- }
- },
- },
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_source_contexts`" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert values["run_id"] == "run_001"
- assert values["demand_content_id"] == 123
- assert json.loads(values["evidence_pack"])["source_post_id"] == "post_001"
- assert json.loads(values["source_context"])["schema_version"] == "runtime_record.v1"
- def test_database_runtime_derives_itemset_ids_from_seed_pack_itemsets():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.write_json(
- "run_001",
- "pattern_seed_pack.json",
- {
- "schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "source_post_id": "post_001",
- "pattern_execution_id": 581,
- "itemsets": [{"itemset_id": 1608352}, {"itemset_id": 1608352}],
- "seed_terms": ["爱国情感", "人物故事"],
- },
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_pattern_seed_packs`" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert json.loads(values["itemset_ids"]) == [1608352]
- assert json.loads(values["pattern_seed_pack"])["schema_version"] == "runtime_record.v1"
- def test_database_runtime_appends_jsonl_with_raw_payload():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.append_jsonl(
- "run_001",
- "search_queries.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "search_query_id": "q_001",
- "search_query": "对比分析",
- "search_query_generation_method": "item_single",
- "pattern_seed_ref": {
- "source_field": "seed_terms",
- "source_index": 0,
- "seed_term": "对比分析",
- },
- "raw_payload": {"run_id": "run_001", "search_query_id": "q_001"},
- }
- ],
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_queries`" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert json.loads(values["pattern_seed_ref"])["seed_term"] == "对比分析"
- assert json.loads(values["raw_payload"])["search_query_id"] == "q_001"
- def test_database_runtime_upserts_content_media_records():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.append_jsonl(
- "run_001",
- "content_media_records.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "platform": "douyin",
- "platform_content_id": "content_001",
- "content_media_status": "oss_uploaded",
- "content_metadata_source": "douyin_keyword_search",
- "play_url": "https://source.example/video.mp4",
- "local_path": None,
- "oss_url": "https://res.example/video.mp4",
- "raw_payload": {
- "run_id": "run_001",
- "platform_content_id": "content_001",
- "oss_object_key": "crawler/video/content_001.mp4",
- },
- "created_at": "2026-06-15T00:00:00+00:00",
- }
- ],
- )
- sql, _ = connection.statements[-1]
- assert "INSERT INTO `content_agent_content_media_records`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert "`oss_url` = VALUES(`oss_url`)" in sql
- def test_database_runtime_upserts_failed_search_query_status():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.append_jsonl(
- "run_001",
- "search_queries.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "search_query_id": "q_001",
- "search_query": "接口失败",
- "search_query_generation_method": "item_single",
- "search_query_effect_status": "failed",
- "raw_payload": {
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "search_query_id": "q_001",
- "search_query_effect_status": "failed",
- "query_failure": {
- "status": "failed",
- "error_code": "PLATFORM_REQUEST_FAILED",
- },
- },
- }
- ],
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_queries`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["search_query_effect_status"] == "failed"
- payload = json.loads(values["raw_payload"])
- assert payload["query_failure"]["error_code"] == "PLATFORM_REQUEST_FAILED"
- def test_database_runtime_preserves_llm_variant_payload_fields():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.append_jsonl(
- "run_001",
- "search_queries.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "search_query_id": "q_002",
- "search_query": "人物叙事素材",
- "search_query_generation_method": "llm_variant",
- "pattern_seed_ref": {
- "source_field": "seed_terms",
- "source_index": 0,
- "seed_term": "人物故事",
- },
- "llm_variant_of": "q_001",
- "raw_payload": {
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "search_query_id": "q_002",
- "search_query_generation_method": "llm_variant",
- "llm_variant_of": "q_001",
- "llm_input_evidence": {"seed_term": "人物故事"},
- "llm_prompt_version": "fake-query-prompt-v1",
- "llm_generation_model": "fake-query-model",
- },
- }
- ],
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_queries`" in sql
- assert "llm_variant_of" not in values
- payload = json.loads(values["raw_payload"])
- assert payload["llm_variant_of"] == "q_001"
- assert payload["llm_input_evidence"]["seed_term"] == "人物故事"
- assert payload["llm_prompt_version"] == "fake-query-prompt-v1"
- assert payload["llm_generation_model"] == "fake-query-model"
- def test_database_runtime_upserts_pattern_recall_evidence():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- # V3 形状: Gemini 直读判定证据(recall_status=judged + evidence_summary 判定字段);
- # decode 时代 7 列已随 V2 数据归档删除,列过滤必须把混入的旧键挡在 DB 外。
- store.append_jsonl(
- "run_001",
- "pattern_recall_evidence.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "recall_evidence_id": "recall_001",
- "content_discovery_id": "content_001",
- "platform": "douyin",
- "platform_content_id": "7390000000000000000",
- "recall_status": "judged",
- "decode_status": "success",
- "matched_terms": ["爱国情感"],
- "evidence_summary": {
- "judge_status": "ok",
- "fit_senior_50plus": True,
- "relevance_score": 0.85,
- },
- "raw_payload": {
- "platform": "douyin",
- "judge_reason": "贴题且适合 50+",
- },
- }
- ],
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_pattern_recall_evidence`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert values["recall_evidence_id"] == "recall_001"
- assert values["recall_status"] == "judged"
- assert "platform" not in values
- assert "decode_status" not in values
- assert "matched_terms" not in values
- assert json.loads(values["evidence_summary"])["fit_senior_50plus"] is True
- assert json.loads(values["raw_payload"])["platform"] == "douyin"
- def test_database_runtime_preserves_p5_rule_decision_fields():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.append_jsonl(
- "run_001",
- "rule_decisions.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "decision_id": "d_001",
- "policy_bundle_id": "policy_bundle_v1",
- "rule_pack_id": "douyin_content_discovery_rule_pack_v1",
- "rule_pack_version": "1.0.0",
- "strategy_version": "V1",
- "decision_target_type": "content",
- "decision_target_id": "content_001",
- "decision_action": "REJECT_CONTENT",
- "decision_reason_code": "pattern_recall_failed",
- "search_query_effect_status": "rule_blocked",
- "score": None,
- "triggered_blocking_rules": ["gate_pattern_recall_failed"],
- "scorecard": {"total_score": None, "score_missing": True},
- "decision_replay_data": {
- "policy_bundle_hash": "hash_001",
- "dispatch_id": "dispatch_content",
- "effect_mapping_id": "map_hard_gate_reject_rule_blocked",
- },
- "raw_payload": {
- "decision_id": "d_001",
- "search_query_effect_status": "rule_blocked",
- "decision_replay_data": {
- "policy_bundle_hash": "hash_001",
- "dispatch_id": "dispatch_content",
- },
- },
- }
- ],
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_rule_decisions`" in sql
- assert values["search_query_effect_status"] == "rule_blocked"
- assert json.loads(values["triggered_blocking_rules"]) == ["gate_pattern_recall_failed"]
- assert json.loads(values["scorecard"])["score_missing"] is True
- assert json.loads(values["decision_replay_data"])["dispatch_id"] == "dispatch_content"
- assert json.loads(values["raw_payload"])["decision_replay_data"]["policy_bundle_hash"] == "hash_001"
- def test_database_runtime_preserves_v4_score_and_walk_json_contract():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.append_jsonl(
- "run_001",
- "rule_decisions.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "decision_id": "d_v4_001",
- "policy_bundle_id": "policy_bundle_v4",
- "rule_pack_id": "douyin_content_discovery_rule_pack_v4",
- "rule_pack_version": "4.0.0",
- "strategy_version": "V4",
- "decision_target_type": "content",
- "decision_target_id": "content_001",
- "decision_action": "ADD_TO_CONTENT_POOL",
- "decision_reason_code": "v4_query_and_platform_pass",
- "search_query_effect_status": "success",
- "score": 80,
- "scorecard": {
- "schema_version": "v4_scorecard.v1",
- "query_relevance_score": 82,
- "platform_performance_score": 78,
- "missing_observable_fields": ["view_count"],
- },
- "decision_replay_data": {
- "policy_bundle_hash": "hash_v4",
- "rule_pack_id": "douyin_content_discovery_rule_pack_v4",
- "rule_pack_version": "4.0.0",
- "dispatch_id": "dispatch_content_v4",
- "strategy_version": "V4",
- "allow_walk": True,
- "walk_gate_snapshot": {
- "query_relevance_score": 82,
- "platform_performance_score": 78,
- "score": 80,
- },
- },
- "raw_payload": {
- "decision_id": "d_v4_001",
- "v4_contract": {
- "query_relevance_score": 82,
- "platform_performance_score": 78,
- },
- },
- }
- ],
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_rule_decisions`" in sql
- scorecard = json.loads(values["scorecard"])
- replay_data = json.loads(values["decision_replay_data"])
- assert scorecard["schema_version"] == "v4_scorecard.v1"
- assert scorecard["query_relevance_score"] == 82
- assert scorecard["platform_performance_score"] == 78
- assert scorecard["missing_observable_fields"] == ["view_count"]
- assert replay_data["allow_walk"] is True
- assert replay_data["walk_gate_snapshot"]["score"] == 80
- def test_database_runtime_preserves_p5_search_clue_aggregation_in_raw_payload():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.append_jsonl(
- "run_001",
- "search_clues.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "clue_id": "clue_001",
- "search_query_id": "q_001",
- "search_query": "爱国情感",
- "discovery_start_source": "pattern_seed",
- "previous_discovery_step": "search_query_generated",
- "result_count": 1,
- "pooled_content_count": 0,
- "review_content_count": 0,
- "pending_content_count": 0,
- "rejected_content_count": 1,
- "search_query_effect_status": "rule_blocked",
- "query_aggregation_id": "agg_query_rule_blocked",
- "walk_next_step": "stop_search_query",
- "raw_payload": {
- "clue_id": "clue_001",
- "query_aggregation_id": "agg_query_rule_blocked",
- "effect_status_counts": {"rule_blocked": 1},
- },
- }
- ],
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_search_clues`" in sql
- assert "query_aggregation_id" not in values
- assert values["search_query_effect_status"] == "rule_blocked"
- assert json.loads(values["raw_payload"])["query_aggregation_id"] == "agg_query_rule_blocked"
- def test_database_runtime_writes_failed_search_clue_and_platform_query_failed_event():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- query_failure = {
- "search_query_id": "q_001",
- "search_query": "接口失败",
- "search_query_generation_method": "item_single",
- "status": "failed",
- "error_code": "PLATFORM_REQUEST_FAILED",
- "message": "platform query failed",
- }
- store.append_jsonl(
- "run_001",
- "search_clues.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "clue_id": "clue_001",
- "search_query_id": "q_001",
- "search_query": "接口失败",
- "discovery_start_source": "pattern_itemset",
- "previous_discovery_step": "pattern_search_query",
- "result_count": 0,
- "pooled_content_count": 0,
- "review_content_count": 0,
- "pending_content_count": 0,
- "rejected_content_count": 0,
- "search_query_effect_status": "failed",
- "query_aggregation_id": "platform_query_failure",
- "walk_next_step": "stop_search_query",
- "raw_payload": {
- "clue_id": "clue_001",
- "query_failure": query_failure,
- },
- }
- ],
- )
- store.append_jsonl(
- "run_001",
- "run_events.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "event_id": "evt_001",
- "event_type": "platform_query_failed",
- "status": "failed",
- "input_ref": "search_queries.jsonl:q_001",
- "output_ref": "search_clues.jsonl",
- "error_code": "PLATFORM_REQUEST_FAILED",
- "message": "platform query failed",
- "raw_payload": {
- "event_id": "evt_001",
- "query_failure": query_failure,
- },
- }
- ],
- )
- clue_sql, clue_params = connection.statements[-2]
- clue_values = _insert_values(clue_sql, clue_params)
- assert "INSERT INTO `content_agent_search_clues`" in clue_sql
- assert "ON DUPLICATE KEY UPDATE" in clue_sql
- assert clue_values["search_query_effect_status"] == "failed"
- assert clue_values["walk_next_step"] == "stop_search_query"
- assert json.loads(clue_values["raw_payload"])["query_failure"]["search_query_id"] == "q_001"
- event_sql, event_params = connection.statements[-1]
- event_values = _insert_values(event_sql, event_params)
- assert "INSERT INTO `content_agent_run_events`" in event_sql
- assert "ON DUPLICATE KEY UPDATE" in event_sql
- assert event_values["event_type"] == "platform_query_failed"
- assert event_values["status"] == "failed"
- assert event_values["error_code"] == "PLATFORM_REQUEST_FAILED"
- def test_database_runtime_writes_publish_jobs_db_only_records():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.write_publish_jobs(
- "run_001",
- "policy_run_001",
- [
- {
- "publish_job_id": "publish_job_001",
- "platform_content_id": "7390000000000000000",
- "job_status": "created",
- "trigger_mode": "manual_review",
- "request_payload": {
- "decision_id": "decision_001",
- "source_path_record_ids": ["path_001"],
- },
- "response_payload": {},
- }
- ],
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_publish_jobs`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert values["run_id"] == "run_001"
- assert values["policy_run_id"] == "policy_run_001"
- assert values["publish_job_id"] == "publish_job_001"
- assert values["job_status"] == "created"
- assert values["trigger_mode"] == "manual_review"
- assert json.loads(values["request_payload"])["decision_id"] == "decision_001"
- def test_database_runtime_writes_author_assets():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.write_author_assets(
- [
- {
- "author_asset_id": "author_asset_001",
- "platform": "douyin",
- "platform_author_id": "author_001",
- "author_display_name": "作者一号",
- "asset_status": "active",
- "source_type": "runtime_author_work",
- "validation_status": "validated",
- "eligible_as_source": 1,
- "elderly_ratio": 0.72,
- "elderly_tgi": 138,
- "content_tags": ["人物故事"],
- "source_run_id": "run_001",
- "source_policy_run_id": "policy_run_001",
- "profile_snapshot": {"sample_count": 9},
- "evidence_refs": {"decision_ids": ["d_001"]},
- "raw_payload": {"author_asset_id": "author_asset_001"},
- }
- ]
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_author_assets`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert values["author_asset_id"] == "author_asset_001"
- assert values["platform_author_id"] == "author_001"
- assert values["eligible_as_source"] == 1
- assert values["elderly_ratio"] == 0.72
- assert values["elderly_tgi"] == 138
- assert json.loads(values["content_tags"]) == ["人物故事"]
- assert json.loads(values["profile_snapshot"])["sample_count"] == 9
- assert json.loads(values["evidence_refs"])["decision_ids"] == ["d_001"]
- def test_database_runtime_writes_author_asset_roles():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.write_author_asset_roles(
- [
- {
- "author_asset_id": "author_asset_001",
- "role": "source_seed",
- "role_status": "active",
- "role_reason_code": "p7_author_asset_eligible",
- "assigned_by": "system",
- "source_run_id": "run_001",
- "raw_payload": {"role": "source_seed"},
- }
- ]
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_author_asset_roles`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert values["author_asset_id"] == "author_asset_001"
- assert values["role"] == "source_seed"
- assert values["assigned_by"] == "system"
- assert json.loads(values["raw_payload"])["role"] == "source_seed"
- def test_database_runtime_writes_search_clue_assets():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.write_search_clue_assets(
- [
- {
- "search_clue_asset_id": "search_clue_asset_001",
- "platform": "douyin",
- "clue_type": "search_query",
- "normalized_clue_text": "银发旅行",
- "display_clue_text": "银发旅行",
- "promotion_status": "promoted",
- "reusable_priority": 1,
- "can_seed_next_run": 1,
- "first_seen_run_id": "run_001",
- "first_seen_policy_run_id": "policy_run_001",
- "summary_metrics": {"pooled_content_count": 1},
- "raw_payload": {"promotion_reason": "success_search_clue"},
- }
- ]
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_search_clue_assets`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert values["search_clue_asset_id"] == "search_clue_asset_001"
- assert values["can_seed_next_run"] == 1
- assert json.loads(values["summary_metrics"])["pooled_content_count"] == 1
- def test_database_runtime_writes_search_clue_asset_evidence():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.write_search_clue_asset_evidence(
- [
- {
- "evidence_id": "search_clue_evidence_001",
- "search_clue_asset_id": "search_clue_asset_001",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "clue_id": "clue_001",
- "search_query_id": "q_001",
- "pooled_content_count": 1,
- "source_path_record_ids": ["path_001"],
- "decision_ids": ["decision_001"],
- "performance_feedback_refs": [],
- "raw_payload": {"clue_id": "clue_001"},
- }
- ]
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_search_clue_asset_evidence`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert values["run_id"] == "run_001"
- assert values["clue_id"] == "clue_001"
- assert json.loads(values["source_path_record_ids"]) == ["path_001"]
- assert json.loads(values["decision_ids"]) == ["decision_001"]
- def test_database_runtime_update_final_output_upserts_validation_status():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.update_json(
- "run_001",
- "final_output.json",
- {
- "schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "summary": {"run_path_complete": True},
- "validation_status": "pass",
- },
- )
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert "INSERT INTO `content_agent_final_outputs`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["validation_status"] == "pass"
- assert json.loads(values["summary"])["run_path_complete"] is True
- def test_database_runtime_preserves_m5_explanation_payloads():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- explanation = {
- "schema_version": "v4_decision_explanation.v1",
- "query_relevance_score": 80,
- "platform_performance_score": 70,
- "score": 75,
- "allow_walk": True,
- "walk_gate_snapshot": {"score": 75},
- }
- v4_summary = {
- "schema_version": "v4_strategy_review_summary.v1",
- "score_buckets": {"pool": 1},
- "allow_walk_distribution": {"allowed": 1, "denied": 0, "missing": 0},
- "walk_gate_review": {"v4_gate_distribution": {"allowed": 1}},
- }
- store.update_json(
- "run_001",
- "final_output.json",
- {
- "schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "summary": {},
- "validation_status": "pass",
- "decision_records": [
- {
- "decision_id": "decision_001",
- "v4_explanation": explanation,
- }
- ],
- },
- )
- store.update_json(
- "run_001",
- "strategy_review.json",
- {
- "schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "review_id": "review_001",
- "review_status": "generated",
- "summary": {},
- "v4_summary": v4_summary,
- "effective_search_queries": [],
- "weak_search_queries": [],
- "top_reject_reasons": [],
- "productive_paths": [],
- "suggestions": [],
- "raw_payload": {"v4_summary": v4_summary},
- },
- )
- final_values = _insert_values(*connection.statements[-2])
- final_payload = json.loads(final_values["final_output"])
- assert final_payload["decision_records"][0]["v4_explanation"] == explanation
- review_values = _insert_values(*connection.statements[-1])
- review_payload = json.loads(review_values["raw_payload"])
- assert review_payload["v4_summary"] == v4_summary
- def test_database_runtime_reads_performance_feedback_payloads():
- connection = FakeConnection()
- connection.select_all_result = [
- {
- "raw_payload": json.dumps(
- {
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "feedback_id": "feedback_001",
- "completion_rate": 0.8,
- }
- )
- }
- ]
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- rows = store.read_performance_feedback("run_001", "policy_run_001")
- sql, params = connection.statements[-1]
- assert "FROM `content_agent_performance_feedback`" in sql
- assert params == ["run_001", "policy_run_001"]
- assert rows[0]["feedback_id"] == "feedback_001"
- assert rows[0]["raw_payload"]["completion_rate"] == 0.8
- def test_database_runtime_read_jsonl_reconstructs_runtime_payload():
- connection = FakeConnection()
- connection.select_all_result = [
- {
- "raw_payload": json.dumps(
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "search_query_id": "q_001",
- }
- )
- }
- ]
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- rows = store.read_jsonl("run_001", "search_queries.jsonl")
- assert rows[0]["search_query_id"] == "q_001"
- assert rows[0]["raw_payload"]["search_query_id"] == "q_001"
- def test_database_runtime_read_jsonl_reconstructs_rule_decision_from_formal_columns():
- connection = FakeConnection()
- connection.select_all_result = [
- {
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "decision_id": "d_001",
- "policy_bundle_id": "bundle_001",
- "rule_pack_id": "rule_pack_001",
- "rule_pack_version": "4.0.0",
- "strategy_version": "V4",
- "decision_target_type": "content",
- "decision_target_id": "content_001",
- "decision_action": "ADD_TO_CONTENT_POOL",
- "decision_reason_code": "v4_query_and_platform_pass",
- "search_query_effect_status": "success",
- "score": Decimal("75.50"),
- "triggered_blocking_rules": json.dumps([]),
- "scorecard": json.dumps({"schema_version": "v4_scorecard.v1", "total_score": 75.5}),
- "source_evidence": json.dumps({"source_post_id": "post_001"}),
- "decision_replay_data": json.dumps({"allow_walk": True}),
- "raw_payload": json.dumps(
- {
- "record_schema_version": "runtime_record.v1",
- "strategy_id": "strategy_001",
- "policy_bundle_hash": "hash_001",
- }
- ),
- }
- ]
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- rows = store.read_jsonl("run_001", "rule_decisions.jsonl")
- sql, params = connection.statements[-1]
- assert "FROM `content_agent_rule_decisions`" in sql
- assert "`scorecard`" in sql
- assert params == ["run_001"]
- row = rows[0]
- assert row["record_schema_version"] == "runtime_record.v1"
- assert row["decision_id"] == "d_001"
- assert row["score"] == 75.5
- assert row["scorecard"]["schema_version"] == "v4_scorecard.v1"
- assert row["source_evidence"]["source_post_id"] == "post_001"
- assert row["decision_replay_data"]["allow_walk"] is True
- assert row["raw_payload"] == {
- "record_schema_version": "runtime_record.v1",
- "strategy_id": "strategy_001",
- "policy_bundle_hash": "hash_001",
- }
- def test_database_runtime_rejects_forbidden_raw_payload_keys_in_lists():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- try:
- store.append_jsonl(
- "run_001",
- "search_queries.jsonl",
- [
- {
- "record_schema_version": "runtime_record.v1",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "search_query_id": "q_001",
- "search_query": "对比分析",
- "raw_payload": {"items": [{"dsn": "should_not_be_stored"}]},
- }
- ],
- )
- except ValueError as exc:
- assert "forbidden key" in str(exc)
- else:
- raise AssertionError("expected forbidden raw_payload key to be rejected")
- def test_database_runtime_update_run_record_ignores_empty_sanitized_updates():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.update_run_record("run_001", {"unknown_field": "ignored", "status": None})
- assert connection.statements == []
- def test_database_runtime_update_run_record_persists_platform_failure_detail():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.update_run_record(
- "run_001",
- {
- "status": "failed",
- "error_code": "PLATFORM_REQUEST_FAILED",
- "error_detail": {
- "query_failures": [
- {
- "search_query_id": "q_001",
- "status": "failed",
- "error_code": "PLATFORM_REQUEST_FAILED",
- }
- ]
- },
- },
- )
- sql, params = connection.statements[-1]
- assert "UPDATE `content_agent_runs` SET" in sql
- assert "`error_detail` = %s" in sql
- assert json.loads(params[2])["query_failures"][0]["search_query_id"] == "q_001"
- assert params[-1] == "run_001"
- def test_database_runtime_upserts_media_pipeline_task():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- task = {
- "task_id": "mpt_001",
- "task_type": "oss_upload",
- "idempotency_key": "run_001:q_001_p01:douyin:c1:oss_upload",
- "run_id": "run_001",
- "policy_run_id": "policy_run_001",
- "batch_id": "q_001_p01_initial_top3",
- "platform": "douyin",
- "platform_content_id": "c1",
- "status": "queued",
- "attempt_count": 0,
- "lease_until": None,
- "next_retry_at": None,
- "input_payload": {"play_url_present": True},
- "result_payload": None,
- "error_summary": None,
- }
- stored = store.enqueue_media_pipeline_task(task)
- sql, params = connection.statements[-1]
- values = _insert_values(sql, params)
- assert stored == task
- assert "INSERT INTO `content_agent_media_pipeline_tasks`" in sql
- assert "ON DUPLICATE KEY UPDATE" in sql
- assert values["schema_version"] == "content_agent.v1"
- assert values["task_id"] == "mpt_001"
- assert json.loads(values["input_payload"])["play_url_present"] is True
- def test_database_runtime_updates_media_pipeline_task_status():
- connection = FakeConnection()
- store = DatabaseRuntimeStore(_config(), connection_factory=lambda: connection)
- store.update_media_pipeline_task(
- "mpt_001",
- {
- "status": "completed",
- "attempt_count": 1,
- "result_payload": {"oss_url_present": True},
- },
- )
- sql, params = connection.statements[-1]
- assert "UPDATE `content_agent_media_pipeline_tasks` SET" in sql
- assert "`status` = %s" in sql
- assert "`result_payload` = %s" in sql
- assert json.loads(params[2])["oss_url_present"] is True
- assert params[-1] == "mpt_001"
- def test_business_modules_do_not_import_or_name_database_tables():
- root = Path("content_agent/business_modules")
- text = "\n".join(path.read_text(encoding="utf-8") for path in root.rglob("*.py"))
- assert not re.search(
- r"pymysql|sqlalchemy|psycopg|sqlite3|SELECT |INSERT |UPDATE |DELETE |SHOW |CREATE |ALTER |content_agent_",
- text,
- )
- def _config():
- return ContentSupplyDbConfig(
- host="127.0.0.1",
- port=3306,
- user="content_rw",
- password="dummy_password",
- database="content-deconstruction-supply",
- )
- def _insert_values(sql, params):
- match = re.search(r"\((.*?)\) VALUES", sql)
- assert match, sql
- columns = [part.strip().strip("`") for part in match.group(1).split(",")]
- return dict(zip(columns, params))
|