test_p8_performance_feedback.py 4.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125
  1. from content_agent.business_modules import learning_review
  2. from content_agent.integrations.runtime_files import LocalRuntimeFileStore
  3. from tests.test_p8_strategy_review import _write_minimal_runtime
  4. from tests.test_v5_golden_corpus_baseline import _load_historical_decisions
  5. class FeedbackRuntime(LocalRuntimeFileStore):
  6. def __init__(self, *args, feedback_rows=None, **kwargs):
  7. super().__init__(*args, **kwargs)
  8. self.feedback_rows = feedback_rows or []
  9. def read_performance_feedback(self, run_id: str, policy_run_id: str):
  10. return [
  11. row
  12. for row in self.feedback_rows
  13. if row["run_id"] == run_id and row["policy_run_id"] == policy_run_id
  14. ]
  15. def test_strategy_review_summarizes_fake_performance_feedback(tmp_path):
  16. run_id = "run_feedback"
  17. policy_run_id = "policy_feedback"
  18. runtime = FeedbackRuntime(
  19. tmp_path / "runtime",
  20. feedback_rows=[
  21. {
  22. "run_id": run_id,
  23. "policy_run_id": policy_run_id,
  24. "feedback_id": "feedback_001",
  25. "feedback_status": "available",
  26. "completion_rate": 0.72,
  27. "share_rate": 0.08,
  28. "average_watch_seconds": 18.5,
  29. }
  30. ],
  31. )
  32. runtime.prepare_run(run_id)
  33. _write_minimal_runtime(runtime, run_id, policy_run_id)
  34. review = learning_review.run(run_id, policy_run_id, runtime)
  35. assert review["performance_feedback"]["performance_feedback_status"] == "available"
  36. assert review["performance_feedback"]["feedback_count"] == 1
  37. assert review["performance_feedback"]["average_completion_rate"] == 0.72
  38. assert any(
  39. item["recommendation_type"] == "performance_feedback"
  40. and item["suggested_action"] == "review_feedback_before_strategy_change"
  41. for item in review["recommendations"]
  42. )
  43. assert review["rule_review"]["decision_distribution"]["ADD_TO_CONTENT_POOL"] == 1
  44. def _fake_feedback_for_decision(decision, platform="douyin"):
  45. platform_content_id = decision["decision_target_id"]
  46. return {
  47. "run_id": decision["run_id"],
  48. "policy_run_id": decision["policy_run_id"],
  49. "feedback_id": f"fake_m4_{decision['policy_run_id']}_{platform_content_id}",
  50. "feedback_source": "m4_fake_feedback",
  51. "feedback_status": "available",
  52. "platform": platform,
  53. "platform_content_id": platform_content_id,
  54. "completion_rate": 0.72,
  55. }
  56. def _join_feedback(feedback_rows, discovered_rows, *, use_content_discovery_id=False):
  57. key_field = "content_discovery_id" if use_content_discovery_id else "platform_content_id"
  58. discovered_keys = {
  59. (row["run_id"], row["policy_run_id"], row.get(key_field))
  60. for row in discovered_rows
  61. }
  62. return [
  63. row
  64. for row in feedback_rows
  65. if (row["run_id"], row["policy_run_id"], row["platform_content_id"]) in discovered_keys
  66. ]
  67. def test_fake_feedback_join_matches_m0_golden_samples():
  68. decisions = _load_historical_decisions("v4_douyin_57663")
  69. wanted_actions = {
  70. "ADD_TO_CONTENT_POOL",
  71. "KEEP_CONTENT_FOR_REVIEW",
  72. "REJECT_CONTENT",
  73. "TECHNICAL_RETRY_REQUIRED",
  74. }
  75. selected = []
  76. seen_actions = set()
  77. for decision in decisions:
  78. action = decision["decision_action"]
  79. if action in wanted_actions and action not in seen_actions:
  80. selected.append(decision)
  81. seen_actions.add(action)
  82. if seen_actions == wanted_actions:
  83. break
  84. feedback_rows = [_fake_feedback_for_decision(decision) for decision in selected]
  85. discovered_rows = [
  86. {
  87. "run_id": decision["run_id"],
  88. "policy_run_id": decision["policy_run_id"],
  89. "platform_content_id": decision["decision_target_id"],
  90. "content_discovery_id": f"discovery_{decision['decision_target_id']}",
  91. }
  92. for decision in selected
  93. ]
  94. assert seen_actions == wanted_actions
  95. assert _join_feedback(feedback_rows, discovered_rows) == feedback_rows
  96. def test_fake_feedback_join_uses_platform_content_id_not_content_discovery_id():
  97. decision = _load_historical_decisions("v4_douyin_57663")[0]
  98. feedback = [_fake_feedback_for_decision(decision)]
  99. discovered = [
  100. {
  101. "run_id": decision["run_id"],
  102. "policy_run_id": decision["policy_run_id"],
  103. "platform_content_id": decision["decision_target_id"],
  104. "content_discovery_id": f"discovery_{decision['decision_target_id']}",
  105. }
  106. ]
  107. assert _join_feedback(feedback, discovered, use_content_discovery_id=True) == []