getMessage(), $code) || ($e->errorCode ?? '') === $code, 'expected ' . $code . ', got ' . $e->getMessage()); return; } throw new RuntimeException('Expected ' . $code); }; $f = followupAudioTestDatabase(['asr_chunk_seconds' => 1]); $calls = ['asr' => 0, 'extraction' => 0]; $mode = 'success'; $active = 0; $allTasks = []; $transport = static function (array $spec, callable $heartbeat) use (&$calls, &$mode, &$active, $expect): array { $stage = $spec['stage']; $calls[$stage]++; $task = Store::task($active); $state = Store::open($active, 'pipeline', $task['pipeline_cipher']); $expect(Checkpoint::hasIntent($state) && $task['upstream_started_at'] > 0, 'durable encrypted intent exists BEFORE transport'); if ($stage === 'asr') { $expect(isset($spec['multipart']['file']) && $spec['multipart']['file'] instanceof CURLFile && !isset($spec['json']), 'ASR binary multipart only'); $chunk = $spec['multipart']['file']->getFilename(); $expect(is_file($chunk) && filesize($chunk) > 44 && filesize($chunk) < 40000, 'real bounded FFmpeg WAV chunk'); $expect((fileperms($chunk) & 0777) === 0600, 'private processing chunk'); $index = count(Checkpoint::segments($state)); if ($mode === 'partial' && $index === 1) { return ['http_code' => 429, 'errno' => 0, 'body' => '{}']; } if ($mode === 'unknown') { return ['http_code' => 0, 'errno' => 28, 'body' => '']; } if ($mode === 'revoke') { Db::name('admin')->where('id', 1)->update(['disable' => 1]); } $texts = ['今天早上空腹血糖六点一。', '昨天晚上散步二十分钟。', '结尾已完整转写。']; return ['http_code' => 200, 'errno' => 0, 'body' => json_encode(['text' => $texts[$index]], JSON_UNESCAPED_UNICODE)]; } $expect(!isset($spec['multipart']) && is_string($spec['json']['messages'][0]['content']), 'extraction receives text, never audio_processed claim'); if ($mode === 'extract-reject') { return ['http_code' => 429, 'errno' => 0, 'body' => '{}']; } if ($mode === 'extract-unknown') { return ['http_code' => 502, 'errno' => 0, 'body' => '{}']; } $segment = Checkpoint::segments($state)[0]; $citation = \app\common\service\followupaudio\FollowupAudioTranscriptPrompt::citations(Checkpoint::segments($state))[0]; $expect(($spec['json']['response_format']['type'] ?? '') === 'json_schema' && $spec['json']['max_tokens'] === 8192, 'actual worker sends configured V2 structured request'); $answer = ['schema_version' => 'followup-audio-transcript-v2', 'summary' => '合成随访', 'uncertainties' => [], 'items' => [[ 'kind' => 'blood', 'values' => ['fasting_blood_sugar' => 6.1], 'record_date' => '2026-09-29', 'record_time' => null, 'date_text' => '今天', 'time_text' => '早上', 'time_period' => '早晨', 'time_estimated' => true, 'needs_review' => true, 'evidence_ids' => [$citation['id']], ]]]; if ($mode === 'invalid-evidence') { $answer['items'][0]['evidence_ids'] = ['c_' . str_repeat('0', 24)]; } if ($mode === 'old-version') { $answer['schema_version'] = 'followup-audio-transcript-v1'; unset($answer['items'][0]['evidence_ids']); $answer['items'][0]['evidence'] = [['segment_id' => $segment['id'], 'text' => $segment['text']]]; } $message = ['content' => json_encode($answer, JSON_UNESCAPED_UNICODE), 'reasoning_content' => 'PRIVATE_REASONING_NOT_RETAINED']; if ($mode === 'refusal') { $message['refusal'] = 'synthetic refusal'; } return ['http_code' => 200, 'errno' => 0, 'body' => json_encode(['choices' => [['finish_reason' => $mode === 'length' ? 'length' : 'stop', 'message' => $message]]])]; }; $make = static function () use ($f, &$active, &$allTasks): int { $samples = str_repeat(pack('v', 1000), 48000 - 1) . random_bytes(2); $wav = 'RIFF' . pack('V', 36 + strlen($samples)) . 'WAVEfmt ' . pack('VvvVVvv', 16, 1, 1, 16000, 32000, 2, 16) . 'data' . pack('V', strlen($samples)) . $samples; $source = $f['private'] . '/fixture.wav'; file_put_contents($source, $wav); chmod($source, 0600); $upload = followupAudioTestUpload($f, $source, 'synthetic.wav'); $created = Store::create($upload, '2026-09-29 20:00:00', 'qwen', 1, $f['actor']); $active = $created['id']; $allTasks[] = $active; return $active; }; $run = static function () use ($transport): bool { return (new Worker(new Dify($transport)))->runOnce(); }; try { $expect((int) Db::name('followup_audio_mutex')->count() === 2, 'migration twice is idempotent'); $id = $make(); $expect($run(), 'actual Worker consumes task'); $task = Store::task($id); $detail = Store::detail($id); $expect($detail['status'] === 'review' && $calls === ['asr' => 3, 'extraction' => 1], 'three full chunks then one extraction: ' . json_encode([$detail['status'], $detail['error_code'], $calls])); $expect($detail['pipeline_progress'] === ['transcribed' => 3, 'total' => 3] && $detail['transcript_available'], 'safe public progress'); $expect($detail['transcript'] === "今天早上空腹血糖六点一。\n昨天晚上散步二十分钟。\n结尾已完整转写。", 'canonical full ordered ASR text'); $expect(!str_contains(json_encode(Store::open($id, 'pipeline', $task['pipeline_cipher'])), 'PRIVATE_REASONING_NOT_RETAINED'), 'model reasoning is never persisted in checkpoint'); $expect(!str_contains($task['pipeline_cipher'], '血糖') && !str_contains($task['upstream_ids_json'], '血糖'), 'transcript encrypted, opaque task metadata only'); $expect($detail['items'][0]['evidence'][0]['position_type'] === 'segment' && $detail['items'][0]['evidence'][0]['end_ms'] === 1000, 'server supplied real chunk bounds'); $items = $detail['items']; $items[0]['selected'] = true; $items[0]['needs_review'] = false; $draft = Store::saveDraft($id, $detail['version'], $items, 1); $reject(fn () => Store::saveDraft($id, $detail['version'], $items, 1), 'VERSION_CONFLICT'); $forged = $draft['items']; $forged[0]['evidence'][0]['start_ms'] = 500; $reject(fn () => Store::saveDraft($id, $draft['version'], $forged, 1), 'IMMUTABLE_FIELD'); $applied = Apply::apply($id, $draft['version'], $draft['items'], 1, $f['actor']); $expect($applied['status'] === 'applied' && (int) Db::name('tcm_blood_record')->where('diagnosis_id', 1)->count() === 1, 'review/apply actual local test patient only'); $before = $calls; Apply::apply($id, $draft['version'], $draft['items'], 1, $f['actor']); $expect($calls === $before && (int) Db::name('tcm_blood_record')->count() === 1, 'apply idempotent, no supplier reissue'); $mode = 'partial'; $id = $make(); $before = $calls; $run(); $failed = Store::detail($id); $expect($failed['status'] === 'failed' && $failed['can_retry'] && $failed['pipeline_progress']['transcribed'] === 1, 'known ASR rejection preserves completed chunk'); Store::retry($id); $mode = 'success'; $run(); $expect(Store::detail($id)['status'] === 'review' && $calls['asr'] - $before['asr'] === 4, 'retry sends rejected chunk+tail only, not completed ASR'); $mode = 'extract-reject'; $id = $make(); $run(); $failed = Store::detail($id); $before = $calls; $expect($failed['status'] === 'failed' && $failed['can_retry'] && $failed['transcript_available'] && $failed['transcript'] !== '', 'cached transcript available on extraction failure'); Store::retry($id); $mode = 'success'; $run(); $expect(Store::detail($id)['status'] === 'review' && $calls['asr'] === $before['asr'] && $calls['extraction'] === $before['extraction'] + 1, 'extraction-only explicit retry'); foreach (['unknown', 'extract-unknown'] as $mode) { $id = $make(); $run(); $failed = Store::detail($id); $expect($failed['status'] === 'needs_reconciliation' && !$failed['can_retry'], 'unknown network outcome fenced: ' . $mode); $reject(fn () => Store::retry($id), 'RETRY_NOT_SAFE'); } $mode = 'invalid-evidence'; $id = $make(); $run(); $before = $calls; $expect(Store::detail($id)['error_code'] === 'UPSTREAM_SCHEMA_INVALID', 'fabricated segment evidence rejected'); Store::retry($id); $run(); $expect($calls === $before && Store::detail($id)['error_code'] === 'UPSTREAM_SCHEMA_INVALID', 'known completed invalid extraction validates cached answer without resending'); foreach (['length', 'refusal', 'old-version'] as $mode) { $id = $make(); $run(); $before = $calls; $row = Store::detail($id); $expect($row['status'] === 'failed' && $row['error_code'] === 'UPSTREAM_SCHEMA_INVALID' && $row['transcript_available'], 'length/refusal known failed and ASR cached'); $state = Store::open($id, 'pipeline', Store::task($id)['pipeline_cipher']); $expect($state['extraction']['state'] === 'complete' && ($mode === 'old-version' || $state['extraction']['answer'] === '') && !str_contains(json_encode($state), 'PRIVATE_REASONING_NOT_RETAINED'), 'only final accepted content retained, no reasoning'); Store::retry($id); $run(); $expect($calls === $before, 'length/refusal never silently reissued'); } $mode = 'success'; $id = $make(); $claim = Store::claim(); $task = Store::task($id); $state = Checkpoint::initial($task, Provider::resolve('qwen')['fingerprint'], 1000); $state['revision'] = 1; $expect(Store::checkpoint($id, $claim['lease_token'], ['pipeline_checkpoint' => ['expected_revision' => 0, 'state' => $state]]), 'initial checkpoint persisted'); $intent = $state; $intent['revision'] = 2; $intent['chunks'][0]['state'] = 'intent'; $intent['chunks'][0]['request_id'] = 'fa-' . str_repeat('a', 32); $expect(Store::checkpoint($id, $claim['lease_token'], ['upstream_started_at' => time(), 'pipeline_checkpoint' => ['expected_revision' => 1, 'state' => $intent]]), 'intent CAS'); $reject(fn () => Store::checkpoint($id, $claim['lease_token'], ['pipeline_checkpoint' => ['expected_revision' => 1, 'state' => $intent]]), 'PIPELINE_CHECKPOINT_CONFLICT'); $complete = $intent; $complete['revision'] = 3; $complete['chunks'][0]['state'] = 'complete'; $complete['chunks'][0]['text'] = '今天早上空腹血糖六点一。'; Store::checkpoint($id, $claim['lease_token'], ['pipeline_checkpoint' => ['expected_revision' => 2, 'state' => $complete]]); Db::name('followup_audio_task')->where('id', $id)->update(['lease_until' => time() - 1]); Store::claim(); $expect(Store::detail($id)['status'] === 'failed' && Store::detail($id)['can_retry'], 'crash after persisted completed chunk is safely recoverable'); Store::retry($id); $before = $calls; $run(); $expect(Store::detail($id)['status'] === 'review' && $calls['asr'] - $before['asr'] === 2, 'completed-state recovery retains chunk'); $mode = 'revoke'; $id = $make(); $before = $calls; $run(); Db::name('admin')->where('id', 1)->update(['disable' => 0]); $expect(Store::detail($id)['status'] === 'needs_reconciliation' && $calls['asr'] === $before['asr'] + 1 && $calls['extraction'] === $before['extraction'], 'revoked actor cannot persist response or send next request'); Db::name('admin')->where('id', 1)->update(['disable' => 1]); $reject(fn () => Access::task($id, 1, $f['actor']), '账号已停用'); Db::name('admin')->where('id', 1)->update(['disable' => 0]); foreach ($allTasks as $id) { Db::name('followup_audio_task')->where('id', $id)->update(['expires_at' => time() - 1, 'lease_until' => 0]); Db::name('followup_audio_upload')->where('id', Store::task($id)['upload_id'])->update(['expires_at' => time() - 1]); } $expect(Store::detail($allTasks[2])['transcript'] === '' && !Store::detail($allTasks[2])['transcript_available'], 'retention enforced before cleanup'); $purged = Store::cleanupExpired(); $expect($purged['tasks_purged'] === count($allTasks) && $purged['errors'] === 0, 'all expired encrypted states purged'); foreach ($allTasks as $id) { $row = Store::task($id); $expect($row['pipeline_cipher'] === null && $row['extraction_cipher'] === '', '90d erases new ASR/LLM text'); } echo 'FOLLOWUP_AUDIO_PIPELINE_DATABASE assertions=' . $checks . ' PASS real_mysql=1 actual_upload_worker_review_apply=1 encrypted_checkpoints=1 partial_retry=1 unknown_fenced=1 CAS=1 permission_recheck=1 retention=1' . PHP_EOL; } finally { followupAudioTestDatabaseCleanup($f); }