Files
zyt/server/app/common/service/followupaudio/FollowupAudioStore.php
T

442 lines
25 KiB
PHP

<?php
declare(strict_types=1);
namespace app\common\service\followupaudio;
use app\common\service\prescriptionai\PrescriptionAiCipher;
use DomainException;
use think\facade\Db;
/** Database queue. No plaintext transcript, credentials or review data in task metadata/logs. */
final class FollowupAudioStore
{
public static function enabled(): bool { return (bool) config('followup_audio.enabled', false); }
public static function verified(?string $profile = null): bool
{
$profiles = self::verifiedProfiles();
return (bool) config('followup_audio.audio_verified', false) && ($profile === null ? $profiles !== [] : in_array($profile, $profiles, true));
}
private static function verifiedProfiles(): array
{
return array_values(array_intersect(['qwen', 'openai'], (array) config('followup_audio.verified_profiles', [])));
}
public static function assertEnabled(?string $profile = null): void
{
if (!self::enabled() || !self::verified($profile)) { throw new DomainException('FOLLOWUP_AUDIO_DISABLED_OR_UNVERIFIED'); }
}
public static function create(array $upload, string $recordedAt, string $modelKey, int $actor, array $info): array
{
self::assertEnabled($modelKey);
FollowupAudioPolicy::strictRecordedAt($recordedAt);
if (!in_array($modelKey, ['qwen', 'openai'], true)) { throw new DomainException('FOLLOWUP_AUDIO_MODEL_INVALID'); }
$created = Db::transaction(static function () use ($upload, $recordedAt, $modelKey, $actor, $info): array {
$stored = Db::name('followup_audio_upload')->where('id', (string) ($upload['id'] ?? ''))->lock(true)->find();
if (!$stored || (int) $stored['actor_id'] !== $actor || $stored['status'] !== 'complete' || (int) $stored['expires_at'] <= time()) {
throw new DomainException('FOLLOWUP_AUDIO_UPLOAD_UNAVAILABLE');
}
if ((float) $stored['duration_seconds'] <= 0 || (float) $stored['duration_seconds'] > 3600
|| (int) $stored['total_bytes'] <= 0 || (int) $stored['total_bytes'] > 524288000) {
throw new DomainException('FOLLOWUP_AUDIO_UPLOAD_METADATA_INVALID');
}
$diagnosisId = (int) $stored['diagnosis_id'];
FollowupAudioAccess::diagnosis($diagnosisId, $actor, $info);
$diagnosis = FollowupAudioApply::diagnosis($diagnosisId, true);
$existing = Db::name('followup_audio_task')->where('upload_id', $stored['id'])->lock(true)->find();
if ($existing && $existing['model_key'] !== $modelKey) {
throw new DomainException('FOLLOWUP_AUDIO_UPLOAD_ALREADY_USED');
}
// Diagnose+content+profile is immutable even if a retry changes upload ID, filename or recordedAt.
$existing = $existing ?: Db::name('followup_audio_task')->where('diagnosis_id', $diagnosisId)
->where('sha256', $stored['sha256'])->where('model_key', $modelKey)->lock(true)->find();
if ($existing) { return self::reuse($existing, $recordedAt); }
$path = FollowupAudioUpload::path($stored);
if (!is_file($path) || (int) filesize($path) !== (int) $stored['total_bytes']
|| !hash_equals((string) $stored['sha256'], (string) hash_file('sha256', $path))) {
throw new DomainException('FOLLOWUP_AUDIO_UPLOAD_INTEGRITY_FAILED');
}
$now = time();
$expiresAt = $now + max(1, min(90, (int) config('followup_audio.retention_days', 90))) * 86400;
// The task and its audio have one retention deadline; upload staging TTL must not erase an active task.
Db::name('followup_audio_upload')->where('id', $stored['id'])->update(['expires_at' => $expiresAt]);
try {
$id = (int) Db::name('followup_audio_task')->insertGetId([
'diagnosis_id' => $diagnosisId, 'patient_id' => (int) $diagnosis['patient_id'], 'actor_id' => $actor,
'upload_id' => $stored['id'], 'file_name' => $stored['file_name'], 'sha256' => $stored['sha256'],
'duration_seconds' => $stored['duration_seconds'], 'recorded_at' => $recordedAt, 'model_key' => $modelKey,
'status' => 'queued', 'stage' => 'queued', 'version' => 1, 'attempts' => 0,
'lease_token' => '', 'lease_until' => 0, 'upstream_started_at' => 0,
'upstream_run_id' => '', 'upstream_file_id' => '', 'upstream_ids_json' => '{}',
'error_code' => '', 'error_message' => '', 'extraction_cipher' => '', 'review_cipher' => '', 'applied_cipher' => '',
'created_at' => $now, 'updated_at' => $now, 'expires_at' => $expiresAt, 'applied_at' => 0, 'purged_at' => 0,
]);
return ['id' => $id, 'reused' => false, 'reuse_message' => ''];
} catch (\think\db\exception\PDOException $exception) {
// The unique content key is the final arbiter even for a concurrent caller outside this service.
if ((int) ($exception->getData()['PDO Error Info']['Driver Error Code'] ?? 0) !== 1062) { throw $exception; }
$winner = Db::name('followup_audio_task')->where('diagnosis_id', $diagnosisId)
->where('sha256', $stored['sha256'])->where('model_key', $modelKey)->lock(true)->find();
if (!$winner) { throw $exception; }
Db::name('followup_audio_upload')->where('id', $stored['id'])->update(['expires_at' => $stored['expires_at']]);
return self::reuse($winner, $recordedAt);
}
});
return ['task_id' => $created['id'], 'reused' => $created['reused'], 'reuse_message' => $created['reuse_message']]
+ self::summary(self::task($created['id']));
}
private static function reuse(array $task, string $recordedAt): array
{
$message = '同一诊单、录音内容及模型已有任务,已复用原任务,不会再次调用模型。';
if ($task['recorded_at'] !== $recordedAt) {
$message .= '沿用原任务的录音时间,请在未采用的审阅项中更正记录日期;已采用任务不会重跑。';
}
if ($task['status'] === 'needs_reconciliation' || (int) $task['upstream_started_at'] > 0 && $task['status'] === 'failed') {
$message .= '原上游结果待核对,重复上传不能触发重试。';
}
return ['id' => (int) $task['id'], 'reused' => true, 'reuse_message' => $message];
}
public static function task(int $taskId): array
{
$task = Db::name('followup_audio_task')->where('id', $taskId)->find();
if (!$task) { throw new DomainException('FOLLOWUP_AUDIO_TASK_UNAVAILABLE'); }
return $task;
}
public static function lists(int $diagnosisId): array
{
$rows = Db::name('followup_audio_task')->where('diagnosis_id', $diagnosisId)->order('id', 'desc')->limit(200)->select()->toArray();
return array_map([self::class, 'summary'], $rows);
}
public static function summary(array $task): array
{
$result = [];
foreach (['id', 'diagnosis_id', 'file_name', 'recorded_at', 'status', 'stage', 'error_code', 'error_message',
'version', 'created_at', 'expires_at'] as $key) { $result[$key] = $task[$key]; }
foreach (['id', 'diagnosis_id', 'version', 'created_at', 'expires_at'] as $key) { $result[$key] = (int) $result[$key]; }
$alive = (int) $task['expires_at'] > time() && !(int) $task['purged_at'];
if (!$alive && $task['status'] !== 'applied') {
$result['status'] = 'expired';
$result['stage'] = 'expired';
$result['error_code'] = 'FOLLOWUP_AUDIO_EXPIRED';
$result['error_message'] = '录音及审阅已到保留期限,不能采用';
}
$result['can_retry'] = self::enabled() && self::verified((string) $task['model_key']) && $alive && $task['status'] === 'failed'
&& (int) $task['upstream_started_at'] === 0 && (int) $task['attempts'] < 3;
$result['audio_available'] = $alive && (string) Db::name('followup_audio_upload')->where('id', $task['upload_id'])->value('status') === 'complete';
return $result;
}
public static function detail(int $taskId): array
{
$task = self::task($taskId);
$data = self::summary($task) + ['summary' => '', 'transcript' => '', 'uncertainties' => [], 'items' => [], 'applied_items' => []];
// Expiry is enforced at read/apply time, not only when a scheduled cleanup eventually runs.
if ((int) $task['expires_at'] > time() && !(int) $task['purged_at'] && $task['extraction_cipher'] !== '') {
$extraction = self::open($taskId, 'extraction', $task['extraction_cipher']);
$review = self::open($taskId, 'review', $task['review_cipher']);
$data['summary'] = $extraction['summary'];
$data['transcript'] = $extraction['transcript'];
$data['uncertainties'] = $extraction['uncertainties'];
$data['items'] = $review['items'];
}
if ($task['applied_cipher'] !== '') {
$data['applied_items'] = self::open($taskId, 'applied', $task['applied_cipher'])['items'];
}
return $data;
}
public static function saveDraft(int $taskId, int $version, array $items, int $actor): array
{
self::assertEnabled();
$stale = Db::transaction(static function () use ($taskId, $version, $items, $actor): bool {
$task = self::lockTask($taskId);
self::assertEnabled((string) $task['model_key']);
FollowupAudioAccess::task($taskId, $actor, []);
self::assertReview($task, $version);
$review = self::open($taskId, 'review', $task['review_cipher']);
$extraction = self::open($taskId, 'extraction', $task['extraction_cipher']);
$merged = FollowupAudioApply::mergeItems($review['items'], $items, $extraction['items']);
$diagnosis = FollowupAudioApply::diagnosis((int) $task['diagnosis_id'], true);
FollowupAudioApply::assertPatient($task, $diagnosis);
$refreshed = FollowupAudioApply::refresh($task, $merged, $diagnosis, false);
$changed = false;
foreach ($merged as $index => $item) {
if (!hash_equals($item['expected_hash'], $refreshed[$index]['expected_hash'])) {
$changed = true;
break;
}
}
if ($changed) {
// A refresh is committed, but adoption must be an explicit subsequent request/version.
$merged = $refreshed;
foreach ($merged as &$item) { $item['selected'] = false; $item['needs_review'] = true; }
unset($item);
}
self::writeReview($task, $merged);
return $changed;
});
$detail = self::detail($taskId);
if ($stale) { $detail['review_refreshed'] = true; }
return $detail;
}
public static function retry(int $taskId): array
{
self::assertEnabled();
Db::transaction(static function () use ($taskId): void {
$task = self::lockTask($taskId);
self::assertEnabled((string) $task['model_key']);
if ($task['status'] !== 'failed' || (int) $task['upstream_started_at'] !== 0 || (int) $task['attempts'] >= 3
|| (int) $task['expires_at'] <= time() || (int) $task['purged_at']) {
throw new DomainException('FOLLOWUP_AUDIO_RETRY_NOT_SAFE');
}
Db::name('followup_audio_task')->where('id', $taskId)->update([
'status' => 'queued', 'stage' => 'queued', 'error_code' => '', 'error_message' => '',
'version' => (int) $task['version'] + 1, 'updated_at' => time(),
]);
});
return self::summary(self::task($taskId));
}
/** A singleton DB mutex enforces concurrency across independent CLI processes. */
public static function claim(): ?array
{
if (!self::enabled() || !self::verified()) { return null; }
return Db::transaction(static function (): ?array {
if (!Db::name('followup_audio_mutex')->where('id', 1)->lock(true)->find()) {
throw new DomainException('FOLLOWUP_AUDIO_MIGRATION_REQUIRED');
}
$now = time();
$expired = Db::name('followup_audio_task')->where('status', 'running')->where('lease_until', '<=', $now)->lock(true)->select()->toArray();
foreach ($expired as $task) {
$uncertain = (int) $task['upstream_started_at'] > 0;
Db::name('followup_audio_task')->where('id', $task['id'])->update([
'status' => $uncertain ? 'needs_reconciliation' : 'failed', 'stage' => $uncertain ? 'needs_reconciliation' : 'failed',
'error_code' => $uncertain ? 'FOLLOWUP_AUDIO_UPSTREAM_UNCERTAIN' : 'FOLLOWUP_AUDIO_LEASE_EXPIRED',
'error_message' => $uncertain ? '上游结果待核对,禁止重复提交' : '处理进程中断,请检查后重试',
'lease_token' => '', 'lease_until' => 0, 'updated_at' => $now, 'version' => (int) $task['version'] + 1,
]);
}
// A locking CURRENT read is essential: a plain COUNT can reuse a pre-mutex REPEATABLE READ snapshot.
$active = Db::name('followup_audio_task')->where('status', 'running')->where('lease_until', '>', $now)
->field('id')->lock(true)->select()->toArray();
if (count($active) >= max(1, min(8, (int) config('followup_audio.concurrency', 1)))) { return null; }
$task = Db::name('followup_audio_task')->where('status', 'queued')->where('upstream_started_at', 0)->whereIn('model_key', self::verifiedProfiles())
->where('expires_at', '>', $now)->where('purged_at', 0)->order('id', 'asc')->lock(true)->find();
if (!$task) { return null; }
$changes = ['status' => 'running', 'stage' => 'preparing', 'lease_token' => bin2hex(random_bytes(32)),
'lease_until' => $now + self::leaseSeconds(), 'attempts' => (int) $task['attempts'] + 1,
'updated_at' => $now, 'version' => (int) $task['version'] + 1];
Db::name('followup_audio_task')->where('id', $task['id'])->update($changes);
return array_replace($task, $changes);
});
}
public static function heartbeat(int $id, string $token): bool
{
return Db::transaction(static function () use ($id, $token): bool {
$task = self::lockTask($id);
if (!self::hasLease($task, $token)) { return false; }
Db::name('followup_audio_task')->where('id', $id)->update(['lease_until' => time() + self::leaseSeconds(), 'updated_at' => time()]);
return true;
});
}
/** Only opaque upstream identifiers; never accepts a URL, prompt, audio, response or credentials. */
public static function checkpoint(int $id, string $token, array $fields): bool
{
return Db::transaction(static function () use ($id, $token, $fields): bool {
$task = self::lockTask($id);
if (!self::hasLease($task, $token)) { return false; }
$changes = ['updated_at' => time(), 'lease_until' => time() + self::leaseSeconds()];
foreach ($fields as $key => $value) {
if ($key === 'upstream_started_at') {
if (!is_int($value) || $value <= 0) { throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID'); }
$changes[$key] = (int) $task[$key] > 0 ? (int) $task[$key] : time();
} elseif ($key === 'stage') {
if (!in_array($value, ['preparing', 'uploading', 'analyzing', 'validating'], true)) { throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID'); }
$changes[$key] = $value;
} elseif (in_array($key, ['upstream_run_id', 'upstream_file_id'], true)) {
$changes[$key] = self::opaqueId($value);
} elseif ($key === 'upstream_ids_json') {
$ids = is_string($value) ? json_decode($value, true, 16, JSON_THROW_ON_ERROR) : $value;
if (!is_array($ids)) { throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID'); }
$previous = json_decode($task['upstream_ids_json'] ?: '{}', true, 16, JSON_THROW_ON_ERROR);
foreach ($ids as $name => $identifier) {
if (!in_array($name, ['request_id', 'task_id', 'message_id', 'conversation_id', 'upstream_request_id'], true)) {
throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID');
}
$previous[$name] = self::opaqueId($identifier);
}
$changes[$key] = FollowupAudioPolicy::canonical($previous);
} else { throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID'); }
}
Db::name('followup_audio_task')->where('id', $id)->update($changes);
return true;
});
}
public static function complete(int $id, string $token, array $extraction): bool
{
return Db::transaction(static function () use ($id, $token, $extraction): bool {
$task = self::lockTask($id);
if (!self::hasLease($task, $token)) { return false; }
$normalized = FollowupAudioPolicy::normalizeExtraction($extraction, $task['recorded_at']);
$durationMs = (int) ceil((float) $task['duration_seconds'] * 1000);
foreach ($normalized['items'] as $item) {
foreach ($item['evidence'] as $evidence) {
if (($evidence['start_ms'] ?? 0) > $durationMs || ($evidence['end_ms'] ?? 0) > $durationMs) {
throw new DomainException('FOLLOWUP_AUDIO_EVIDENCE_OUTSIDE_AUDIO');
}
}
}
$diagnosis = FollowupAudioApply::diagnosis((int) $task['diagnosis_id'], true);
FollowupAudioApply::assertPatient($task, $diagnosis);
$items = FollowupAudioApply::refresh($task, $normalized['items'], $diagnosis, true);
Db::name('followup_audio_task')->where('id', $id)->update([
'extraction_cipher' => self::seal($id, 'extraction', $normalized),
'review_cipher' => self::seal($id, 'review', ['items' => $items]),
'status' => 'review', 'stage' => 'review', 'lease_token' => '', 'lease_until' => 0,
'updated_at' => time(), 'version' => (int) $task['version'] + 1, 'error_code' => '', 'error_message' => '',
]);
return true;
});
}
public static function fail(int $id, string $token, string $code, string $message, bool $uncertain = false): bool
{
return Db::transaction(static function () use ($id, $token, $code, $uncertain): bool {
$task = self::lockTask($id);
if (!self::hasLease($task, $token)) { return false; }
// An arbitrary exception/message can contain patient text or credentials: never persist it.
$code = preg_match('/^[A-Z][A-Z0-9_]{2,95}$/D', $code) ? $code : 'FOLLOWUP_AUDIO_PROCESS_FAILED';
$status = $uncertain ? 'needs_reconciliation' : 'failed';
Db::name('followup_audio_task')->where('id', $id)->update([
'status' => $status, 'stage' => $status, 'error_code' => $code,
'error_message' => $uncertain ? '上游结果待核对,禁止重复提交' : '处理未完成,请查看错误代码;不会自动重复调用',
'lease_token' => '', 'lease_until' => 0, 'updated_at' => time(), 'version' => (int) $task['version'] + 1,
]);
return true;
});
}
public static function cleanupExpired(): array
{
$counts = ['tasks_purged' => 0, 'uploads_deleted' => 0, 'active_skipped' => 0, 'errors' => 0];
$attemptedUploads = [];
$ids = Db::name('followup_audio_task')->where('expires_at', '<=', time())->where('purged_at', 0)->order('id')->limit(500)->column('id');
foreach ($ids as $id) {
try {
$upload = Db::transaction(static function () use ($id, &$counts): ?string {
$task = self::lockTask((int) $id);
if ((int) $task['lease_until'] > time()) { $counts['active_skipped']++; return null; }
if (!(int) $task['purged_at']) {
Db::name('followup_audio_task')->where('id', $id)->update([
'extraction_cipher' => '', 'review_cipher' => '', 'file_name' => '已清理录音', 'purged_at' => time(),
'status' => $task['status'] === 'applied' ? 'applied' : 'expired',
'stage' => $task['status'] === 'applied' ? 'applied' : 'expired',
'lease_token' => '', 'lease_until' => 0, 'updated_at' => time(), 'version' => (int) $task['version'] + 1,
]);
$counts['tasks_purged']++;
}
return (string) $task['upload_id'];
});
if ($upload !== null) {
$attemptedUploads[$upload] = true;
if (self::cleanupUpload($upload, $counts)) { $counts['uploads_deleted']++; }
}
} catch (\Throwable $exception) { $counts['errors']++; }
}
// Retry previously purged file deletions and clean unsubmitted staging uploads. Completed older rows never starve newer expiry work.
$orphans = Db::name('followup_audio_upload')->where('expires_at', '<=', time())->where('status', '<>', 'deleted')->limit(500)->column('id');
foreach ($orphans as $id) {
try {
if (isset($attemptedUploads[$id])) { continue; }
$task = Db::name('followup_audio_task')->where('upload_id', $id)->field('purged_at')->find();
if ((!$task || (int) $task['purged_at']) && self::cleanupUpload((string) $id, $counts)) { $counts['uploads_deleted']++; }
} catch (\Throwable $exception) { $counts['errors']++; }
}
return $counts;
}
/** Match create's upload-row lock order; an expiry/create race must never unlink newly retained audio. */
private static function cleanupUpload(string $id, array &$counts): bool
{
return Db::transaction(static function () use ($id, &$counts): bool {
$upload = Db::name('followup_audio_upload')->where('id', $id)->lock(true)->find();
if (!$upload || $upload['status'] === 'deleted' || (int) $upload['expires_at'] > time()) { return false; }
$tasks = Db::name('followup_audio_task')->where('upload_id', $id)->lock(true)->select()->toArray();
foreach ($tasks as $task) {
if ((int) $task['lease_until'] > time() || !(int) $task['purged_at']) {
$counts['active_skipped']++;
return false;
}
}
FollowupAudioUpload::cleanup($id);
Db::name('followup_audio_upload')->where('id', $id)->update(['file_name' => '已清理录音']);
return true;
});
}
public static function lockTask(int $id): array
{
$task = Db::name('followup_audio_task')->where('id', $id)->lock(true)->find();
if (!$task) { throw new DomainException('FOLLOWUP_AUDIO_TASK_UNAVAILABLE'); }
return $task;
}
public static function assertReview(array $task, int $version): void
{
if ((int) $task['expires_at'] <= time() || (int) $task['purged_at']) { throw new DomainException('FOLLOWUP_AUDIO_EXPIRED'); }
if ($task['status'] !== 'review') { throw new DomainException('FOLLOWUP_AUDIO_NOT_REVIEWABLE'); }
if ((int) $task['version'] !== $version) { throw new DomainException('FOLLOWUP_AUDIO_VERSION_CONFLICT'); }
}
public static function writeReview(array $task, array $items): void
{
Db::name('followup_audio_task')->where('id', $task['id'])->update([
'review_cipher' => self::seal((int) $task['id'], 'review', ['items' => $items]),
'version' => (int) $task['version'] + 1, 'updated_at' => time(),
]);
}
public static function seal(int $id, string $purpose, array $payload): string
{
return self::cipher()->encrypt($payload, 'followup-audio:' . $id . ':' . $purpose);
}
public static function open(int $id, string $purpose, string $ciphertext): array
{
return self::cipher()->decrypt($ciphertext, 'followup-audio:' . $id . ':' . $purpose);
}
private static function cipher(): PrescriptionAiCipher
{
$key = (string) config('followup_audio.encryption_key', '');
return new PrescriptionAiCipher($key === '' ? null : $key);
}
private static function hasLease(array $task, string $token): bool
{
return self::enabled() && self::verified((string) $task['model_key']) && $task['status'] === 'running' && $token !== ''
&& hash_equals((string) $task['lease_token'], $token) && (int) $task['lease_until'] > time()
&& (int) $task['expires_at'] > time() && !(int) $task['purged_at'];
}
private static function leaseSeconds(): int { return max(30, (int) config('followup_audio.lease_seconds', 600)); }
private static function opaqueId($value): string
{
if (!is_string($value) || !preg_match('/^[a-zA-Z0-9._:-]{1,191}$/D', $value)) {
throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID');
}
return $value;
}
}