feat: configure audio providers and add local sample acceptance

This commit is contained in:
2026-10-07 18:15:32 +08:00
parent 70be2fc70f
commit fea33285e8
21 changed files with 1103 additions and 131 deletions
@@ -4,7 +4,7 @@ declare(strict_types=1);
namespace app\common\service\followupaudio;
/** Dedicated, fail-closed local_file audio transport; intentionally does not use DifyChatService. */
/** Fail-closed audio transport; explicit Dify/OpenAI driver, never text or provider fallback. */
final class FollowupAudioDify
{
private array $settings;
@@ -24,7 +24,7 @@ final class FollowupAudioDify
if (empty($this->settings['enabled'])) {
throw new FollowupAudioException('FEATURE_DISABLED');
}
if (empty($this->settings['audio_verified']) || !in_array((string) ($task['model_key'] ?? ''), (array) ($this->settings['verified_profiles'] ?? []), true)) {
if (!FollowupAudioProviderConfig::verified((string) ($task['model_key'] ?? ''), $this->settings, $this->provider)) {
throw new FollowupAudioException('AUDIO_NOT_VERIFIED');
}
if ((int) ($task['upstream_started_at'] ?? 0) > 0) {
@@ -53,25 +53,28 @@ final class FollowupAudioDify
'upload_id' => 'synthetic', 'sha256' => $fixture['sha256'],
'recorded_at' => '2026-09-29 10:00:00', 'model_key' => $profile,
'duration_seconds' => (float) ($fixture['duration_seconds'] ?? 0),
], $heartbeat);
], $heartbeat, true);
}
/** Returns status and an irreversible app-identity fingerprint, never configuration values/credentials. */
public function configurationStatus(string $profile): array
{
try {
[$base, $key] = $this->resolveProfile($profile);
return ['configured' => true, 'profile' => $profile, 'code' => 'OK',
'application_fingerprint' => hash('sha256', $base . "\0" . $key)];
} catch (FollowupAudioException $e) {
return ['configured' => false, 'profile' => in_array($profile, ['qwen', 'openai'], true) ? $profile : 'invalid',
'code' => $e->errorCode];
}
return FollowupAudioProviderConfig::status($profile, $this->settings, $this->provider);
}
private function analyzeFile(string $path, array $task, callable $heartbeat): array
private function analyzeFile(string $path, array $task, callable $heartbeat, bool $synthetic = false): array
{
[$base, $key, $timeout] = $this->resolveProfile((string) ($task['model_key'] ?? ''));
[$base, $key, $timeout, $provider] = $this->resolveProfile((string) ($task['model_key'] ?? ''));
if (!str_starts_with($base, 'https://') && (!$synthetic || empty($this->settings['allow_insecure_synthetic']))) {
throw new FollowupAudioException('HTTPS_REQUIRED');
}
if (!$synthetic) {
$snapshot = json_decode((string) ($task['upstream_ids_json'] ?? ''), true);
if (!is_array($snapshot) || !is_string($snapshot['provider_fingerprint'] ?? null)
|| !hash_equals($provider['fingerprint'], $snapshot['provider_fingerprint'])) {
throw new FollowupAudioException('PROVIDER_CONFIGURATION_CHANGED');
}
}
$audio = $this->inspectAudio($path, (string) ($task['sha256'] ?? ''));
if (isset($task['duration_seconds']) && abs((float) $task['duration_seconds'] - $audio['duration']) > 1.0) {
throw new FollowupAudioException('AUDIO_INVALID');
@@ -81,36 +84,52 @@ final class FollowupAudioDify
if (!$date || $date->format('Y-m-d H:i:s') !== $recordedAt) {
throw new FollowupAudioException('RECORDED_AT_INVALID');
}
$processing = $this->prepareAudio($audio, $heartbeat);
$processing = $this->prepareAudio($audio, $heartbeat, $provider['driver'] === 'openai_audio');
try {
$query = $this->prompt($recordedAt, $audio['duration']);
$ids = ['request_id' => 'fa-' . bin2hex(random_bytes(16))];
$ids = ['request_id' => 'fa-' . bin2hex(random_bytes(16)), 'provider_fingerprint' => $provider['fingerprint']];
$user = 'followup-audio-' . substr(hash('sha256', (string) ($task['id'] ?? '') . ':' . $ids['request_id']), 0, 32);
// Persist intent BEFORE any network side effect, including upload. A crashed worker cannot resend.
$this->checkpoint($heartbeat, ['stage' => 'uploading', 'upstream_started_at' => time(),
'upstream_ids_json' => json_encode($ids, JSON_THROW_ON_ERROR)], false);
$uploaded = $this->request([
'url' => $base . '/files/upload', 'api_key' => $key, 'timeout' => $timeout, 'request_ids' => $ids,
'multipart' => ['user' => $user, 'file' => new \CURLFile($processing['path'], $processing['mime'], 'followup-audio.' . $processing['extension'])],
], $heartbeat);
$fileId = $this->identifier($uploaded['id'] ?? null);
if ($fileId === '') {
throw new FollowupAudioException('UPSTREAM_UPLOAD_INVALID', true);
if ($provider['driver'] === 'openai_audio') {
// Build and validate the complete request before persisting the single billable intent.
$payload = $this->openAiPayload($processing, $query, $provider['model']);
$this->checkpoint($heartbeat, ['stage' => 'analyzing', 'upstream_started_at' => time(),
'upstream_ids_json' => json_encode($ids, JSON_THROW_ON_ERROR)], false);
$response = $this->request(['url' => $base . '/chat/completions', 'api_key' => $key,
'timeout' => $timeout, 'json' => $payload, 'request_ids' => $ids], $heartbeat);
$response['message_id'] = $this->identifier($response['id'] ?? null);
$choices = is_array($response['choices'] ?? null) ? $response['choices'] : [];
$choice = is_array($choices[0] ?? null) ? $choices[0] : [];
$choice['message'] = is_array($choice['message'] ?? null) ? $choice['message'] : [];
$response['answer'] = count($choices) === 1
&& ($choice['finish_reason'] ?? '') === 'stop' && empty($choice['message']['refusal'])
&& empty($choice['message']['tool_calls']) ? ($choice['message']['content'] ?? null) : null;
} else {
// Persist intent BEFORE any network side effect, including upload. A crashed worker cannot resend.
$this->checkpoint($heartbeat, ['stage' => 'uploading', 'upstream_started_at' => time(),
'upstream_ids_json' => json_encode($ids, JSON_THROW_ON_ERROR)], false);
$uploaded = $this->request([
'url' => $base . '/files/upload', 'api_key' => $key, 'timeout' => $timeout, 'request_ids' => $ids,
'multipart' => ['user' => $user, 'file' => new \CURLFile($processing['path'], $processing['mime'], 'followup-audio.' . $processing['extension'])],
], $heartbeat);
$fileId = $this->identifier($uploaded['id'] ?? null);
if ($fileId === '') {
throw new FollowupAudioException('UPSTREAM_UPLOAD_INVALID', true);
}
// Preserve the observed file ID even if the upstream mislabeled/rejected its media type.
$this->checkpoint($heartbeat, ['stage' => 'analyzing', 'upstream_file_id' => $fileId], true);
if (isset($uploaded['mime_type']) && (!is_string($uploaded['mime_type']) || !str_starts_with($uploaded['mime_type'], 'audio/'))) {
throw new FollowupAudioException('UPSTREAM_AUDIO_REJECTED');
}
$payload = [
'inputs' => new \stdClass(),
'query' => $query,
'response_mode' => 'blocking', 'user' => $user, 'auto_generate_name' => false,
'files' => [['type' => 'audio', 'transfer_method' => 'local_file', 'upload_file_id' => $fileId]],
];
self::assertAudioPayload($payload, $fileId);
$response = $this->request(['url' => $base . '/chat-messages', 'api_key' => $key,
'timeout' => $timeout, 'json' => $payload, 'request_ids' => $ids], $heartbeat);
}
// Preserve the observed file ID even if the upstream mislabeled/rejected its media type.
$this->checkpoint($heartbeat, ['stage' => 'analyzing', 'upstream_file_id' => $fileId], true);
if (isset($uploaded['mime_type']) && (!is_string($uploaded['mime_type']) || !str_starts_with($uploaded['mime_type'], 'audio/'))) {
throw new FollowupAudioException('UPSTREAM_AUDIO_REJECTED');
}
$payload = [
'inputs' => new \stdClass(),
'query' => $query,
'response_mode' => 'blocking', 'user' => $user, 'auto_generate_name' => false,
'files' => [['type' => 'audio', 'transfer_method' => 'local_file', 'upload_file_id' => $fileId]],
];
self::assertAudioPayload($payload, $fileId);
$response = $this->request(['url' => $base . '/chat-messages', 'api_key' => $key,
'timeout' => $timeout, 'json' => $payload, 'request_ids' => $ids], $heartbeat);
foreach (['task_id', 'message_id', 'conversation_id', 'upstream_request_id'] as $name) {
$id = $this->identifier($response[$name] ?? null);
if ($id !== '') { $ids[$name] = $id; }
@@ -156,38 +175,25 @@ final class FollowupAudioDify
private function resolveProfile(string $profile): array
{
if (!in_array($profile, ['qwen', 'openai'], true)) {
throw new FollowupAudioException('INVALID_PROFILE');
}
$base = rtrim((string) ($this->provider['base_url'] ?? ''), '/');
$key = (string) ($this->provider['models'][$profile]['api_key'] ?? '');
if ($base === '' || trim($key) === '') {
throw new FollowupAudioException('CONFIG_MISSING');
}
$parts = parse_url($base);
if (!is_array($parts) || !in_array($parts['scheme'] ?? '', ['https', 'http'], true)
|| empty($parts['host']) || isset($parts['user']) || isset($parts['pass']) || isset($parts['query'])
|| isset($parts['fragment']) || preg_match('/[\x00-\x20\x7f]/', $base)
|| preg_match('/[\x00-\x20\x7f]/', $key)) {
throw new FollowupAudioException('CONFIG_INVALID');
}
$path = (string) ($parts['path'] ?? '');
if (str_ends_with($path, '/chat/completions')) {
throw new FollowupAudioException('DIFY_APPLICATION_REQUIRED');
}
if (str_ends_with($path, '/chat-messages')) {
$base = substr($base, 0, -strlen('/chat-messages'));
} elseif (!str_ends_with($path, '/v1')) {
$base .= '/v1';
}
$provider = FollowupAudioProviderConfig::resolve($profile, $this->settings, $this->provider);
$timeout = (int) ($this->settings['request_timeout'] ?? 240);
if ($timeout < 1 || $timeout > 300) {
throw new FollowupAudioException('CONFIG_INVALID');
if ($timeout < 1 || $timeout > 300) { throw new FollowupAudioException('CONFIG_INVALID'); }
if (!function_exists('curl_init')) { throw new FollowupAudioException('CURL_UNAVAILABLE'); }
return [$provider['base_url'], $provider['api_key'], $timeout, $provider];
}
/** Only MP3/WAV input_audio is supported, with bytes from the verified complete processing copy. */
private function openAiPayload(array $audio, string $query, string $model): array
{
if (!in_array($audio['extension'], ['mp3', 'wav'], true)) { throw new FollowupAudioException('AUDIO_INVALID'); }
$bytes = file_get_contents($audio['path']);
if (!is_string($bytes) || $bytes === '' || strlen($bytes) > (int) ($this->settings['upstream_max_bytes'] ?? 20971520)) {
throw new FollowupAudioException('UPSTREAM_AUDIO_LIMIT');
}
if (!function_exists('curl_init')) {
throw new FollowupAudioException('CURL_UNAVAILABLE');
}
return [$base, $key, $timeout];
return ['model' => $model, 'stream' => false, 'messages' => [['role' => 'user', 'content' => [
['type' => 'text', 'text' => $query],
['type' => 'input_audio', 'input_audio' => ['data' => base64_encode($bytes), 'format' => $audio['extension']]],
]]]];
}
private function inspectAudio(string $path, string $sha256, bool $processing = false): array
@@ -240,7 +246,7 @@ final class FollowupAudioDify
}
/** Preserve the exact original. Only a private, bounded audio copy may be sent to the same Dify app. */
private function prepareAudio(array $audio, callable $heartbeat): array
private function prepareAudio(array $audio, callable $heartbeat, bool $openAi = false): array
{
// Dify has a separate audio-upload ceiling (default 50 MiB); use a conservative 20 MiB local budget.
// Deployment must verify its own app/model limits via the synthetic probe before enabling a profile.
@@ -249,7 +255,8 @@ final class FollowupAudioDify
if ($limit < 1 || $limit > 52428800 || $timeout < 1 || $timeout > 600) {
throw new FollowupAudioException('CONFIG_INVALID');
}
if (filesize($audio['path']) <= $limit && $audio['extension'] !== 'amr') { return $audio; }
if (filesize($audio['path']) <= $limit && $audio['extension'] !== 'amr'
&& (!$openAi || in_array($audio['extension'], ['mp3', 'wav'], true))) { return $audio; }
// 32 kbit/s mono speech preserves the entire hour within ~14.5 MB, without splitting model requests.
if ($audio['duration'] * 4000 + 2048 > $limit) {
throw new FollowupAudioException('UPSTREAM_AUDIO_LIMIT');
@@ -416,8 +423,9 @@ final class FollowupAudioDify
$http = (int) ($response['http_code'] ?? 0);
// Even failed/uncertain replies can contain a billable task ID. Persist it before raising a sanitized error.
$candidate = json_decode((string) ($response['body'] ?? ''), true);
if (is_array($candidate) && isset($candidate['id'], $spec['json']['messages'])) { $candidate['message_id'] = $candidate['id']; }
if ((int) ($response['errno'] ?? 0) !== 0 || $http < 200 || $http >= 300
|| !is_array($candidate) || !empty($candidate['code']) || ($candidate['event'] ?? '') === 'error') {
|| !is_array($candidate) || !empty($candidate['code']) || !empty($candidate['error']) || ($candidate['event'] ?? '') === 'error') {
$observed = $spec['request_ids'] ?? [];
foreach (['task_id', 'message_id', 'conversation_id'] as $name) {
$id = $this->identifier($candidate[$name] ?? null);
@@ -441,7 +449,7 @@ final class FollowupAudioDify
try { $body = json_decode((string) ($response['body'] ?? ''), true, 64, JSON_THROW_ON_ERROR); }
catch (\Throwable $e) { throw new FollowupAudioException('UPSTREAM_UNCERTAIN', true); }
if (!is_array($body)) { throw new FollowupAudioException('UPSTREAM_UNCERTAIN', true); }
if (!empty($body['code']) || ($body['event'] ?? '') === 'error') {
if (!empty($body['code']) || !empty($body['error']) || ($body['event'] ?? '') === 'error') {
throw new FollowupAudioException('UPSTREAM_UNCERTAIN', true);
}
$requestId = $this->identifier($response['request_id'] ?? null);
@@ -15,14 +15,16 @@ final class FollowupAudioException extends \RuntimeException
$this->errorCode = $errorCode;
$this->uncertain = $uncertain;
$messages = [
'CONFIG_MISSING' => '当前 Dify 应用凭据未配置,音频能力尚未验证',
'CONFIG_MISSING' => '当前音频模型服务未完整配置,音频能力尚未验证',
'HTTPS_REQUIRED' => '真实音频只允许通过有效 HTTPS 服务传输',
'PROVIDER_CONFIGURATION_CHANGED' => '音频模型配置已变化或缺少绑定,任务未发送,请重新核对',
'FEATURE_DISABLED' => '随访音频功能未启用',
'AUDIO_NOT_VERIFIED' => '当前应用尚未通过音频能力验证',
'AUDIO_NOT_PROCESSED' => '上游未确认读取原始音频,未生成可采用结果',
'UPSTREAM_UNCERTAIN' => '上游结果未知,请先核对任务或计费,禁止重复提交',
'RECONCILIATION_REQUIRED' => '任务已有上游请求,须先人工核对,禁止重复提交',
'UPSTREAM_SCHEMA_INVALID' => '上游未返回完整、可核验的结构化音频事实',
'UPSTREAM_AUDIO_REJECTED' => '当前 Dify 应用拒绝音频附件,未降级为纯文本',
'UPSTREAM_AUDIO_REJECTED' => '当前模型服务拒绝音频附件,未降级为纯文本',
'LEASE_LOST' => '任务租约或访问权限已失效',
'ACCESS_REVOKED' => '执行权限已撤销',
'UPSTREAM_AUDIO_LIMIT' => '录音处理副本仍超出已配置的上游限制,未发送,请联系管理员',
@@ -0,0 +1,86 @@
<?php
declare(strict_types=1);
namespace app\common\service\followupaudio;
/** Server-only provider identity. Changing any request identity invalidates the audio gate. */
final class FollowupAudioProviderConfig
{
public static function resolve(string $profile, ?array $settings = null, ?array $legacy = null): array
{
if (!in_array($profile, ['qwen', 'openai'], true)) { throw new FollowupAudioException('INVALID_PROFILE'); }
$settings = $settings ?? (array) config('followup_audio', []);
$legacy = $legacy ?? (array) config('prescription_ai', []);
$provider = (array) ($settings['providers'][$profile] ?? []);
$driver = (string) ($provider['driver'] ?? 'dify');
if (!in_array($driver, ['dify', 'openai_audio'], true)) { throw new FollowupAudioException('CONFIG_INVALID'); }
$base = (string) ($provider['base_url'] ?? '');
$key = (string) ($provider['api_key'] ?? '');
if ($driver === 'dify') {
if ($base === '') { $base = (string) ($legacy['base_url'] ?? ''); }
if ($key === '') { $key = (string) ($legacy['models'][$profile]['api_key'] ?? ''); }
}
$model = (string) ($provider['model'] ?? '');
$label = trim((string) ($provider['label'] ?? ''));
if ($label === '') { $label = $model !== '' ? $model : $profile; }
if ($base === '' || $key === '' || ($driver === 'openai_audio' && $model === '')) { throw new FollowupAudioException('CONFIG_MISSING'); }
$parts = parse_url($base);
if (!is_array($parts) || !in_array($parts['scheme'] ?? '', ['http', 'https'], true)
|| empty($parts['host']) || isset($parts['user']) || isset($parts['pass']) || isset($parts['query']) || isset($parts['fragment'])
|| preg_match('/[\x00-\x20\x7f\\\\]/', $base) || preg_match('/[\x00-\x20\x7f]/', $key)
|| strlen($model) > 200 || preg_match('/[\x00-\x20\x7f]/', $model)
|| trim($label) === '' || strlen($label) > 200 || preg_match('/[\x00-\x1f\x7f]/', $label)) {
throw new FollowupAudioException('CONFIG_INVALID');
}
$base = rtrim($base, '/');
$path = rtrim((string) ($parts['path'] ?? ''), '/');
if ($driver === 'dify' && str_ends_with($path, '/chat/completions')) {
throw new FollowupAudioException('DIFY_APPLICATION_REQUIRED');
}
if ($driver === 'openai_audio' && str_ends_with($path, '/chat-messages')) {
throw new FollowupAudioException('CONFIG_INVALID');
}
$endpoint = $driver === 'dify' ? '/chat-messages' : '/chat/completions';
if (str_ends_with($base, $endpoint)) { $base = substr($base, 0, -strlen($endpoint)); }
elseif (!str_ends_with($base, '/v1')) { $base .= '/v1'; }
$fingerprint = hash('sha256', json_encode([$driver, $base, $model, $key], JSON_UNESCAPED_SLASHES | JSON_THROW_ON_ERROR));
return ['driver' => $driver, 'base_url' => $base, 'api_key' => $key, 'model' => $model, 'label' => $label, 'fingerprint' => $fingerprint];
}
/** No endpoint, model credential or plaintext provider data in capability responses. */
public static function status(string $profile, ?array $settings = null, ?array $legacy = null): array
{
try {
$provider = self::resolve($profile, $settings, $legacy);
return ['configured' => true, 'profile' => $profile, 'code' => 'OK', 'application_fingerprint' => $provider['fingerprint']];
} catch (FollowupAudioException $error) {
return ['configured' => false, 'profile' => in_array($profile, ['qwen', 'openai'], true) ? $profile : 'invalid', 'code' => $error->errorCode];
}
}
public static function isVerified(string $profile): bool
{
return self::verified($profile, (array) config('followup_audio', []), (array) config('prescription_ai', []));
}
/** Shared with injected-config transport tests; enabled remains an independent operational switch. */
public static function verified(string $profile, array $settings, array $legacy): bool
{
if (empty($settings['audio_verified']) || !in_array($profile, (array) ($settings['verified_profiles'] ?? []), true)) { return false; }
try { $provider = self::resolve($profile, $settings, $legacy); }
catch (FollowupAudioException $error) { return false; }
$verified = (string) ($settings['providers'][$profile]['verified_fingerprint'] ?? '');
return str_starts_with($provider['base_url'], 'https://') && preg_match('/^[a-f0-9]{64}$/D', $verified)
&& hash_equals($provider['fingerprint'], $verified);
}
public static function publicModels(): array
{
$models = [];
foreach (['qwen', 'openai'] as $profile) {
if (self::isVerified($profile)) { $models[] = ['value' => $profile, 'label' => self::resolve($profile)['label']]; }
}
return $models;
}
}
@@ -14,13 +14,30 @@ 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));
return $profile === null ? self::verifiedProfiles() !== [] : FollowupAudioProviderConfig::isVerified($profile);
}
private static function verifiedProfiles(): array
{
return array_values(array_intersect(['qwen', 'openai'], (array) config('followup_audio.verified_profiles', [])));
return array_values(array_filter(['qwen', 'openai'], [FollowupAudioProviderConfig::class, 'isVerified']));
}
/** A task binds to the exact provider/endpoint/model/key identity approved when it was created. */
public static function providerMatches(array $task): bool
{
$ids = json_decode((string) ($task['upstream_ids_json'] ?? '{}'), true);
$fingerprint = is_array($ids) ? ($ids['provider_fingerprint'] ?? '') : '';
if (!is_string($fingerprint) || !preg_match('/^[a-f0-9]{64}$/D', $fingerprint)) { return false; }
$status = FollowupAudioProviderConfig::status((string) ($task['model_key'] ?? ''));
return !empty($status['configured']) && is_string($status['application_fingerprint'] ?? null)
&& hash_equals($fingerprint, $status['application_fingerprint']);
}
public static function assertTaskProvider(array $task): void
{
if (!self::providerMatches($task)) {
throw new FollowupAudioException('PROVIDER_CONFIGURATION_CHANGED', (int) ($task['upstream_started_at'] ?? 0) > 0);
}
}
public static function assertEnabled(?string $profile = null): void
@@ -69,7 +86,9 @@ final class FollowupAudioStore
'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' => '{}',
'upstream_run_id' => '', 'upstream_file_id' => '', 'upstream_ids_json' => FollowupAudioPolicy::canonical([
'provider_fingerprint' => FollowupAudioProviderConfig::resolve($modelKey)['fingerprint'],
]),
'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,
]);
@@ -91,6 +110,9 @@ final class FollowupAudioStore
private static function reuse(array $task, string $recordedAt): array
{
$message = '同一诊单、录音内容及模型已有任务,已复用原任务,不会再次调用模型。';
if (!self::providerMatches($task)) {
$message .= '原任务的服务配置已变更或缺少配置快照,不能自动改用新服务重跑。';
}
if ($task['recorded_at'] !== $recordedAt) {
$message .= '沿用原任务的录音时间,请在未采用的审阅项中更正记录日期;已采用任务不会重跑。';
}
@@ -127,7 +149,7 @@ final class FollowupAudioStore
$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;
&& self::providerMatches($task) && (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;
}
@@ -192,6 +214,7 @@ final class FollowupAudioStore
Db::transaction(static function () use ($taskId): void {
$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
|| (int) $task['expires_at'] <= time() || (int) $task['purged_at']) {
throw new DomainException('FOLLOWUP_AUDIO_RETRY_NOT_SAFE');
@@ -207,7 +230,7 @@ final class FollowupAudioStore
/** A singleton DB mutex enforces concurrency across independent CLI processes. */
public static function claim(): ?array
{
if (!self::enabled() || !self::verified()) { return null; }
if (!self::enabled()) { 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');
@@ -227,9 +250,22 @@ 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; }
$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; }
$pending = Db::name('followup_audio_task')->where('status', 'queued')->where('upstream_started_at', 0)
->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::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([
'status' => 'failed', 'stage' => 'failed', 'error_code' => 'PROVIDER_CONFIGURATION_CHANGED',
'error_message' => '任务服务配置已变更或缺少配置快照,已停止处理;恢复原配置并核对后再操作',
'updated_at' => $now, 'version' => (int) $candidate['version'] + 1,
]);
continue;
}
if (self::verified((string) $candidate['model_key'])) { $task = $candidate; break; }
}
if ($task === null) { 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];
@@ -269,6 +305,13 @@ final class FollowupAudioStore
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 ($name === 'provider_fingerprint') {
if (!is_string($identifier) || !is_string($previous[$name] ?? null)
|| !hash_equals($previous[$name], $identifier)) {
throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID');
}
continue; // The immutable task snapshot is retained, never replaced by a callback.
}
if (!in_array($name, ['request_id', 'task_id', 'message_id', 'conversation_id', 'upstream_request_id'], true)) {
throw new DomainException('FOLLOWUP_AUDIO_CHECKPOINT_INVALID');
}
@@ -313,13 +356,17 @@ final class FollowupAudioStore
{
return Db::transaction(static function () use ($id, $token, $code, $uncertain): bool {
$task = self::lockTask($id);
if (!self::hasLease($task, $token)) { return false; }
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;
$status = $uncertain ? 'needs_reconciliation' : 'failed';
Db::name('followup_audio_task')->where('id', $id)->update([
'status' => $status, 'stage' => $status, 'error_code' => $code,
'error_message' => $uncertain ? '上游结果待核对,禁止重复提交' : '处理未完成,请查看错误代码;不会自动重复调用',
'error_message' => $uncertain ? '上游结果待核对,禁止重复提交'
: ($code === 'PROVIDER_CONFIGURATION_CHANGED'
? '任务服务配置已变更或缺少配置快照,已停止处理;恢复原配置并核对后再操作'
: '处理未完成,请查看错误代码;不会自动重复调用'),
'lease_token' => '', 'lease_until' => 0, 'updated_at' => time(), 'version' => (int) $task['version'] + 1,
]);
return true;
@@ -422,9 +469,10 @@ final class FollowupAudioStore
return new PrescriptionAiCipher($key === '' ? null : $key);
}
private static function hasLease(array $task, string $token): bool
private static function hasLease(array $task, string $token, bool $processing = true): bool
{
return self::enabled() && self::verified((string) $task['model_key']) && $task['status'] === 'running' && $token !== ''
return (!$processing || self::enabled() && self::verified((string) $task['model_key']) && self::providerMatches($task))
&& $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'];
}
@@ -17,15 +17,17 @@ final class FollowupAudioWorker
public function runOnce(): bool
{
if (!FollowupAudioStore::enabled() || !FollowupAudioStore::verified() || !($task = FollowupAudioStore::claim())) {
if (!FollowupAudioStore::enabled() || !($task = FollowupAudioStore::claim())) {
return false;
}
$id = (int) $task['id'];
$token = (string) $task['lease_token'];
$started = (int) ($task['upstream_started_at'] ?? 0) > 0;
try {
FollowupAudioStore::assertTaskProvider($task);
if (!FollowupAudioStore::verified((string) $task['model_key'])) { throw new FollowupAudioException('AUDIO_NOT_VERIFIED'); }
$heartbeat = function (array $fields = []) use ($task, $id, $token, &$started): bool {
FollowupAudioStore::assertTaskProvider($task);
if (!FollowupAudioStore::enabled() || !FollowupAudioStore::verified((string) $task['model_key'])) { return false; }
$actor = PrescriptionAiAccess::actor((int) $task['actor_id']);
if (!$actor) { return false; }
@@ -42,7 +44,7 @@ final class FollowupAudioWorker
throw new FollowupAudioException('LEASE_LOST', true);
}
} catch (FollowupAudioException $e) {
FollowupAudioStore::fail($id, $token, $e->errorCode, $e->getMessage(), $e->uncertain);
FollowupAudioStore::fail($id, $token, $e->errorCode, $e->getMessage(), $e->uncertain || $started);
} catch (\Throwable $e) {
// The exception may contain SQL, names, transcript, credentials or a signed URL.
FollowupAudioStore::fail($id, $token, 'INTERNAL_ERROR', '音频任务执行异常,请核对任务状态', $started);