feat: add durable ASR-to-patient extraction and guarded review
This commit is contained in:
@@ -149,7 +149,11 @@ final class FollowupAudioStore
|
||||
$result['error_message'] = '录音及审阅已到保留期限,不能采用';
|
||||
}
|
||||
$result['can_retry'] = self::enabled() && self::verified((string) $task['model_key']) && $alive && $task['status'] === 'failed'
|
||||
&& self::providerMatches($task) && (int) $task['upstream_started_at'] === 0 && (int) $task['attempts'] < 3;
|
||||
&& self::providerMatches($task) && self::safeToResume($task) && (int) $task['attempts'] < 3;
|
||||
$pipeline = $alive ? self::pipelineState($task) : [];
|
||||
$segments = $pipeline ? FollowupAudioPipelineCheckpoint::segments($pipeline) : [];
|
||||
$result['pipeline_progress'] = ['transcribed' => count($segments), 'total' => count($pipeline['chunks'] ?? [])];
|
||||
$result['transcript_available'] = $alive && ($task['extraction_cipher'] !== '' || trim(implode("\n", array_column($segments, 'text'))) !== '');
|
||||
$result['audio_available'] = $alive && (string) Db::name('followup_audio_upload')->where('id', $task['upload_id'])->value('status') === 'complete';
|
||||
return $result;
|
||||
}
|
||||
@@ -167,6 +171,10 @@ final class FollowupAudioStore
|
||||
$data['uncertainties'] = $extraction['uncertainties'];
|
||||
$data['items'] = $review['items'];
|
||||
}
|
||||
if ((int) $task['expires_at'] > time() && !(int) $task['purged_at'] && $task['extraction_cipher'] === '') {
|
||||
$pipeline = self::pipelineState($task);
|
||||
if ($pipeline) { $data['transcript'] = implode("\n", array_column(FollowupAudioPipelineCheckpoint::segments($pipeline), 'text')); }
|
||||
}
|
||||
if ($task['applied_cipher'] !== '') {
|
||||
$data['applied_items'] = self::open($taskId, 'applied', $task['applied_cipher'])['items'];
|
||||
}
|
||||
@@ -215,7 +223,7 @@ final class FollowupAudioStore
|
||||
$task = self::lockTask($taskId);
|
||||
self::assertEnabled((string) $task['model_key']);
|
||||
self::assertTaskProvider($task);
|
||||
if ($task['status'] !== 'failed' || (int) $task['upstream_started_at'] !== 0 || (int) $task['attempts'] >= 3
|
||||
if ($task['status'] !== 'failed' || !self::safeToResume($task) || (int) $task['attempts'] >= 3
|
||||
|| (int) $task['expires_at'] <= time() || (int) $task['purged_at']) {
|
||||
throw new DomainException('FOLLOWUP_AUDIO_RETRY_NOT_SAFE');
|
||||
}
|
||||
@@ -238,7 +246,7 @@ final class FollowupAudioStore
|
||||
$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;
|
||||
$uncertain = !self::safeToResume($task);
|
||||
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',
|
||||
@@ -250,10 +258,18 @@ final class FollowupAudioStore
|
||||
$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; }
|
||||
$pending = Db::name('followup_audio_task')->where('status', 'queued')->where('upstream_started_at', 0)
|
||||
$pending = Db::name('followup_audio_task')->where('status', 'queued')
|
||||
->where('expires_at', '>', $now)->where('purged_at', 0)->order('id', 'asc')->limit(200)->lock(true)->select()->toArray();
|
||||
$task = null;
|
||||
foreach ($pending as $candidate) {
|
||||
if (!self::safeToResume($candidate)) {
|
||||
Db::name('followup_audio_task')->where('id', $candidate['id'])->update([
|
||||
'status' => 'needs_reconciliation', 'stage' => 'needs_reconciliation',
|
||||
'error_code' => 'RECONCILIATION_REQUIRED', 'error_message' => '上游结果待核对,禁止重复提交',
|
||||
'updated_at' => $now, 'version' => (int) $candidate['version'] + 1,
|
||||
]);
|
||||
continue;
|
||||
}
|
||||
if (!self::providerMatches($candidate)) {
|
||||
// Old rows without a binding and changed configurations are visible failures, never silently re-routed.
|
||||
Db::name('followup_audio_task')->where('id', $candidate['id'])->update([
|
||||
@@ -296,10 +312,18 @@ final class FollowupAudioStore
|
||||
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'); }
|
||||
if (!in_array($value, ['preparing', 'uploading', 'analyzing', 'transcribing', 'extracting', '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 === 'pipeline_checkpoint') {
|
||||
if (!is_array($value) || !is_int($value['expected_revision'] ?? null) || !is_array($value['state'] ?? null)
|
||||
|| FollowupAudioProviderConfig::resolve((string) $task['model_key'])['driver'] !== 'asr_then_llm') {
|
||||
throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID');
|
||||
}
|
||||
$previous = empty($task['pipeline_cipher']) ? [] : self::open($id, 'pipeline', $task['pipeline_cipher']);
|
||||
FollowupAudioPipelineCheckpoint::transition($previous, $value['state'], $task, $value['expected_revision']);
|
||||
$changes['pipeline_cipher'] = self::seal($id, 'pipeline', $value['state']);
|
||||
} 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'); }
|
||||
@@ -359,7 +383,11 @@ final class FollowupAudioStore
|
||||
if (!self::hasLease($task, $token, false)) { 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';
|
||||
$uncertain = $uncertain || (int) $task['upstream_started_at'] > 0;
|
||||
// For the explicit pipeline, durable state distinguishes a known response from an unresolved intent.
|
||||
// Legacy requests keep their historical sticky uncertainty policy.
|
||||
$pipeline = self::pipelineState($task);
|
||||
$uncertain = $pipeline !== [] ? FollowupAudioPipelineCheckpoint::hasIntent($pipeline)
|
||||
: ($uncertain || (int) $task['upstream_started_at'] > 0);
|
||||
$status = $uncertain ? 'needs_reconciliation' : 'failed';
|
||||
Db::name('followup_audio_task')->where('id', $id)->update([
|
||||
'status' => $status, 'stage' => $status, 'error_code' => $code,
|
||||
@@ -385,7 +413,7 @@ final class FollowupAudioStore
|
||||
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(),
|
||||
'extraction_cipher' => '', 'review_cipher' => '', 'pipeline_cipher' => null, '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,
|
||||
@@ -477,6 +505,26 @@ final class FollowupAudioStore
|
||||
&& (int) $task['expires_at'] > time() && !(int) $task['purged_at'];
|
||||
}
|
||||
|
||||
/** Corrupt/unbound encrypted state is never an excuse to resend a billable request. */
|
||||
private static function pipelineState(array $task): array
|
||||
{
|
||||
if (empty($task['pipeline_cipher'])) { return []; }
|
||||
try {
|
||||
$state = self::open((int) $task['id'], 'pipeline', $task['pipeline_cipher']);
|
||||
FollowupAudioPipelineCheckpoint::validate($state, $task);
|
||||
return $state;
|
||||
} catch (\Throwable $error) { return []; }
|
||||
}
|
||||
|
||||
private static function safeToResume(array $task): bool
|
||||
{
|
||||
if (!empty($task['pipeline_cipher'])) {
|
||||
$state = self::pipelineState($task);
|
||||
return $state !== [] && !FollowupAudioPipelineCheckpoint::hasIntent($state);
|
||||
}
|
||||
return (int) $task['upstream_started_at'] === 0;
|
||||
}
|
||||
|
||||
private static function leaseSeconds(): int { return max(30, (int) config('followup_audio.lease_seconds', 600)); }
|
||||
|
||||
private static function opaqueId($value): string
|
||||
|
||||
Reference in New Issue
Block a user