Commit 6cce9eaa 6cce9eaa10ceeca9c2d7280f60f884e189f639f8 by 沈秋雨

feat(review): 实现多人协作审核及状态同步功能

- 新增审核人身份录入和本地缓存功能,支持切换审核人并生成客户端令牌
- 引入审核领取机制,支持防止多人同时修改同一条记录
- 实现审核状态的乐观锁版本控制,防止并发更新冲突
- 后端新增审核相关字段和索引,支持多审核人操作和审核状态管理
- 在审核列表和详情页中展示审核状态、领取状态和审核人信息
- 支持撤销审核、删除标注和入库操作,增加详细错误提示和冲突提示
- 定时同步审核状态,自动刷新审核列表和详情,保持界面数据一致
- 优化审核筛选条件,新增待审核、已审核、已提交过滤项
- 前端增加同步通知提示,展示审核冲突和同步错误信息
- 修改导出和导入逻辑,确保入库操作仅在 staging-db 环境进行
1 parent 7d33be0f
......@@ -534,6 +534,19 @@
min-width: 220px;
}
.sync-notice {
display: none;
padding: 8px 12px;
background: #fff4d6;
border-bottom: 1px solid #e8cf8a;
color: #725500;
font-size: 13px;
}
.sync-notice.show { display: block; }
#reviewerInput { width: 130px; }
@media (max-width: 1100px) {
main { grid-template-columns: 1fr; height: auto; }
aside { border-right: 0; border-bottom: 1px solid var(--line); max-height: 420px; }
......@@ -554,13 +567,14 @@
<select id="topKSelect"></select>
<label for="decisionFilter">类型</label>
<select id="decisionFilter">
<option value="pending" selected>待审核</option>
<option value="">全部</option>
<option value="l1_hit">L1 命中</option>
<option value="l2_duplicate">L2 重复</option>
<option value="l2_review">L2 审核</option>
<option value="conflict">L1/L2 冲突</option>
<option value="new">新歌</option>
<option value="skip">已跳过</option>
<option value="reviewed">已审核</option>
<option value="submitted">已提交</option>
</select>
<input id="searchInput" type="search" placeholder="搜索歌名 / ID">
<label for="pageSizeSelect">每页</label>
......@@ -572,9 +586,12 @@
</select>
<button id="exportBtn">导出标注</button>
<button id="importReviewedBtn">入库审核通过</button>
<label for="reviewerInput">审核人</label>
<input id="reviewerInput" maxlength="64" placeholder="姓名或工号">
<button id="reloadBtn" class="secondary">刷新</button>
</div>
</header>
<div id="syncNotice" class="sync-notice"></div>
<main>
<aside>
......@@ -623,10 +640,16 @@
pageSize: 50,
mode: 'retrieval',
topK: '',
decisionFilter: '',
decisionFilter: 'pending',
queryId: '',
candidateId: '',
lyricsCache: new Map()
lyricsCache: new Map(),
reviewer: '',
clientToken: '',
syncBusy: false,
detailGeneration: 0,
queryLoadGeneration: 0,
searchTimer: null
};
const els = {
......@@ -638,6 +661,7 @@
pageSizeSelect: document.getElementById('pageSizeSelect'),
exportBtn: document.getElementById('exportBtn'),
importReviewedBtn: document.getElementById('importReviewedBtn'),
reviewerInput: document.getElementById('reviewerInput'),
reloadBtn: document.getElementById('reloadBtn'),
prevPageBtn: document.getElementById('prevPageBtn'),
nextPageBtn: document.getElementById('nextPageBtn'),
......@@ -646,9 +670,32 @@
sampleCount: document.getElementById('sampleCount'),
queryList: document.getElementById('queryList'),
detail: document.getElementById('scrollDetail'),
stickyInfo: document.getElementById('stickyInfo')
stickyInfo: document.getElementById('stickyInfo'),
syncNotice: document.getElementById('syncNotice')
};
function newClientToken() {
if (window.crypto?.randomUUID) return window.crypto.randomUUID();
return `${Date.now()}-${Math.random().toString(16).slice(2)}`;
}
function initReviewerIdentity() {
state.reviewer = (localStorage.getItem('l2ReviewerName') || '').trim();
if (!state.reviewer) {
state.reviewer = (window.prompt('请输入审核人姓名或工号') || '').trim();
}
if (!state.reviewer) throw new Error('必须填写审核人姓名或工号后才能开始审核');
localStorage.setItem('l2ReviewerName', state.reviewer);
els.reviewerInput.value = state.reviewer;
state.clientToken = sessionStorage.getItem('l2ReviewClientToken') || newClientToken();
sessionStorage.setItem('l2ReviewClientToken', state.clientToken);
}
function showSyncNotice(message) {
els.syncNotice.textContent = message || '';
els.syncNotice.classList.toggle('show', Boolean(message));
}
function esc(value) {
return String(value ?? '').replace(/[&<>"']/g, ch => ({
'&': '&amp;', '<': '&lt;', '>': '&gt;', '"': '&quot;', "'": '&#39;'
......@@ -744,6 +791,83 @@
});
}
async function postJSON(url, body) {
const res = await fetch(url, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(body)
});
const payload = await res.json().catch(() => ({}));
if (!res.ok) {
const error = new Error(payload.error || `${res.status} ${res.statusText}`);
error.status = res.status;
error.payload = payload;
throw error;
}
return payload;
}
function applyReviewSnapshot(row, record) {
if (!row || !record) return;
row.review_status = record.biz_review_status ?? row.review_status;
row.review_note = record.biz_review_note ?? row.review_note;
row.reviewed_by = record.reviewed_by ?? '';
row.reviewed_at = record.reviewed_at ?? '';
row.review_version = String(record.review_version ?? row.review_version ?? 0);
row.review_claimed_by = record.review_claimed_by ?? '';
row.review_claim_expires_at = record.review_claim_expires_at ?? '';
row.staging_status = record.staging_status ?? row.staging_status;
row.claimed_by_me = Boolean(record.claimed_by_me) || row.review_claimed_by === state.reviewer;
}
function applySnapshotToCurrentGroup(record) {
const group = state.groups.find(item => item.id === state.queryId);
if (!group || !record) return;
group.review_version = Number(record.review_version || 0);
group.review_claimed_by = record.review_claimed_by || '';
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';
}
function currentQuery() {
return state.queryRows[0] || null;
}
async function claimCurrentReview({ allowReviewed = false, allowNonReview = false } = {}) {
const query = currentQuery();
if (!query || state.run?.type !== 'staging' || !query.staging_id) return true;
const status = query.review_status || '';
const isPendingReview = rowDecision(query) === 'review' && ['pending', 'unsure', ''].includes(status);
const isOwnReviewedRecord = allowReviewed && query.reviewed_by === state.reviewer;
const isNonReviewRecord = allowNonReview && status === 'not_required';
const reviewable = isPendingReview || isOwnReviewedRecord || isNonReviewRecord;
// new / merge / skip 等记录只浏览,不应触发领取请求。
if (!reviewable) return true;
try {
const payload = await postJSON('/api/review-claim', {
// bigint 必须按字符串传输,不能转成会丢精度的 JavaScript Number。
staging_id: String(query.staging_id),
reviewer: state.reviewer,
client_token: state.clientToken
});
for (const row of state.queryRows) applyReviewSnapshot(row, payload.record);
applySnapshotToCurrentGroup(payload.record);
showSyncNotice('');
return true;
} catch (error) {
if (error.status === 409) {
for (const row of state.queryRows) applyReviewSnapshot(row, error.payload?.current_record);
applySnapshotToCurrentGroup(error.payload?.current_record);
const owner = error.payload?.current_record?.review_claimed_by;
showSyncNotice(owner ? `该记录已由 ${owner} 领取;你仍可浏览,但不能覆盖对方的审核结果。` : error.message);
return false;
}
throw error;
}
}
function dataPath() {
if (!state.run) return '';
return state.run.retrieval;
......@@ -757,13 +881,14 @@
const row = state.summary.find(item => String(item.top_k) === String(state.topK)) || {};
els.metrics.innerHTML = [
metric('总数', row.total_count || '-'),
metric('已入库 new', row.new_count || '-'),
metric('耗时(s)', row.elapsed_seconds || '-'),
metric('吞吐(条/s)', row.throughput_per_second || '-'),
metric('平均召回', row.avg_recalled_candidates || '-'),
metric('命中数', row.hit_count || '0'),
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.deleted_count || '0'),
metric('new', row.new_count || '-'),
metric('merge', row.duplicate_count || '0'),
metric('review', row.review_count || '0'),
metric('skip', row.skip_count || '0')
].join('');
}
......@@ -786,7 +911,12 @@
const active = item.id === state.queryId ? ' active' : '';
const badges = [
item.has_conflict ? '<span class="pill conflict">L1/L2 冲突</span>' : '',
item.has_l1_hint ? '<span class="pill l1">L1 命中候选</span>' : ''
item.has_l1_hint ? '<span class="pill l1">L1 命中候选</span>' : '',
item.submitted ? '<span class="pill new">已提交</span>' : '',
item.reviewed && !item.submitted ? '<span class="pill">已审核</span>' : '',
item.claimed_by_me ? '<span class="pill new">我正在审核</span>' : '',
item.review_claimed_by && !item.claimed_by_me
? `<span class="pill">${esc(item.review_claimed_by)} 正在审核</span>` : ''
].join(' ');
return `<button class="query-item${active}${hitClass}${conflict}" data-query-id="${esc(item.id)}">
<div class="query-title">${esc(item.query_name || '(无歌名)')}</div>
......@@ -799,8 +929,12 @@
btn.addEventListener('click', async () => {
state.queryId = btn.dataset.queryId;
state.candidateId = '';
showSyncNotice('');
renderList();
await loadQueryRows();
renderAll();
await claimCurrentReview();
renderList();
renderDetail();
});
}
els.prevPageBtn.disabled = state.page <= 1;
......@@ -858,21 +992,19 @@
</div>`;
}
async function updateStagingReview(sourceId, decision, note) {
const res = await fetch('/api/staging-review', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ source_id: sourceId, decision, note })
async function updateStagingReview(query, decision, note) {
return postJSON('/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
});
if (!res.ok) {
const err = await res.json().catch(() => ({}));
throw new Error(err.error || `${res.status}`);
}
return res.json();
}
const _DB_REVIEW_LABELS = {
approved_import: '确认不重复(入库)',
approved_import: '确认不重复(入库)',
rejected_duplicate: '确认重复',
unsure: '待确认',
deleted: '已删除',
......@@ -943,6 +1075,7 @@
}
async function renderDetail() {
const generation = ++state.detailGeneration;
const rows = selectedQueryRows();
if (!rows.length) {
els.stickyInfo.innerHTML = '';
......@@ -951,9 +1084,12 @@
}
const query = rows[0];
const candidate = selectedCandidate(rows);
const queryText = await readText(query.query_lyrics_path);
const candidatePath = candidate?.candidate_lyrics_path;
const candidateText = await readText(candidatePath);
const [queryText, candidateText] = await Promise.all([
readText(query.query_lyrics_path),
readText(candidatePath)
]);
if (generation !== state.detailGeneration) return;
const hitDecision = candidate?.candidate_decision;
const candidateName = candidate?.candidate_name;
const candidateLyricist = candidate?.candidate_lyricist;
......@@ -964,6 +1100,12 @@
const l1 = candidate ? l1Match(candidate) : false;
const queryAudioUrl = query.audio_url || '';
const candidateAudioUrl = candidate?.candidate_audio_url || '';
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' &&
query.reviewed_by === state.reviewer;
const canDeleteRecord = query.staging_status !== 'imported';
// Sticky area: 新入库歌词 + 召回候选(各含音频播放器)
els.stickyInfo.innerHTML = `
......@@ -1007,22 +1149,24 @@
<div class="section">
<div class="section-head">
<h2 class="section-title">人工标注</h2>
${query.review_status && query.review_status !== 'pending' && query.review_status !== 'not_required'
? `<span class="pill ${query.review_status === 'approved_import' ? 'new' : query.review_status === 'rejected_duplicate' ? 'duplicate' : query.review_status === 'deleted' ? 'duplicate' : ''}">${query.review_status === 'deleted' ? '已删除' : '已审核'}</span>`
${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>`
: ''}
</div>
<div class="section-body">
${query.review_status && query.review_status !== 'pending' && query.review_status !== 'not_required' ? `
${hasFinalReview ? `
<div class="review-bar" style="align-items:center">
<span style="flex:1">
<b>${esc(_DB_REVIEW_LABELS[query.review_status] || query.review_status)}</b>
<b>${esc(query.review_status === 'approved_import' && query.staging_status === 'imported'
? '确认不重复(已提交)'
: (_DB_REVIEW_LABELS[query.review_status] || query.review_status))}</b>
${query.reviewed_by ? ` · 审核人 ${esc(query.reviewed_by)}` : ''}
${query.reviewed_at ? ` · ${esc(query.reviewed_at)}` : ''}
${query.review_note ? ` · 备注: ${esc(query.review_note)}` : ''}
</span>
<button id="undoReviewBtn" class="secondary">撤销审核</button>
${canUndoReview ? '<button id="undoReviewBtn" class="secondary">撤销审核</button>' : ''}
</div>
` : `
` : needsManualReview ? `
<div class="review-bar">
${['duplicate', 'not_duplicate', 'unsure'].map(value => `
<button class="review-choice ${selectedReview.final_decision === value ? 'active' : ''}" data-review="${value}">
......@@ -1031,6 +1175,11 @@
<input id="reviewNote" value="${esc(selectedReview.note || '')}" placeholder="人工备注">
<button id="deleteReviewBtn" class="secondary" style="color:#c0392b;border-color:#c0392b">删除</button>
</div>
` : `
<div class="review-bar" style="align-items:center">
<span style="flex:1">该记录的自动判定为 <b>${esc(rowDecision(query) || '-')}</b>,当前仅供浏览,无需领取。</span>
${canDeleteRecord ? '<button id="deleteReviewBtn" class="secondary" style="color:#c0392b;border-color:#c0392b">删除</button>' : ''}
</div>
`}
</div>
</div>
......@@ -1069,14 +1218,24 @@
setReview(candidate, { final_decision: localDecision, note });
btn.disabled = true;
try {
await updateStagingReview(query.query_source_id, dbDecision, note);
// 重新加载当前样本数据以刷新 DB 状态
await loadQueryRows(query.query_source_id);
if (!query.claimed_by_me && !await claimCurrentReview()) {
renderDetail();
return;
}
const payload = await updateStagingReview(query, dbDecision, note);
for (const row of state.queryRows) applyReviewSnapshot(row, payload.record);
await reloadSummary();
// 重新加载列表并跳转到下一条
await advanceToNext();
} catch (e) {
console.error('DB同步失败:', e);
alert(`标注已保存到本地,但数据库同步失败: ${e.message}`);
if (e.status === 409) {
showSyncNotice(`提交冲突:${e.message}。备注草稿已保留,请刷新后确认。`);
} else {
showSyncNotice(`数据库同步失败:${e.message}。备注草稿已保留,可稍后重试。`);
}
renderDetail();
}
});
}
const deleteBtn = document.getElementById('deleteReviewBtn');
......@@ -1085,13 +1244,18 @@
if (!confirm('确认删除此样本?删除后将不再导入。')) return;
deleteBtn.disabled = true;
try {
await updateStagingReview(query.query_source_id, 'deleted',
if (!query.claimed_by_me && !await claimCurrentReview({ allowNonReview: true })) {
renderDetail();
return;
}
await updateStagingReview(query, 'deleted',
document.getElementById('reviewNote')?.value || '');
await loadQueryRows(query.query_source_id);
await reloadSummary();
await advanceToNext();
} catch (e) {
alert(`删除失败: ${e.message}`);
}
renderDetail();
}
});
}
const undoBtn = document.getElementById('undoReviewBtn');
......@@ -1100,9 +1264,14 @@
if (!confirm('确认撤销此样本的人工审核结果?')) return;
undoBtn.disabled = true;
try {
await updateStagingReview(query.query_source_id, 'pending', '');
if (!query.claimed_by_me && !await claimCurrentReview({ allowReviewed: true })) {
renderDetail();
return;
}
await updateStagingReview(query, 'pending', '');
if (candidate) setReview(candidate, { final_decision: undefined, note: '' });
await loadQueryRows(query.query_source_id);
await loadQueryRows();
await reloadSummary();
} catch (e) {
alert(`撤销失败: ${e.message}`);
}
......@@ -1122,6 +1291,12 @@
const rows = [];
for (const row of state.queryRows) {
const review = reviews[reviewKey(row)] || {};
const databaseDecision = {
approved_import: 'not_duplicate',
rejected_duplicate: 'duplicate',
unsure: 'unsure',
deleted: 'deleted'
}[row.review_status] || '';
rows.push({
run_id: state.run?.id || '',
top_k: state.topK,
......@@ -1136,9 +1311,9 @@
l2_decision: rowDecision(row),
l1_metadata_match: l1Match(row) ? '1' : '0',
l1_l2_conflict: isConflict(row) ? '1' : '0',
final_decision: review.final_decision || '',
note: review.note || '',
updated_at: review.updated_at || ''
final_decision: state.run?.type === 'staging' ? databaseDecision : (review.final_decision || ''),
note: state.run?.type === 'staging' ? (row.review_note || '') : (review.note || ''),
updated_at: state.run?.type === 'staging' ? (row.reviewed_at || '') : (review.updated_at || '')
});
}
const fields = Object.keys(rows[0] || {
......@@ -1161,41 +1336,19 @@
async function importReviewed() {
if (!state.run) return;
const reviews = loadReviews();
const prefix = `${state.run.id}::${state.topK || ''}::`;
const rows = [];
for (const [key, review] of Object.entries(reviews)) {
if (!key.startsWith(prefix)) continue;
if (review.final_decision !== 'not_duplicate') continue;
const parts = key.split('::');
const sourceId = parts[2] || '';
if (!sourceId) continue;
rows.push({
source_id: sourceId,
review_decision: 'import',
review_note: review.note || ''
});
}
if (!rows.length) {
// 诊断信息:显示当前结果下实际找到的标注情况
const allForRun = Object.entries(reviews).filter(([k]) => k.startsWith(prefix));
const decisionSummary = allForRun.map(([, v]) => v.final_decision || '未标注').join(', ') || '无';
alert(`当前结果中没有标注为“确认不重复”的 review 样本。\n\n` +
`当前结果下找到 ${allForRun.length} 条标注,状态: ${decisionSummary}`);
if (state.run.type !== 'staging') {
alert('历史报表仅供查看,请在 staging-db 中执行入库。');
return;
}
if (!confirm(`确认入库 ${rows.length} 条审核通过样本?`)) return;
if (!confirm('确认入库数据库中所有“确认不重复(待入库)”样本?')) return;
els.importReviewedBtn.disabled = true;
els.importReviewedBtn.textContent = '入库中...';
try {
const res = await fetch('/api/import-reviewed', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ run_id: state.run.id, rows })
});
const payload = await res.json();
if (!res.ok) throw new Error(payload.error || payload.output || `${res.status} ${res.statusText}`);
alert(`入库完成。结果文件:${payload.review_csv}`);
const payload = await postJSON('/api/import-reviewed', { reviewer: state.reviewer });
alert(`入库完成:${payload.inserted_count} 条。结果文件:${payload.review_csv}`);
await loadGroups();
await reloadSummary();
renderAll();
} catch (err) {
alert(`入库失败:${err.message}`);
} finally {
......@@ -1221,7 +1374,8 @@
decision: state.decisionFilter,
q: els.searchInput.value.trim(),
page: String(state.page),
page_size: String(state.pageSize)
page_size: String(state.pageSize),
client_token: state.clientToken
});
const data = await getJSON(`/api/groups?${params.toString()}`);
state.groups = data.groups || [];
......@@ -1232,20 +1386,103 @@
state.candidateId = '';
}
await loadQueryRows();
// 仅 pending/unsure 的 review 记录会真正发起领取;其他分类只浏览。
await claimCurrentReview();
}
async function loadQueryRows() {
state.queryRows = [];
const generation = ++state.queryLoadGeneration;
const requestedQueryId = state.queryId;
if (!state.queryId || !state.run || !state.topK) return;
const params = new URLSearchParams({
path: dataPath(),
top_k: state.topK,
query_id: state.queryId
});
state.queryRows = [];
const data = await getJSON(`/api/query?${params.toString()}`);
if (generation !== state.queryLoadGeneration || requestedQueryId !== state.queryId) return;
state.queryRows = data.rows || [];
}
async function reloadSummary() {
if (!state.run) return;
const summary = await getJSON(`/api/csv?path=${encodeURIComponent(state.run.summary)}`);
const nextSummary = summary.rows || [];
if (JSON.stringify(nextSummary) !== JSON.stringify(state.summary)) {
state.summary = nextSummary;
renderMetrics();
}
}
async function advanceToNext() {
// 重新加载列表(应用当前筛选,已审核的会被过滤掉)
await loadGroups();
renderAll();
}
async function syncVisibleReviewStatuses() {
if (state.syncBusy || state.run?.type !== 'staging' || !state.groups.length) return;
const ids = state.groups.map(group => group.staging_id).filter(Boolean);
if (!ids.length) return;
state.syncBusy = true;
try {
const params = new URLSearchParams({
staging_ids: ids.join(','),
client_token: state.clientToken
});
const payload = await getJSON(`/api/review-statuses?${params.toString()}`);
const statuses = new Map((payload.rows || []).map(row => [String(row.staging_id), row]));
let listChanged = false;
for (const group of state.groups) {
const remote = statuses.get(String(group.staging_id));
if (!remote) continue;
if (Number(remote.review_version || 0) !== Number(group.review_version || 0) ||
Boolean(remote.claimed_by_me) !== Boolean(group.claimed_by_me)) {
listChanged = true;
}
group.review_version = Number(remote.review_version || 0);
group.review_claimed_by = remote.review_claimed_by || '';
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';
}
const query = currentQuery();
const selectedRemote = query ? statuses.get(String(query.staging_id)) : null;
if (selectedRemote && Number(selectedRemote.review_version || 0) !== Number(query.review_version || 0)) {
if (!selectedRemote.claimed_by_me) {
showSyncNotice(`该记录已由 ${selectedRemote.reviewed_by || selectedRemote.review_claimed_by || '其他审核人'} 更新;当前输入未被覆盖,请刷新确认。`);
}
}
if (listChanged) renderList();
} catch (error) {
console.error('同步审核状态失败:', error);
} finally {
state.syncBusy = false;
}
}
async function heartbeatCurrentReview() {
const query = currentQuery();
if (state.run?.type !== 'staging' || !query?.staging_id || !query.review_claimed_by) return;
const heartbeatStagingId = String(query.staging_id);
try {
const payload = await postJSON('/api/review-heartbeat', {
staging_id: String(query.staging_id),
client_token: state.clientToken
});
if (String(currentQuery()?.staging_id || '') === heartbeatStagingId) {
for (const row of state.queryRows) applyReviewSnapshot(row, payload.record);
}
} catch (error) {
if (error.status === 409 && String(currentQuery()?.staging_id || '') === heartbeatStagingId) {
showSyncNotice(error.message);
}
else console.error('审核领取续期失败:', error);
}
}
async function loadRun(runId) {
state.run = state.runs.find(run => run.id === runId) || state.runs[0];
if (!state.run) return;
......@@ -1280,6 +1517,24 @@
}
els.runSelect.addEventListener('change', () => loadRun(els.runSelect.value));
els.reviewerInput.addEventListener('change', async () => {
const nextReviewer = els.reviewerInput.value.trim();
if (!nextReviewer) {
els.reviewerInput.value = state.reviewer;
return;
}
if (currentQuery()?.review_claimed_by === state.reviewer) {
alert('请先完成当前审核,或等待领取过期后再切换审核人。');
els.reviewerInput.value = state.reviewer;
return;
}
state.reviewer = nextReviewer;
state.clientToken = newClientToken();
localStorage.setItem('l2ReviewerName', state.reviewer);
sessionStorage.setItem('l2ReviewClientToken', state.clientToken);
await loadGroups();
renderAll();
});
els.topKSelect.addEventListener('change', async () => {
state.topK = els.topKSelect.value;
state.page = 1;
......@@ -1288,12 +1543,15 @@
await loadGroups();
renderAll();
});
els.searchInput.addEventListener('input', async () => {
els.searchInput.addEventListener('input', () => {
window.clearTimeout(state.searchTimer);
state.searchTimer = window.setTimeout(async () => {
state.page = 1;
state.queryId = '';
state.candidateId = '';
await loadGroups();
renderAll();
}, 350);
});
els.decisionFilter.addEventListener('change', async () => {
state.decisionFilter = els.decisionFilter.value;
......@@ -1331,6 +1589,21 @@
renderAll();
});
try {
initReviewerIdentity();
} catch (err) {
els.detail.innerHTML = `<div class="empty">${esc(err.message)}</div>`;
throw err;
}
window.setInterval(() => {
if (!document.hidden) syncVisibleReviewStatuses();
}, 15000);
window.setInterval(() => {
if (!document.hidden) reloadSummary().catch(console.error);
}, 60000);
window.setInterval(() => {
if (!document.hidden) heartbeatCurrentReview();
}, 60000);
loadRuns().catch(err => {
els.detail.innerHTML = `<div class="empty">${esc(err.message)}</div>`;
......
......@@ -26,6 +26,7 @@ DASHBOARD = ROOT / "l2_review_dashboard.html"
REPORT_DIR.mkdir(parents=True, exist_ok=True)
GROUP_INDEX_CACHE: dict[tuple[str, int, int, str], list[dict[str, object]]] = {}
load_dotenv(ROOT / ".env")
REVIEW_CLAIM_TTL_SECONDS = max(60, int(os.getenv("REVIEW_CLAIM_TTL_SECONDS", "600")))
TARGET_DB_CONFIG = {
"host": os.getenv("TARGET_DB_HOST"),
......@@ -52,6 +53,33 @@ TARGET_COLUMNS = [
"commit_desc", "price", "source_table_name", "source_song_id",
"lyric_archive_element_id", "melody_archive_element_id", "audio_fingerprint",
]
REVIEW_SCHEMA_COLUMNS = {
"review_claimed_by": "varchar(64) DEFAULT NULL COMMENT '当前领取审核人'",
"review_claim_token": "varchar(64) DEFAULT NULL COMMENT '领取会话凭证'",
"review_claimed_at": "datetime DEFAULT NULL COMMENT '领取时间'",
"review_claim_expires_at": "datetime DEFAULT NULL COMMENT '领取过期时间'",
"review_version": "bigint(20) NOT NULL DEFAULT '0' COMMENT '审核乐观锁版本'",
}
class ReviewConflictError(Exception):
"""The staging row changed or is owned by another review session."""
def __init__(self, message: str, current_record: dict[str, object] | None = None) -> None:
super().__init__(message)
self.current_record = current_record or {}
def _conflict_response(handler: BaseHTTPRequestHandler, exc: ReviewConflictError) -> None:
_json_response(
handler,
{
"error": str(exc),
"code": "review_conflict",
"current_record": exc.current_record,
},
status=409,
)
def _json_response(handler: BaseHTTPRequestHandler, payload: object, status: int = 200) -> None:
......@@ -59,6 +87,7 @@ def _json_response(handler: BaseHTTPRequestHandler, payload: object, status: int
try:
handler.send_response(status)
handler.send_header("Content-Type", "application/json; charset=utf-8")
handler.send_header("Cache-Control", "no-store")
handler.send_header("Content-Length", str(len(body)))
handler.end_headers()
handler.wfile.write(body)
......@@ -97,6 +126,38 @@ def _target_conn():
return pymysql.connect(**TARGET_DB_CONFIG)
def _ensure_review_schema(apply_migration: bool = False) -> None:
"""Validate or add the columns required by collaborative reviewing."""
with _target_conn() as conn, conn.cursor() as cursor:
cursor.execute(f"SHOW COLUMNS FROM {TARGET_TABLE_NAME_TMP}")
existing_columns = {str(row["Field"]) for row in cursor.fetchall()}
missing = [name for name in REVIEW_SCHEMA_COLUMNS if name not in existing_columns]
if missing and not apply_migration:
names = ", ".join(missing)
raise RuntimeError(
f"暂存表缺少多人审核字段: {names}。"
"请先使用 --migrate-review-schema 启动一次。"
)
for name in missing:
cursor.execute(
f"ALTER TABLE {TARGET_TABLE_NAME_TMP} "
f"ADD COLUMN `{name}` {REVIEW_SCHEMA_COLUMNS[name]}"
)
cursor.execute(f"SHOW INDEX FROM {TARGET_TABLE_NAME_TMP} WHERE Key_name = %s", ("idx_staging_review_claim",))
if not cursor.fetchone():
if not apply_migration:
raise RuntimeError(
"暂存表缺少索引 idx_staging_review_claim。"
"请先使用 --migrate-review-schema 启动一次。"
)
cursor.execute(
f"ALTER TABLE {TARGET_TABLE_NAME_TMP} "
"ADD KEY `idx_staging_review_claim` (`biz_review_status`, `review_claim_expires_at`)"
)
if apply_migration:
conn.commit()
def _staging_summary() -> list[dict[str, str]]:
sql = f"""
SELECT
......@@ -104,9 +165,15 @@ def _staging_summary() -> list[dict[str, str]]:
SUM(CASE WHEN dedup_action = 'new' THEN 1 ELSE 0 END) AS new_count,
SUM(CASE WHEN dedup_action = 'merge' THEN 1 ELSE 0 END) AS merge_count,
SUM(CASE WHEN dedup_action = 'review' THEN 1 ELSE 0 END) AS review_count,
SUM(CASE WHEN dedup_action = 'skip' THEN 1 ELSE 0 END) AS skip_count
SUM(CASE WHEN dedup_action = 'skip' THEN 1 ELSE 0 END) AS skip_count,
SUM(CASE WHEN dedup_action = 'review' AND (biz_review_status IS NULL OR biz_review_status IN ('pending', 'unsure')) THEN 1 ELSE 0 END) AS pending_review_count,
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
FROM (
SELECT source_song_id, dedup_action
SELECT source_song_id, dedup_action, biz_review_status, staging_status
FROM {TARGET_TABLE_NAME_TMP}
WHERE (source_song_id, staging_id) IN (
SELECT source_song_id, MAX(staging_id)
......@@ -129,6 +196,12 @@ def _staging_summary() -> list[dict[str, str]]:
"skip_count": str(row.get("skip_count") or 0),
"new_count": str(row.get("new_count") or 0),
"total_count": str(row.get("total_count") or 0),
"approved_count": str(row.get("approved_count") or 0),
"rejected_count": str(row.get("rejected_count") or 0),
"deleted_count": str(row.get("deleted_count") or 0),
"reviewed_count": str(row.get("reviewed_count") or 0),
"submitted_count": str(row.get("submitted_count") or 0),
"pending_review_count": str(row.get("pending_review_count") or 0),
}]
......@@ -161,52 +234,83 @@ def _staging_dashboard_row(row: dict[str, object]) -> dict[str, str]:
"review_note": str(row.get("biz_review_note") or ""),
"reviewed_by": str(row.get("reviewed_by") or ""),
"reviewed_at": str(row.get("reviewed_at") or ""),
"review_version": str(row.get("review_version") or 0),
"review_claimed_by": str(row.get("review_claimed_by") or ""),
"review_claim_expires_at": str(row.get("review_claim_expires_at") or ""),
"staging_status": str(row.get("staging_status") or ""),
"l1_metadata_match": "1" if row.get("l1_matched_id") else "0",
"l1_l2_conflict": "0",
"audio_url": str(row.get("audio_url") or ""),
}
def _staging_groups_response(mode: str, term: str, page: int, page_size: int, decision: str = '') -> dict[str, object]:
def _staging_groups_response(
mode: str,
term: str,
page: int,
page_size: int,
decision: str = "",
client_token: str = "",
) -> dict[str, object]:
where_clauses = []
params: list[object] = []
if term:
where_clauses.append("(CAST(source_song_id AS CHAR) LIKE %s OR name LIKE %s OR lyricist LIKE %s OR composer LIKE %s)")
where_clauses.append("(CAST(s.source_song_id AS CHAR) LIKE %s OR s.name LIKE %s OR s.lyricist LIKE %s OR s.composer LIKE %s)")
like = f"%{term}%"
params.extend([like, like, like, like])
# hits 模式下只展示有命中的记录(merge/review)
if mode == "hits":
where_clauses.append("dedup_action IN ('merge', 'review')")
where_clauses.append("s.dedup_action IN ('merge', 'review')")
# 决策筛选:l1_hit / l2_hit / conflict / new / skip
if decision == 'l1_hit':
where_clauses.append("l1_matched_id IS NOT NULL")
where_clauses.append("s.l1_matched_id IS NOT NULL")
elif decision == 'l2_duplicate':
where_clauses.append("dedup_action = 'merge'")
where_clauses.append("s.dedup_action = 'merge'")
elif decision == 'l2_review':
where_clauses.append("dedup_action = 'review'")
where_clauses.append("s.dedup_action = 'review'")
elif decision == 'conflict':
where_clauses.append("(l1_matched_id IS NOT NULL AND dedup_action = 'new') OR (l1_matched_id IS NULL AND dedup_action IN ('merge', 'review'))")
where_clauses.append("(s.l1_matched_id IS NOT NULL AND s.dedup_action = 'new') OR (s.l1_matched_id IS NULL AND s.dedup_action IN ('merge', 'review'))")
elif decision == 'new':
where_clauses.append("dedup_action = 'new' AND l1_matched_id IS NULL")
where_clauses.append("s.dedup_action = 'new' AND s.l1_matched_id IS NULL")
elif decision == 'skip':
where_clauses.append("dedup_action = 'skip'")
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'")
elif decision == 'submitted':
where_clauses.append("s.biz_review_status IN ('approved_import', 'rejected_duplicate', 'deleted')")
where_clauses.append("s.staging_status = 'imported'")
elif decision == 'pending':
where_clauses.append("(s.biz_review_status IS NULL OR s.biz_review_status IN ('pending', 'unsure'))")
where_clauses.append("s.dedup_action = 'review'")
where_clauses.append("s.staging_status <> 'imported'")
where_clauses.append(
"(s.review_claim_expires_at IS NULL OR s.review_claim_expires_at <= NOW() OR s.review_claim_token = %s)"
)
params.append(client_token)
# 兼容旧的 dedup_action 筛选
elif decision and decision in ('merge', 'review'):
where_clauses.append("dedup_action = %s")
where_clauses.append("s.dedup_action = %s")
params.append(decision)
where = "WHERE " + " AND ".join(where_clauses) if where_clauses else ""
# 同一首歌可能因多次导入批次产生多行,按 source_song_id 去重取最新一行
count_sql = f"SELECT COUNT(DISTINCT source_song_id) AS n FROM {TARGET_TABLE_NAME_TMP} {where}"
sql = f"""
SELECT s.staging_id, s.source_song_id, s.name, s.lyricist, s.composer, s.dedup_action,
s.biz_review_status, s.l1_matched_id
latest_join = f"""
FROM {TARGET_TABLE_NAME_TMP} s
INNER JOIN (
SELECT source_song_id, MAX(staging_id) AS max_staging_id
FROM {TARGET_TABLE_NAME_TMP}
{where}
GROUP BY source_song_id
) latest ON s.staging_id = latest.max_staging_id
"""
count_sql = f"SELECT COUNT(*) AS n {latest_join} {where}"
sql = f"""
SELECT s.staging_id, s.source_song_id, s.name, s.lyricist, s.composer, s.dedup_action,
s.biz_review_status, s.staging_status, s.l1_matched_id, s.review_version,
CASE WHEN s.review_claim_expires_at > NOW() THEN s.review_claimed_by ELSE NULL END AS review_claimed_by,
s.review_claim_expires_at,
CASE WHEN s.review_claim_token = %s AND s.review_claim_expires_at > NOW() THEN 1 ELSE 0 END AS claimed_by_me
{latest_join}
{where}
ORDER BY
CASE s.dedup_action WHEN 'review' THEN 0 WHEN 'merge' THEN 1 ELSE 2 END,
s.staging_create_time DESC
......@@ -215,7 +319,7 @@ def _staging_groups_response(mode: str, term: str, page: int, page_size: int, de
with _target_conn() as conn, conn.cursor() as cursor:
cursor.execute(count_sql, params)
total = (cursor.fetchone() or {}).get("n", 0)
cursor.execute(sql, [*params, page_size, max(0, (page - 1) * page_size)])
cursor.execute(sql, [client_token, *params, page_size, max(0, (page - 1) * page_size)])
raw_rows = cursor.fetchall()
# Python 层面再去重,防止 SQL 层面因数据库版本或数据异常未能完全去重
seen_source_ids: set[str] = set()
......@@ -237,6 +341,16 @@ def _staging_groups_response(mode: str, term: str, page: int, page_size: int, de
"has_conflict": False,
"has_l1_hint": bool(row.get("l1_matched_id")),
"decision": row.get("dedup_action") or "",
"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"
),
"staging_id": str(row.get("staging_id") or ""),
"review_version": int(row.get("review_version") or 0),
"review_claimed_by": str(row.get("review_claimed_by") or ""),
"review_claim_expires_at": str(row.get("review_claim_expires_at") or ""),
"claimed_by_me": bool(row.get("claimed_by_me")),
}
for row in rows
]
......@@ -371,107 +485,241 @@ def _staging_query_rows(query_id: str) -> dict[str, object]:
return {"rows": rows}
def _update_staging_review(source_song_id: str, decision: str, note: str, reviewer: str = "dashboard") -> dict[str, object]:
"""更新暂存表的人工审核状态。
def _review_snapshot(cursor, staging_id: int) -> dict[str, object]:
cursor.execute(
f"""
SELECT staging_id, source_song_id, dedup_action, biz_review_status, biz_review_note,
reviewed_by, reviewed_at, review_version, review_claimed_by,
review_claim_expires_at, staging_status
FROM {TARGET_TABLE_NAME_TMP}
WHERE staging_id = %s
""",
(staging_id,),
)
return cursor.fetchone() or {}
decision: 'approved_import' | 'rejected_duplicate' | 'unsure' | 'pending' (撤销) | 'deleted'
"""
if decision == "pending":
# 撤销审核
sql = f"""
def _claim_staging_review(staging_id: int, reviewer: str, client_token: str) -> dict[str, object]:
if not reviewer or not client_token:
raise ValueError("reviewer and client_token are required")
with _target_conn() as conn, conn.cursor() as cursor:
cursor.execute(
f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET biz_review_status = 'pending',
reviewed_by = NULL,
reviewed_at = NULL,
biz_review_note = NULLIF(%s, '')
WHERE source_song_id = %s AND dedup_action = 'review'
SET review_claimed_by = %s,
review_claim_token = %s,
review_claimed_at = NOW(),
review_claim_expires_at = DATE_ADD(NOW(), INTERVAL %s SECOND),
review_version = review_version + 1
WHERE staging_id = %s
AND staging_status <> 'imported'
AND (
biz_review_status IN ('pending', 'unsure')
OR biz_review_status = 'not_required'
OR (reviewed_by = %s AND biz_review_status IN ('approved_import', 'rejected_duplicate', 'deleted'))
)
AND (
review_claim_expires_at IS NULL
OR review_claim_expires_at <= NOW()
OR review_claim_token = %s
)
""",
(reviewer, client_token, REVIEW_CLAIM_TTL_SECONDS, staging_id, reviewer, client_token),
)
if cursor.rowcount != 1:
current = _review_snapshot(cursor, staging_id)
conn.rollback()
if not current:
message = f"找不到 staging_id={staging_id},请刷新页面后重试"
elif current.get("staging_status") == "imported":
message = "该记录已经入库,只能浏览,不能再领取审核"
elif current.get("review_claimed_by"):
message = f"该记录已由 {current['review_claimed_by']} 领取"
else:
message = f"该记录当前状态为 {current.get('biz_review_status') or '-'},不能领取"
raise ReviewConflictError(message, current)
current = _review_snapshot(cursor, staging_id)
conn.commit()
return current
def _heartbeat_staging_review(staging_id: int, client_token: str) -> dict[str, object]:
with _target_conn() as conn, conn.cursor() as cursor:
cursor.execute(
f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET review_claim_expires_at = DATE_ADD(NOW(), INTERVAL %s SECOND)
WHERE staging_id = %s
AND review_claim_token = %s
AND review_claim_expires_at > NOW()
""",
(REVIEW_CLAIM_TTL_SECONDS, staging_id, client_token),
)
if cursor.rowcount != 1:
current = _review_snapshot(cursor, staging_id)
conn.rollback()
raise ReviewConflictError("审核领取已过期或已转给其他审核人", current)
current = _review_snapshot(cursor, staging_id)
conn.commit()
return current
def _review_statuses(staging_ids: list[int], client_token: str = "") -> list[dict[str, object]]:
if not staging_ids:
return []
placeholders = ",".join(["%s"] * len(staging_ids))
with _target_conn() as conn, conn.cursor() as cursor:
cursor.execute(
f"""
SELECT staging_id, source_song_id, biz_review_status, staging_status, reviewed_by, reviewed_at,
review_version,
CASE WHEN review_claim_expires_at > NOW() THEN review_claimed_by ELSE NULL END AS review_claimed_by,
review_claim_expires_at,
CASE WHEN review_claim_token = %s AND review_claim_expires_at > NOW() THEN 1 ELSE 0 END AS claimed_by_me
FROM {TARGET_TABLE_NAME_TMP}
WHERE staging_id IN ({placeholders})
""",
[client_token, *staging_ids],
)
return cursor.fetchall()
def _update_staging_review(
staging_id: int,
decision: str,
note: str,
reviewer: str,
expected_version: int,
client_token: str,
) -> dict[str, object]:
"""Persist one claimed review using optimistic concurrency control."""
if decision == "pending":
status_fields = """
biz_review_status = 'pending', reviewed_by = NULL, reviewed_at = NULL,
staging_status = CASE WHEN staging_status = 'deleted' THEN 'staged' ELSE staging_status END
"""
params = [note, source_song_id]
elif decision == "deleted":
# 删除:不限定 dedup_action,对 review/skip 等各种 action 均可执行
sql = f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET staging_status = 'deleted',
biz_review_status = 'deleted',
reviewed_by = %s,
reviewed_at = NOW(),
biz_review_note = NULLIF(%s, '')
WHERE source_song_id = %s
status_fields = """
biz_review_status = 'deleted', reviewed_by = %s, reviewed_at = NOW(),
staging_status = 'deleted'
"""
params = [reviewer, note, source_song_id]
else:
sql = f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET biz_review_status = %s,
reviewed_by = %s,
reviewed_at = NOW(),
biz_review_note = NULLIF(%s, '')
WHERE source_song_id = %s AND dedup_action = 'review'
"""
params = [decision, reviewer, note, source_song_id]
status_fields = "biz_review_status = %s, reviewed_by = %s, reviewed_at = NOW()"
params: list[object]
if decision == "pending":
params = [note, staging_id, expected_version, client_token, reviewer]
elif decision == "deleted":
params = [reviewer, 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:
cursor.execute(sql, params)
affected = cursor.rowcount
cursor.execute(
f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET {status_fields},
biz_review_note = NULLIF(%s, ''),
review_version = review_version + 1,
review_claimed_by = NULL,
review_claim_token = NULL,
review_claimed_at = NULL,
review_claim_expires_at = NULL
WHERE staging_id = %s
AND review_version = %s
AND review_claim_token = %s
AND review_claimed_by = %s
AND review_claim_expires_at > NOW()
""",
params,
)
if cursor.rowcount != 1:
current = _review_snapshot(cursor, staging_id)
conn.rollback()
raise ReviewConflictError("记录已被其他人修改,当前提交未覆盖数据库", current)
current = _review_snapshot(cursor, staging_id)
conn.commit()
return {"source_song_id": source_song_id, "decision": decision, "affected": affected}
return current
def _import_approved_staging(source_ids: list[str], reviewer: str = "dashboard") -> dict[str, object]:
if not source_ids:
raise ValueError("没有可入库的 source_id")
placeholders = ",".join(["%s"] * len(source_ids))
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."""
columns = ", ".join(f"`{column}`" for column in TARGET_COLUMNS)
select_columns = ", ".join(f"s.`{column}`" for column in TARGET_COLUMNS)
update_imported_sql = f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET biz_review_status = 'approved_import',
reviewed_by = %s,
reviewed_at = NOW()
WHERE source_song_id IN ({placeholders})
AND dedup_action = 'review'
AND biz_review_status IN ('pending', 'unsure')
"""
insert_sql = f"""
with _target_conn() as conn, conn.cursor() as cursor:
source_filter = ""
lock_params: list[object] = []
if source_ids:
source_placeholders = ",".join(["%s"] * len(source_ids))
source_filter = f" AND s.source_song_id IN ({source_placeholders})"
lock_params.extend(source_ids)
cursor.execute(
f"""
SELECT s.staging_id, s.source_song_id
FROM {TARGET_TABLE_NAME_TMP} s
INNER JOIN (
SELECT source_song_id, MAX(staging_id) AS max_staging_id
FROM {TARGET_TABLE_NAME_TMP}
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.staging_status = 'staged'
AND (s.review_claim_expires_at IS NULL OR s.review_claim_expires_at <= NOW())
{source_filter}
ORDER BY staging_id
FOR UPDATE
""",
lock_params,
)
approved_rows = cursor.fetchall()
if not approved_rows:
conn.rollback()
raise ValueError("没有可入库的审核通过样本")
staging_ids = [row["staging_id"] for row in approved_rows]
approved_source_ids = [str(row["source_song_id"]) for row in approved_rows]
staging_placeholders = ",".join(["%s"] * len(staging_ids))
cursor.execute(
f"""
INSERT INTO {TARGET_TABLE_NAME} ({columns})
SELECT {select_columns}
FROM {TARGET_TABLE_NAME_TMP} s
WHERE s.source_song_id IN ({placeholders})
AND s.dedup_action = 'review'
AND s.biz_review_status = 'approved_import'
AND s.staging_status <> 'imported'
"""
# 入库后查询目标表实际 ID,正确回填 imported_song_id(不能用暂存表自己的 id)
lookup_target_id_sql = f"""
WHERE s.staging_id IN ({staging_placeholders})
""",
staging_ids,
)
inserted_count = cursor.rowcount
source_placeholders = ",".join(["%s"] * len(approved_source_ids))
cursor.execute(
f"""
SELECT source_song_id, id AS target_id
FROM {TARGET_TABLE_NAME}
WHERE source_song_id IN ({placeholders})
"""
mark_sql_template = f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET staging_status = 'imported',
imported_song_id = %s,
error_message = NULL
WHERE source_song_id = %s
AND dedup_action = 'review'
AND biz_review_status = 'approved_import'
"""
with _target_conn() as conn, conn.cursor() as cursor:
cursor.execute(update_imported_sql, [reviewer, *source_ids])
approved_count = cursor.rowcount
cursor.execute(insert_sql, source_ids)
inserted_count = cursor.rowcount
# 从目标表查出实际生成的主键 ID,逐条回填 imported_song_id
cursor.execute(lookup_target_id_sql, source_ids)
WHERE source_song_id IN ({source_placeholders})
""",
approved_source_ids,
)
id_map = {str(row["source_song_id"]): row["target_id"] for row in cursor.fetchall()}
imported_updated = 0
for sid, target_id in id_map.items():
cursor.execute(mark_sql_template, [target_id, sid])
for approved_row in approved_rows:
sid = str(approved_row["source_song_id"])
target_id = id_map.get(sid)
if target_id is None:
raise RuntimeError(f"入库后未找到目标记录: source_song_id={sid}")
cursor.execute(
f"""
UPDATE {TARGET_TABLE_NAME_TMP}
SET staging_status = 'imported', imported_song_id = %s, error_message = NULL,
review_version = review_version + 1
WHERE staging_id = %s AND staging_status = 'staged'
""",
(target_id, approved_row["staging_id"]),
)
imported_updated += cursor.rowcount
conn.commit()
return {
"approved_count": approved_count,
"approved_count": len(approved_rows),
"inserted_count": inserted_count,
"imported_song_ids_updated": imported_updated,
"source_ids": approved_source_ids,
}
......@@ -671,6 +919,10 @@ def _groups_response(path: Path, top_k: str, mode: str, term: str, page: int, pa
groups = [g for g in groups if g.get("decision") == "new" and not g.get("has_l1_hint")]
elif decision == 'skip':
groups = [g for g in groups if g.get("decision") == "skip"]
elif decision == 'reviewed':
groups = [g for g in groups if g.get("reviewed")]
elif decision == 'pending':
groups = [g for g in groups if not g.get("reviewed")]
elif decision and decision in ('merge', 'review', 'duplicate'):
groups = [g for g in groups if g.get("decision") == decision]
filtered = [group for group in groups if _matches_term(group, term)]
......@@ -779,16 +1031,32 @@ class Handler(BaseHTTPRequestHandler):
mode = params.get("mode", ["retrieval"])[0]
term = _norm(params.get("q", [""])[0])
decision = _norm(params.get("decision", [""])[0])
client_token = params.get("client_token", [""])[0].strip()
page = max(1, int(params.get("page", ["1"])[0]))
page_size = min(200, max(10, int(params.get("page_size", ["50"])[0])))
if raw_path == STAGING_DB_PATH:
_json_response(self, _staging_groups_response(mode, term, page, page_size, decision))
_json_response(
self,
_staging_groups_response(mode, term, page, page_size, decision, client_token),
)
return
path = _safe_path(raw_path)
_json_response(self, _groups_response(path, top_k, mode, term, page, page_size, decision))
except Exception as exc: # noqa: BLE001 - local diagnostic API
_error(self, str(exc), status=404)
return
if parsed.path == "/api/review-statuses":
params = parse_qs(parsed.query)
try:
raw_ids = params.get("staging_ids", [""])[0]
staging_ids = [int(value) for value in raw_ids.split(",") if value.strip()]
if len(staging_ids) > 200:
raise ValueError("最多同步 200 条记录")
client_token = params.get("client_token", [""])[0].strip()
_json_response(self, {"rows": _review_statuses(staging_ids, client_token)})
except Exception as exc: # noqa: BLE001
_error(self, str(exc), status=400)
return
if parsed.path == "/api/query":
params = parse_qs(parsed.query)
try:
......@@ -825,20 +1093,47 @@ class Handler(BaseHTTPRequestHandler):
def do_POST(self) -> None:
parsed = urlparse(self.path)
if parsed.path in {"/api/review-claim", "/api/review-heartbeat"}:
try:
length = int(self.headers.get("Content-Length", "0"))
payload = json.loads(self.rfile.read(length).decode("utf-8") or "{}")
staging_id = int(payload.get("staging_id") or 0)
client_token = str(payload.get("client_token") or "").strip()
if not staging_id:
raise ValueError("staging_id is required")
if parsed.path == "/api/review-claim":
reviewer = str(payload.get("reviewer") or "").strip()
current = _claim_staging_review(staging_id, reviewer, client_token)
else:
current = _heartbeat_staging_review(staging_id, client_token)
_json_response(self, {"record": current})
except ReviewConflictError as exc:
_conflict_response(self, exc)
except Exception as exc: # noqa: BLE001
_error(self, str(exc), status=400)
return
if parsed.path == "/api/staging-review":
try:
length = int(self.headers.get("Content-Length", "0"))
payload = json.loads(self.rfile.read(length).decode("utf-8") or "{}")
source_id = str(payload.get("source_id", "")).strip()
staging_id = int(payload.get("staging_id") or 0)
decision = str(payload.get("decision", "")).strip()
note = str(payload.get("note", "")).strip()
reviewer = str(payload.get("reviewer", "dashboard")).strip() or "dashboard"
if not source_id:
raise ValueError("source_id is required")
reviewer = str(payload.get("reviewer") or "").strip()
client_token = str(payload.get("client_token") or "").strip()
expected_version = int(payload.get("expected_version", -1))
if not staging_id:
raise ValueError("staging_id is required")
if not reviewer or not client_token or expected_version < 0:
raise ValueError("reviewer, client_token and expected_version are required")
if decision not in ("approved_import", "rejected_duplicate", "unsure", "pending", "deleted"):
raise ValueError(f"invalid decision: {decision}")
result = _update_staging_review(source_id, decision, note, reviewer)
_json_response(self, result)
result = _update_staging_review(
staging_id, decision, note, reviewer, expected_version, client_token
)
_json_response(self, {"record": result})
except ReviewConflictError as exc:
_conflict_response(self, exc)
except Exception as exc: # noqa: BLE001
_error(self, str(exc), status=400)
return
......@@ -848,23 +1143,19 @@ class Handler(BaseHTTPRequestHandler):
try:
length = int(self.headers.get("Content-Length", "0"))
payload = json.loads(self.rfile.read(length).decode("utf-8") or "{}")
rows = payload.get("rows") or []
if not isinstance(rows, list):
raise ValueError("rows must be a list")
source_ids = payload.get("source_ids")
if source_ids is not None and not isinstance(source_ids, list):
raise ValueError("source_ids must be a list")
result = _import_approved_staging(
[str(value).strip() for value in source_ids if str(value).strip()]
if source_ids
else None
)
approved = [
{
"source_id": str(row.get("source_id", "")).strip(),
"review_decision": "import",
"review_note": str(row.get("review_note", "")).strip(),
}
for row in rows
if str(row.get("source_id", "")).strip()
and str(row.get("review_decision", "")).strip().lower() in {"import", "导入", "not_duplicate"}
{"source_id": source_id, "review_decision": "import", "review_note": ""}
for source_id in result["source_ids"]
]
if not approved:
raise ValueError("没有可入库的审核通过样本")
review_csv = _write_review_import_csv(approved)
result = _import_approved_staging([row["source_id"] for row in approved])
_json_response(self, {"review_csv": str(review_csv), **result})
except Exception as exc: # noqa: BLE001 - local diagnostic API
_error(self, str(exc), status=400)
......@@ -877,6 +1168,8 @@ class Handler(BaseHTTPRequestHandler):
body = path.read_bytes()
self.send_response(200)
self.send_header("Content-Type", content_type)
if path.suffix.lower() in {".html", ".js", ".css"}:
self.send_header("Cache-Control", "no-cache")
self.send_header("Content-Length", str(len(body)))
self.end_headers()
self.wfile.write(body)
......@@ -886,8 +1179,14 @@ def main() -> None:
parser = argparse.ArgumentParser(description="Serve L2 lyric review dashboard")
parser.add_argument("--host", default="127.0.0.1")
parser.add_argument("--port", type=int, default=8765)
parser.add_argument(
"--migrate-review-schema",
action="store_true",
help="add the staging-table columns/index required for collaborative reviewing",
)
args = parser.parse_args()
_ensure_review_schema(apply_migration=args.migrate_review_schema)
server = ThreadingHTTPServer((args.host, args.port), Handler)
print(f"L2 dashboard: http://{args.host}:{args.port}")
server.serve_forever()
......
from __future__ import annotations
from pathlib import Path
from unittest.mock import patch
import pytest
import serve_l2_dashboard as dashboard
class FakeCursor:
def __init__(self, update_rowcount: int, snapshot: dict[str, object]) -> None:
self.update_rowcount = update_rowcount
self.snapshot = snapshot
self.rowcount = 0
self.executions: list[tuple[str, object]] = []
def __enter__(self):
return self
def __exit__(self, *_args):
return False
def execute(self, sql, params=None):
self.executions.append((sql, params))
self.rowcount = self.update_rowcount if sql.lstrip().startswith("UPDATE") else 0
def fetchone(self):
return self.snapshot
class FakeConnection:
def __init__(self, cursor: FakeCursor) -> None:
self._cursor = cursor
self.committed = False
self.rolled_back = False
def __enter__(self):
return self
def __exit__(self, *_args):
return False
def cursor(self):
return self._cursor
def commit(self):
self.committed = True
def rollback(self):
self.rolled_back = True
def test_review_update_uses_staging_version_and_claim_token():
snapshot = {
"staging_id": 42,
"source_song_id": "song-1",
"biz_review_status": "approved_import",
"review_version": 8,
}
cursor = FakeCursor(update_rowcount=1, snapshot=snapshot)
conn = FakeConnection(cursor)
with patch.object(dashboard, "_target_conn", return_value=conn):
result = dashboard._update_staging_review(
42, "approved_import", "ok", "alice", 7, "browser-token"
)
update_sql, params = cursor.executions[0]
assert "WHERE staging_id = %s" in update_sql
assert "review_version = %s" in update_sql
assert "review_claim_token = %s" in update_sql
assert params[-4:] == [42, 7, "browser-token", "alice"]
assert result == snapshot
assert conn.committed
def test_stale_review_returns_conflict_without_overwrite():
snapshot = {
"staging_id": 42,
"biz_review_status": "rejected_duplicate",
"reviewed_by": "bob",
"review_version": 9,
}
cursor = FakeCursor(update_rowcount=0, snapshot=snapshot)
conn = FakeConnection(cursor)
with patch.object(dashboard, "_target_conn", return_value=conn):
with pytest.raises(dashboard.ReviewConflictError) as raised:
dashboard._update_staging_review(
42, "approved_import", "stale", "alice", 7, "browser-token"
)
assert raised.value.current_record == snapshot
assert conn.rolled_back
assert not conn.committed
def test_claim_conflict_exposes_current_owner():
snapshot = {
"staging_id": 42,
"review_claimed_by": "bob",
"review_version": 3,
}
cursor = FakeCursor(update_rowcount=0, snapshot=snapshot)
conn = FakeConnection(cursor)
with patch.object(dashboard, "_target_conn", return_value=conn):
with pytest.raises(dashboard.ReviewConflictError) as raised:
dashboard._claim_staging_review(42, "alice", "alice-token")
assert raised.value.current_record["review_claimed_by"] == "bob"
assert conn.rolled_back
def test_heartbeat_does_not_increment_review_version():
snapshot = {"staging_id": 42, "review_version": 5}
cursor = FakeCursor(update_rowcount=1, snapshot=snapshot)
conn = FakeConnection(cursor)
with patch.object(dashboard, "_target_conn", return_value=conn):
dashboard._heartbeat_staging_review(42, "browser-token")
update_sql, _params = cursor.executions[0]
assert "review_version" not in update_sql
assert conn.committed
class ImportCursor(FakeCursor):
def __init__(self) -> None:
super().__init__(update_rowcount=1, snapshot={})
self.result_sets = [
[{"staging_id": 42, "source_song_id": "song-1"}],
[{"source_song_id": "song-1", "target_id": 9001}],
]
def execute(self, sql, params=None):
self.executions.append((sql, params))
self.rowcount = 1
def fetchall(self):
return self.result_sets.pop(0)
def test_import_locks_only_unclaimed_staged_rows():
cursor = ImportCursor()
conn = FakeConnection(cursor)
with patch.object(dashboard, "_target_conn", return_value=conn):
result = dashboard._import_approved_staging()
lock_sql, _params = cursor.executions[0]
assert "FOR UPDATE" in lock_sql
assert "staging_status = 'staged'" in lock_sql
assert "review_claim_expires_at <= NOW()" in lock_sql
assert result["inserted_count"] == 1
assert result["source_ids"] == ["song-1"]
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'))")
end = html.index("function selectedQueryRows()", start)
click_handler = html[start:end]
assert "loadQueryRows()" in click_handler
assert "claimCurrentReview" in click_handler
assert "loadGroups()" not in click_handler
def test_staging_bigint_is_never_converted_to_javascript_number():
html = Path("l2_review_dashboard.html").read_text(encoding="utf-8")
assert "Number(query.staging_id)" not in html
assert "staging_id: String(query.staging_id)" in html
def test_dashboard_distinguishes_reviewed_from_submitted():
html = Path("l2_review_dashboard.html").read_text(encoding="utf-8")
row = dashboard._staging_dashboard_row(
{
"staging_id": 42,
"source_song_id": "song-1",
"dedup_action": "review",
"biz_review_status": "approved_import",
"staging_status": "imported",
}
)
assert row["staging_status"] == "imported"
assert '<option value="reviewed">已审核</option>' in html
assert '<option value="submitted">已提交</option>' in html