Commit bbb6ce14 bbb6ce1447e193340855290021205a90c07891fe by 沈秋雨

feat(review): 优化审核结果提交和重复数据合并逻辑

1 parent 2393078b
......@@ -646,6 +646,7 @@
<option value="unsure">待确认</option>
<option value="reviewed">已审核</option>
<option value="submitted">已提交</option>
<option value="deleted">已删除</option>
<option value="">全部</option>
<option value="l1_hit">L1 命中</option>
<option value="conflict">L1/L2 冲突</option>
......@@ -661,7 +662,7 @@
<option value="200">200</option>
</select>
<button id="exportBtn">导出标注</button>
<button id="importReviewedBtn">入库审核通过</button>
<button id="importReviewedBtn">提交审核结果</button>
<div class="toolbar-right">
<span>审核人:</span><span id="reviewerDisplay"></span>
<button id="logoutBtn" class="secondary">登出</button>
......@@ -1002,7 +1003,7 @@
group.review_claim_expires_at = record.review_claim_expires_at || '';
group.claimed_by_me = group.review_claimed_by === state.reviewer;
group.reviewed = ['approved_import', 'rejected_duplicate', 'deleted'].includes(record.biz_review_status);
group.submitted = group.reviewed && record.staging_status === 'imported';
group.submitted = group.reviewed && ['imported', 'skipped'].includes(record.staging_status);
}
function currentQuery() {
......@@ -1059,12 +1060,12 @@
metric('待审核', row.pending_review_count || '0'),
metric('已审核', row.reviewed_count || '0'),
metric('已提交', row.submitted_count || '0'),
metric('确认入库', row.approved_count || '0'),
metric('确认不入库', row.rejected_count || '0'),
metric('确认不重复', row.approved_count || '0'),
metric('确认重复', row.rejected_count || '0'),
metric('已删除', row.deleted_count || '0'),
metric('new', row.new_count || '-'),
metric('merge', row.duplicate_count || '0'),
metric('skip', row.skip_count || '0')
metric('新歌', row.new_count || '-'),
metric('重复', row.duplicate_count || '0'),
metric('已跳过', row.skip_count || '0')
].join('');
}
......@@ -1182,20 +1183,21 @@
</div>`;
}
async function updateStagingReview(query, decision, note) {
async function updateStagingReview(query, decision, note, candidateId = '') {
return reviewPostJSON('/api/staging-review', {
staging_id: String(query.staging_id),
expected_version: Number(query.review_version || 0),
client_token: state.clientToken,
reviewer: state.reviewer,
decision,
note
note,
candidate_id: candidateId
});
}
const _DB_REVIEW_LABELS = {
approved_import: '确认入库(待入库)',
rejected_duplicate: '确认不入库',
approved_import: '确认不重复(待入库)',
rejected_duplicate: '确认重复(待合并作者)',
unsure: '待确认',
deleted: '已删除',
};
......@@ -1293,8 +1295,9 @@
const needsManualReview = rowDecision(query) === 'review';
const hasFinalReview = Boolean(query.review_status) &&
!['pending', 'not_required', 'unsure'].includes(query.review_status);
const canUndoReview = hasFinalReview && query.staging_status !== 'imported';
const canDeleteRecord = query.staging_status !== 'imported';
const isSubmitted = ['imported', 'skipped'].includes(query.staging_status);
const canUndoReview = hasFinalReview && query.staging_status === 'staged';
const canDeleteRecord = !['imported', 'skipped'].includes(query.staging_status);
// Sticky area: 新入库歌词 + 召回候选(各含音频播放器)
els.stickyInfo.innerHTML = `
......@@ -1339,7 +1342,7 @@
<div class="section-head">
<h2 class="section-title">人工标注</h2>
${hasFinalReview
? `<span class="pill ${query.review_status === 'approved_import' ? 'new' : query.review_status === 'rejected_duplicate' ? 'duplicate' : query.review_status === 'deleted' ? 'duplicate' : ''}">${query.review_status === 'deleted' ? '已删除' : query.staging_status === 'imported' ? '已提交' : '已审核'}</span>`
? `<span class="pill ${query.review_status === 'approved_import' ? 'new' : query.review_status === 'rejected_duplicate' ? 'duplicate' : query.review_status === 'deleted' ? 'duplicate' : ''}">${query.review_status === 'deleted' ? '已删除' : isSubmitted ? '已提交' : '已审核'}</span>`
: ''}
</div>
<div class="section-body">
......@@ -1347,7 +1350,9 @@
<div class="review-bar" style="align-items:center">
<span style="flex:1">
<b>${esc(query.review_status === 'approved_import' && query.staging_status === 'imported'
? '确认入库(已提交)'
? '确认不重复(已入库)'
: query.review_status === 'rejected_duplicate' && query.staging_status === 'skipped'
? '确认重复(作者已合并)'
: (_DB_REVIEW_LABELS[query.review_status] || query.review_status))}</b>
${query.reviewed_by ? ` · 审核人 ${esc(query.reviewed_by)}` : ''}
${query.reviewed_at ? ` · ${esc(query.reviewed_at)}` : ''}
......@@ -1359,7 +1364,7 @@
<div class="review-bar">
${['duplicate', 'not_duplicate', 'unsure'].map(value => `
<button class="review-choice ${selectedReview.final_decision === value ? 'active' : ''}" data-review="${value}">
${value === 'duplicate' ? '确认不入库' : value === 'not_duplicate' ? '确认入库' : '待确认'}
${value === 'duplicate' ? '确认重复' : value === 'not_duplicate' ? '确认不重复' : '待确认'}
</button>`).join('')}
<input id="reviewNote" value="${esc(selectedReview.note || '')}" placeholder="人工备注">
<button id="deleteReviewBtn" class="secondary" style="color:#c0392b;border-color:#c0392b">删除</button>
......@@ -1411,7 +1416,9 @@
renderDetail();
return;
}
const payload = await updateStagingReview(query, dbDecision, note);
const payload = await updateStagingReview(
query, dbDecision, note, localDecision === 'duplicate' ? (candidate?.candidate_id || '') : ''
);
for (const row of state.queryRows) applyReviewSnapshot(row, payload.record);
// 刷新统计和加载列表并行
await Promise.all([reloadSummary(), loadGroups()]);
......@@ -1530,12 +1537,12 @@
alert('历史报表仅供查看,请在 staging-db 中执行入库。');
return;
}
if (!confirm('确认入库数据库中所有"确认入库(待入库)"样本?')) return;
if (!confirm('确认提交所有审核结果?“确认不重复”将入库,“确认重复”将把词/曲作者增量合并到已有记录。')) return;
els.importReviewedBtn.disabled = true;
els.importReviewedBtn.textContent = '入库中...';
try {
const payload = await reviewPostJSON('/api/import-reviewed');
alert(`入库完成:${payload.inserted_count} 条。结果文件:${payload.review_csv}`);
alert(`提交完成:不重复入库 ${payload.inserted_count} 条,重复样本 ${payload.duplicate_count || 0} 条,其中作者有增量 ${payload.author_merged_count || 0} 条。结果文件:${payload.review_csv}`);
// 刷新统计和加载列表并行
await Promise.all([reloadSummary(), loadGroups()]);
renderAll();
......@@ -1543,14 +1550,17 @@
if (err.status === 409 && err.payload?.code === 'import_conflict') {
const conflicts = err.payload.conflicts || [];
const details = conflicts.slice(0, 20).map((row, index) => {
const types = (row.conflict_types || []).map(type =>
type === 'primary_key' ? '主键 ID 冲突' : '来源键冲突'
).join(' + ') || '唯一键冲突';
const typeLabels = {
primary_key: '主键 ID 冲突',
source_key: '来源键冲突',
duplicate_target_missing: '未找到已入库的重复目标'
};
const types = (row.conflict_types || []).map(type => typeLabels[type] || type).join(' + ') || '数据冲突';
const targetIds = [row.id_conflict_target_id, row.source_conflict_target_id]
.filter(value => value !== null && value !== undefined)
.filter((value, i, values) => values.indexOf(value) === i)
.join('/');
return `${index + 1}. ${row.name || '(无歌名)'} | source=${row.source_table_name || '-'}:${row.source_song_id || '-'} | staging_id=${row.staging_id} | 待入库ID=${row.incoming_id || '-'} | ${types} | 正式表ID=${targetIds || '-'}`;
return `${index + 1}. ${row.name || '(无歌名)'} | source=${row.source_table_name || '-'}:${row.source_song_id || '-'} | staging_id=${row.staging_id} | 待入库ID=${row.incoming_id || '-'} | 重复候选=${row.matched_song_id || '-'} | ${types} | 正式表ID=${targetIds || '-'}`;
});
if (conflicts.length > 20) details.push(`……另有 ${conflicts.length - 20} 条未展开`);
alert([
......@@ -1565,7 +1575,7 @@
}
} finally {
els.importReviewedBtn.disabled = false;
els.importReviewedBtn.textContent = '入库审核通过';
els.importReviewedBtn.textContent = '提交审核结果';
}
}
......@@ -1659,7 +1669,7 @@
group.review_claim_expires_at = remote.review_claim_expires_at || '';
group.claimed_by_me = Boolean(remote.claimed_by_me);
group.reviewed = ['approved_import', 'rejected_duplicate', 'deleted'].includes(remote.biz_review_status);
group.submitted = group.reviewed && remote.staging_status === 'imported';
group.submitted = group.reviewed && ['imported', 'skipped'].includes(remote.staging_status);
}
const query = currentQuery();
const selectedRemote = query ? statuses.get(String(query.staging_id)) : null;
......
......@@ -8,17 +8,20 @@ import csv
import json
import mimetypes
import os
import re
import secrets
import subprocess
import sys
import threading
import time
import unicodedata
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from pathlib import Path
from urllib.parse import parse_qs, unquote, urlparse
import pymysql
import requests
import opencc
from dbutils.pooled_db import PooledDB
from dotenv import load_dotenv
......@@ -79,6 +82,9 @@ REVIEW_SCHEMA_COLUMNS = {
"review_claim_expires_at": "datetime DEFAULT NULL COMMENT '领取过期时间'",
"review_version": "bigint(20) NOT NULL DEFAULT '0' COMMENT '审核乐观锁版本'",
}
_AUTHOR_SEP_RE = re.compile(r"[,,/;;、]")
_AUTHOR_T2S = opencc.OpenCC("t2s")
_AUTHOR_METADATA_PUNCT = set(' \t\n\r,。!?;:、"“”‘’·…—~!¥()【】《》〈〉「」『』﹏,.:;!?()[]{}<>|/\\_-')
class ReviewConflictError(Exception):
......@@ -250,8 +256,8 @@ def _staging_summary() -> list[dict[str, str]]:
SUM(CASE WHEN biz_review_status = 'approved_import' THEN 1 ELSE 0 END) AS approved_count,
SUM(CASE WHEN biz_review_status = 'rejected_duplicate' THEN 1 ELSE 0 END) AS rejected_count,
SUM(CASE WHEN biz_review_status = 'deleted' THEN 1 ELSE 0 END) AS deleted_count,
SUM(CASE WHEN biz_review_status IN ('approved_import', 'rejected_duplicate', 'deleted') AND staging_status <> 'imported' THEN 1 ELSE 0 END) AS reviewed_count,
SUM(CASE WHEN biz_review_status IN ('approved_import', 'rejected_duplicate', 'deleted') AND staging_status = 'imported' THEN 1 ELSE 0 END) AS submitted_count
SUM(CASE WHEN biz_review_status IN ('approved_import', 'rejected_duplicate', 'deleted') AND staging_status NOT IN ('imported', 'skipped') THEN 1 ELSE 0 END) AS reviewed_count,
SUM(CASE WHEN biz_review_status IN ('approved_import', 'rejected_duplicate', 'deleted') AND staging_status IN ('imported', 'skipped') THEN 1 ELSE 0 END) AS submitted_count
FROM (
SELECT source_song_id, dedup_action, biz_review_status, staging_status
FROM {TARGET_TABLE_NAME_TMP}
......@@ -356,10 +362,13 @@ def _staging_groups_response(
where_clauses.append("s.dedup_action = 'skip'")
elif decision == 'reviewed':
where_clauses.append("s.biz_review_status IN ('approved_import', 'rejected_duplicate', 'deleted')")
where_clauses.append("s.staging_status <> 'imported'")
where_clauses.append("s.staging_status NOT IN ('imported', 'skipped')")
elif decision == 'submitted':
where_clauses.append("s.biz_review_status IN ('approved_import', 'rejected_duplicate', 'deleted')")
where_clauses.append("s.staging_status = 'imported'")
where_clauses.append("s.staging_status IN ('imported', 'skipped')")
elif decision == 'deleted':
where_clauses.append("s.biz_review_status = 'deleted'")
where_clauses.append("s.staging_status = 'deleted'")
elif decision == 'pending':
where_clauses.append("(s.biz_review_status IS NULL OR s.biz_review_status = 'pending')")
where_clauses.append("s.dedup_action = 'review'")
......@@ -427,7 +436,7 @@ def _staging_groups_response(
"reviewed": row.get("biz_review_status") in {"approved_import", "rejected_duplicate", "deleted"},
"submitted": (
row.get("biz_review_status") in {"approved_import", "rejected_duplicate", "deleted"}
and row.get("staging_status") == "imported"
and row.get("staging_status") in {"imported", "skipped"}
),
"staging_id": str(row.get("staging_id") or ""),
"review_version": int(row.get("review_version") or 0),
......@@ -674,6 +683,7 @@ def _update_staging_review(
reviewer: str,
expected_version: int,
client_token: str,
candidate_id: str = "",
) -> dict[str, object]:
"""Persist one claimed review using optimistic concurrency control."""
if decision == "pending":
......@@ -686,6 +696,11 @@ def _update_staging_review(
biz_review_status = 'deleted', reviewed_by = %s, reviewed_at = NOW(),
staging_status = 'deleted'
"""
elif decision == "rejected_duplicate":
status_fields = """
biz_review_status = %s, reviewed_by = %s, reviewed_at = NOW(),
matched_song_id = %s
"""
else:
status_fields = "biz_review_status = %s, reviewed_by = %s, reviewed_at = NOW()"
......@@ -694,9 +709,31 @@ def _update_staging_review(
params = [note, staging_id, expected_version, client_token, reviewer]
elif decision == "deleted":
params = [reviewer, note, staging_id, expected_version, client_token, reviewer]
elif decision == "rejected_duplicate":
candidate_id = candidate_id.strip()
if not candidate_id:
raise ValueError("确认重复前请先选择一条召回候选")
params = [decision, reviewer, candidate_id, note, staging_id, expected_version, client_token, reviewer]
else:
params = [decision, reviewer, note, staging_id, expected_version, client_token, reviewer]
with _target_conn() as conn, conn.cursor() as cursor:
if decision == "rejected_duplicate":
cursor.execute(
f"SELECT matched_song_id, recalled_candidates FROM {TARGET_TABLE_NAME_TMP} WHERE staging_id = %s",
(staging_id,),
)
before = cursor.fetchone() or {}
allowed_ids = {str(before.get("matched_song_id") or "")}
recalled = before.get("recalled_candidates")
try:
candidates = json.loads(recalled) if isinstance(recalled, str) else (recalled or [])
allowed_ids.update(str(item.get("id") or "") for item in candidates if isinstance(item, dict))
except (json.JSONDecodeError, TypeError):
pass
allowed_ids.discard("")
if candidate_id not in allowed_ids:
conn.rollback()
raise ValueError("选中的重复候选不在当前召回结果中,请刷新后重试")
cursor.execute(
f"""
UPDATE {TARGET_TABLE_NAME_TMP}
......@@ -724,8 +761,36 @@ def _update_staging_review(
return current
def _merge_author_field(existing: str | None, incoming: str | None) -> str:
"""Match import_hk_songs.py: preserve existing authors and append unique incoming authors."""
if not incoming:
return existing or ""
if not existing:
return incoming
def split_authors(text: str) -> list[str]:
return [author.strip() for author in _AUTHOR_SEP_RE.split(text) if author.strip()]
def normalize(author: str) -> str:
author = unicodedata.normalize("NFKC", author)
author = _AUTHOR_T2S.convert(author.strip().lower())
return "".join(char for char in author if char not in _AUTHOR_METADATA_PUNCT)
existing_authors = split_authors(existing)
incoming_authors = split_authors(incoming)
normalized = {normalize(author) for author in existing_authors}
result = list(existing_authors)
for author in incoming_authors:
key = normalize(author)
if key not in normalized:
result.append(author)
normalized.add(key)
merged = "、".join(result)
return merged if merged != existing else ""
def _import_approved_staging(source_ids: list[str] | None = None) -> dict[str, object]:
"""Import approved rows while holding staging-row locks to prevent duplicate imports."""
"""Finalize reviewed rows: import non-duplicates and merge duplicate authors."""
columns = ", ".join(f"`{column}`" for column in TARGET_COLUMNS)
select_columns = ", ".join(f"s.`{column}`" for column in TARGET_COLUMNS)
with _target_conn() as conn, conn.cursor() as cursor:
......@@ -737,7 +802,8 @@ def _import_approved_staging(source_ids: list[str] | None = None) -> dict[str, o
lock_params.extend(source_ids)
cursor.execute(
f"""
SELECT s.staging_id, s.source_song_id
SELECT s.staging_id, s.id AS incoming_id, s.source_song_id, s.source_table_name, s.name,
s.lyricist, s.composer, s.biz_review_status, s.matched_song_id
FROM {TARGET_TABLE_NAME_TMP} s
INNER JOIN (
SELECT source_song_id, MAX(staging_id) AS max_staging_id
......@@ -745,7 +811,7 @@ def _import_approved_staging(source_ids: list[str] | None = None) -> dict[str, o
GROUP BY source_song_id
) latest ON s.staging_id = latest.max_staging_id
WHERE s.dedup_action = 'review'
AND s.biz_review_status = 'approved_import'
AND s.biz_review_status IN ('approved_import', 'rejected_duplicate')
AND s.staging_status = 'staged'
AND (s.review_claim_expires_at IS NULL OR s.review_claim_expires_at <= NOW())
{source_filter}
......@@ -754,12 +820,42 @@ def _import_approved_staging(source_ids: list[str] | None = None) -> dict[str, o
""",
lock_params,
)
approved_rows = cursor.fetchall()
if not approved_rows:
reviewed_rows = cursor.fetchall()
if not reviewed_rows:
conn.rollback()
raise ValueError("没有可入库的审核通过样本")
raise ValueError("没有可提交的已审核样本")
approved_rows = [row for row in reviewed_rows if row["biz_review_status"] == "approved_import"]
duplicate_rows = [row for row in reviewed_rows if row["biz_review_status"] == "rejected_duplicate"]
staging_ids = [row["staging_id"] for row in approved_rows]
approved_source_ids = [str(row["source_song_id"]) for row in approved_rows]
duplicate_targets: dict[object, int] = {}
unresolved_duplicates = []
for row in duplicate_rows:
target = _resolve_duplicate_target(cursor, row)
if not target:
unresolved_duplicates.append({
"staging_id": row["staging_id"],
"source_song_id": row["source_song_id"],
"source_table_name": row.get("source_table_name"),
"name": row.get("name"),
"incoming_id": row.get("incoming_id"),
"id_conflict_target_id": None,
"source_conflict_target_id": None,
"conflict_types": ["duplicate_target_missing"],
"matched_song_id": row.get("matched_song_id"),
})
else:
duplicate_targets[row["staging_id"]] = int(target["id"])
if unresolved_duplicates:
conn.rollback()
raise ImportConflictError(
f"有 {len(unresolved_duplicates)} 条“确认重复”记录无法找到已入库的重复目标,本次未写入任何数据",
unresolved_duplicates,
)
inserted_count = 0
id_map: dict[str, object] = {}
if approved_rows:
staging_placeholders = ",".join(["%s"] * len(staging_ids))
conflicts = _find_import_conflicts(cursor, staging_ids)
if conflicts:
......@@ -779,7 +875,6 @@ def _import_approved_staging(source_ids: list[str] | None = None) -> dict[str, o
staging_ids,
)
except pymysql.err.IntegrityError as exc:
# 预检与 INSERT 之间仍可能有其他程序写入,回滚后再查一次以返回具体行。
conn.rollback()
conflicts = _find_import_conflicts(cursor, staging_ids)
raise ImportConflictError(
......@@ -798,6 +893,7 @@ def _import_approved_staging(source_ids: list[str] | None = None) -> dict[str, o
approved_source_ids,
)
id_map = {str(row["source_song_id"]): row["target_id"] for row in cursor.fetchall()}
imported_updated = 0
for approved_row in approved_rows:
sid = str(approved_row["source_song_id"])
......@@ -814,15 +910,89 @@ def _import_approved_staging(source_ids: list[str] | None = None) -> dict[str, o
(target_id, approved_row["staging_id"]),
)
imported_updated += cursor.rowcount
duplicate_updated = 0
author_merged = 0
for duplicate_row in duplicate_rows:
target_id = duplicate_targets[duplicate_row["staging_id"]]
cursor.execute(
f"SELECT id, lyricist, composer FROM {TARGET_TABLE_NAME} WHERE id = %s AND deleted = '0' FOR UPDATE",
(target_id,),
)
target = cursor.fetchone()
if not target:
raise RuntimeError(f"重复目标已失效: target_id={target_id}")
merged_lyricist = _merge_author_field(target.get("lyricist"), duplicate_row.get("lyricist"))
merged_composer = _merge_author_field(target.get("composer"), duplicate_row.get("composer"))
changed = False
if merged_lyricist:
cursor.execute(
f"UPDATE {TARGET_TABLE_NAME} SET lyricist = %s WHERE id = %s",
(merged_lyricist, target_id),
)
changed = True
if merged_composer:
cursor.execute(
f"UPDATE {TARGET_TABLE_NAME} SET composer = %s WHERE id = %s",
(merged_composer, target_id),
)
changed = True
if changed:
author_merged += 1
cursor.execute(
f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET staging_status = 'skipped', imported_song_id = %s, error_message = NULL,
merge_authors = %s, review_version = review_version + 1
WHERE staging_id = %s AND staging_status = 'staged'
AND biz_review_status = 'rejected_duplicate'
""",
(target_id, 1 if changed else 0, duplicate_row["staging_id"]),
)
if cursor.rowcount != 1:
raise RuntimeError(f"重复样本终态更新失败: staging_id={duplicate_row['staging_id']}")
duplicate_updated += 1
conn.commit()
return {
"approved_count": len(approved_rows),
"inserted_count": inserted_count,
"imported_song_ids_updated": imported_updated,
"duplicate_count": len(duplicate_rows),
"duplicate_staging_updated": duplicate_updated,
"author_merged_count": author_merged,
"source_ids": approved_source_ids,
"duplicate_source_ids": [str(row["source_song_id"]) for row in duplicate_rows],
}
def _resolve_duplicate_target(cursor, row: dict[str, object]) -> dict[str, object] | None:
"""Resolve a human-confirmed duplicate to an existing, active target record."""
matched_id = str(row.get("matched_song_id") or "").strip()
if not matched_id:
return None
try:
numeric_id = int(matched_id)
except (TypeError, ValueError):
numeric_id = None
if numeric_id is not None:
cursor.execute(
f"SELECT id FROM {TARGET_TABLE_NAME} WHERE id = %s AND deleted = '0' FOR UPDATE",
(numeric_id,),
)
target = cursor.fetchone()
if target:
return target
cursor.execute(
f"""
SELECT id FROM {TARGET_TABLE_NAME}
WHERE source_table_name = %s AND source_song_id = %s AND deleted = '0'
FOR UPDATE
""",
(row.get("source_table_name"), matched_id),
)
return cursor.fetchone()
def _find_import_conflicts(cursor, staging_ids: list[object]) -> list[dict[str, object]]:
"""Return target-table key conflicts for the selected staging rows."""
if not staging_ids:
......@@ -1276,6 +1446,7 @@ class Handler(BaseHTTPRequestHandler):
staging_id = int(payload.get("staging_id") or 0)
decision = str(payload.get("decision", "")).strip()
note = str(payload.get("note", "")).strip()
candidate_id = str(payload.get("candidate_id", "")).strip()
client_token = str(payload.get("client_token") or "").strip()
expected_version = int(payload.get("expected_version", -1))
if not staging_id:
......@@ -1285,7 +1456,7 @@ class Handler(BaseHTTPRequestHandler):
if decision not in ("approved_import", "rejected_duplicate", "unsure", "pending", "deleted"):
raise ValueError(f"invalid decision: {decision}")
result = _update_staging_review(
staging_id, decision, note, reviewer, expected_version, client_token
staging_id, decision, note, reviewer, expected_version, client_token, candidate_id
)
_json_response(self, {"record": result})
except ReviewConflictError as exc:
......@@ -1310,11 +1481,19 @@ class Handler(BaseHTTPRequestHandler):
if source_ids
else None
)
approved = [
finalized = [
{"source_id": source_id, "review_decision": "import", "review_note": ""}
for source_id in result["source_ids"]
]
review_csv = _write_review_import_csv(approved)
finalized.extend(
{
"source_id": source_id,
"review_decision": "duplicate_merged",
"review_note": "确认重复,词/曲作者已增量合并到已有记录",
}
for source_id in result["duplicate_source_ids"]
)
review_csv = _write_review_import_csv(finalized)
_json_response(self, {"review_csv": str(review_csv), **result})
except ReviewAuthError as exc:
_error(self, str(exc), status=401)
......
......@@ -96,6 +96,45 @@ def test_stale_review_returns_conflict_without_overwrite():
assert not conn.committed
def test_rejected_duplicate_persists_the_selected_candidate():
snapshot = {
"staging_id": 42,
"matched_song_id": None,
"recalled_candidates": '[{"id":"existing-song"}]',
"biz_review_status": "rejected_duplicate",
"review_version": 8,
}
cursor = FakeCursor(update_rowcount=1, snapshot=snapshot)
conn = FakeConnection(cursor)
with patch.object(dashboard, "_target_conn", return_value=conn):
dashboard._update_staging_review(
42, "rejected_duplicate", "same", "alice", 7, "browser-token", "existing-song"
)
update_sql, params = cursor.executions[1]
assert "matched_song_id = %s" in update_sql
assert params[:3] == ["rejected_duplicate", "alice", "existing-song"]
assert conn.committed
def test_rejected_duplicate_rejects_a_candidate_outside_recall_results():
cursor = FakeCursor(
update_rowcount=1,
snapshot={"staging_id": 42, "matched_song_id": "candidate-a", "recalled_candidates": "[]"},
)
conn = FakeConnection(cursor)
with patch.object(dashboard, "_target_conn", return_value=conn):
with pytest.raises(ValueError, match="不在当前召回结果"):
dashboard._update_staging_review(
42, "rejected_duplicate", "", "alice", 7, "browser-token", "candidate-b"
)
assert conn.rolled_back
assert not any(sql.lstrip().startswith("UPDATE") for sql, _params in cursor.executions)
def test_pending_review_restores_a_deleted_staging_row():
snapshot = {
"staging_id": 42,
......@@ -185,7 +224,7 @@ class ImportCursor(FakeCursor):
def __init__(self) -> None:
super().__init__(update_rowcount=1, snapshot={})
self.result_sets = [
[{"staging_id": 42, "source_song_id": "song-1"}],
[{"staging_id": 42, "source_song_id": "song-1", "biz_review_status": "approved_import"}],
[],
[{"source_song_id": "song-1", "target_id": 9001}],
]
......@@ -217,7 +256,7 @@ def test_import_locks_only_unclaimed_staged_rows():
def test_import_conflict_rolls_back_and_exposes_the_staging_row():
cursor = ImportCursor()
cursor.result_sets = [
[{"staging_id": 42, "source_song_id": "song-1"}],
[{"staging_id": 42, "source_song_id": "song-1", "biz_review_status": "approved_import"}],
[{
"staging_id": 42,
"source_song_id": "song-1",
......@@ -241,6 +280,69 @@ def test_import_conflict_rolls_back_and_exposes_the_staging_row():
assert not any(sql.lstrip().startswith("INSERT") for sql, _params in cursor.executions)
class DuplicateImportCursor:
def __init__(self) -> None:
self.rowcount = 0
self.executions: list[tuple[str, object]] = []
self.current = None
def __enter__(self):
return self
def __exit__(self, *_args):
return False
def execute(self, sql, params=None):
self.executions.append((sql, params))
compact = " ".join(sql.split())
self.rowcount = 0
if "SELECT s.staging_id" in compact and "FOR UPDATE" in compact:
self.current = [{
"staging_id": 42,
"source_song_id": "new-song",
"source_table_name": "hk_song_platform",
"name": "same song",
"lyricist": "Alice/Bob",
"composer": "Composer B",
"biz_review_status": "rejected_duplicate",
"matched_song_id": "existing-song",
}]
elif "SELECT id FROM" in compact and "source_table_name = %s" in compact:
self.current = {"id": 9001}
elif "SELECT id, lyricist, composer" in compact:
self.current = {"id": 9001, "lyricist": "Alice", "composer": "Composer A"}
elif compact.startswith("UPDATE"):
self.current = None
self.rowcount = 1
else:
raise AssertionError(f"unexpected SQL: {compact}")
def fetchall(self):
return self.current or []
def fetchone(self):
return self.current
def test_confirmed_duplicate_merges_authors_without_inserting_a_new_song():
cursor = DuplicateImportCursor()
conn = FakeConnection(cursor)
with patch.object(dashboard, "_target_conn", return_value=conn):
result = dashboard._import_approved_staging()
sql_params = [(" ".join(sql.split()), params) for sql, params in cursor.executions]
assert not any(sql.startswith("INSERT") for sql, _params in sql_params)
assert any("SET lyricist = %s" in sql and params == ("Alice、Bob", 9001) for sql, params in sql_params)
assert any("SET composer = %s" in sql and params == ("Composer A、Composer B", 9001) for sql, params in sql_params)
assert any("SET staging_status = 'skipped'" in sql for sql, _params in sql_params)
assert result["duplicate_count"] == 1
assert result["author_merged_count"] == 1
assert result["inserted_count"] == 0
assert result["duplicate_source_ids"] == ["new-song"]
assert conn.committed
def test_browsing_a_list_item_can_auto_claim_without_reloading_groups():
html = Path("l2_review_dashboard.html").read_text(encoding="utf-8")
start = html.index("for (const btn of els.queryList.querySelectorAll('.query-item'))")
......@@ -274,7 +376,17 @@ def test_dashboard_distinguishes_reviewed_from_submitted():
assert row["staging_status"] == "imported"
assert '<option value="reviewed">已审核</option>' in html
assert '<option value="submitted">已提交</option>' in html
assert "hasFinalReview && query.staging_status !== 'imported'" in html
assert "hasFinalReview && query.staging_status === 'staged'" in html
def test_dashboard_has_a_dedicated_deleted_filter():
html = Path("l2_review_dashboard.html").read_text(encoding="utf-8")
assert '<option value="deleted">已删除</option>' in html
source = Path("serve_l2_dashboard.py").read_text(encoding="utf-8")
assert "elif decision == 'deleted':" in source
assert "s.biz_review_status = 'deleted'" in source
assert "s.staging_status = 'deleted'" in source
def test_review_access_code_issues_and_validates_edit_session():
......