|
|
@@ -1,8 +1,9 @@
|
|
|
-import { createHash } from 'node:crypto';
|
|
|
+import { randomUUID } from 'node:crypto';
|
|
|
import { mkdir, mkdtemp, rm } from 'node:fs/promises';
|
|
|
import { join } from 'node:path';
|
|
|
import { tmpdir } from 'node:os';
|
|
|
import { AppError } from './errors.mjs';
|
|
|
+import { buildCanonicalTranscript } from './canonical-transcript.mjs';
|
|
|
import { downloadRemoteAudio, probeDurationMs, transcodeToPcmWav } from './media.mjs';
|
|
|
|
|
|
function sleep(ms) {
|
|
|
@@ -14,12 +15,23 @@ function normalizedTranscript(text) {
|
|
|
}
|
|
|
|
|
|
export class TranscriptionWorker {
|
|
|
- constructor({ config, store, provider }) {
|
|
|
+ constructor({
|
|
|
+ config,
|
|
|
+ store,
|
|
|
+ provider,
|
|
|
+ persistence = { enabled: false, async commit() { return { persisted: false }; } },
|
|
|
+ media = { downloadRemoteAudio, probeDurationMs, transcodeToPcmWav },
|
|
|
+ sleepFn = sleep,
|
|
|
+ }) {
|
|
|
this.config = config;
|
|
|
this.store = store;
|
|
|
this.provider = provider;
|
|
|
+ this.persistence = persistence;
|
|
|
+ this.media = media;
|
|
|
+ this.sleep = sleepFn;
|
|
|
this.queue = [];
|
|
|
this.queuedIds = new Set();
|
|
|
+ this.retryTimers = new Map();
|
|
|
this.active = 0;
|
|
|
}
|
|
|
|
|
|
@@ -28,6 +40,22 @@ export class TranscriptionWorker {
|
|
|
}
|
|
|
|
|
|
enqueue(jobId) {
|
|
|
+ const job = this.store.get(jobId);
|
|
|
+ if (!job || ['completed', 'failed', 'cancelled'].includes(job.status)) return;
|
|
|
+ const retryAt = Date.parse(job.nextRetryAt || '');
|
|
|
+ if (job.status === 'retry_wait' && Number.isFinite(retryAt) && retryAt > Date.now()) {
|
|
|
+ if (this.retryTimers.has(jobId)) return;
|
|
|
+ const timer = setTimeout(() => {
|
|
|
+ this.retryTimers.delete(jobId);
|
|
|
+ this.enqueue(jobId);
|
|
|
+ }, retryAt - Date.now());
|
|
|
+ timer.unref?.();
|
|
|
+ this.retryTimers.set(jobId, timer);
|
|
|
+ return;
|
|
|
+ }
|
|
|
+ const timer = this.retryTimers.get(jobId);
|
|
|
+ if (timer) clearTimeout(timer);
|
|
|
+ this.retryTimers.delete(jobId);
|
|
|
if (this.queuedIds.has(jobId)) return;
|
|
|
this.queuedIds.add(jobId);
|
|
|
this.queue.push(jobId);
|
|
|
@@ -37,6 +65,9 @@ export class TranscriptionWorker {
|
|
|
async cancel(jobId) {
|
|
|
const job = this.store.get(jobId);
|
|
|
if (!job || ['completed', 'failed', 'cancelled'].includes(job.status)) return job;
|
|
|
+ const timer = this.retryTimers.get(jobId);
|
|
|
+ if (timer) clearTimeout(timer);
|
|
|
+ this.retryTimers.delete(jobId);
|
|
|
return this.store.update(jobId, {
|
|
|
status: 'cancelled',
|
|
|
stage: 'cancelled',
|
|
|
@@ -45,6 +76,48 @@ export class TranscriptionWorker {
|
|
|
});
|
|
|
}
|
|
|
|
|
|
+ queuePosition(jobId) {
|
|
|
+ const index = this.queue.indexOf(jobId);
|
|
|
+ return index >= 0 ? index + 1 : null;
|
|
|
+ }
|
|
|
+
|
|
|
+ async retry(jobId) {
|
|
|
+ const previous = this.store.get(jobId);
|
|
|
+ if (
|
|
|
+ !previous ||
|
|
|
+ !['failed', 'cancelled'].includes(previous.status) ||
|
|
|
+ !previous.request?.audioUrl
|
|
|
+ ) {
|
|
|
+ throw new AppError('当前任务不能重新执行', {
|
|
|
+ code: 'JOB_NOT_RETRYABLE',
|
|
|
+ status: 409,
|
|
|
+ });
|
|
|
+ }
|
|
|
+ const now = new Date();
|
|
|
+ const job = await this.store.create({
|
|
|
+ id: randomUUID(),
|
|
|
+ ownerKey: previous.ownerKey,
|
|
|
+ idempotencyKey: previous.idempotencyKey,
|
|
|
+ parentJobId: previous.id,
|
|
|
+ status: 'queued',
|
|
|
+ stage: 'queued',
|
|
|
+ progress: 0,
|
|
|
+ attempt: Math.max(0, Number(previous.attempt || 0)),
|
|
|
+ runAttempt: 0,
|
|
|
+ request: previous.request,
|
|
|
+ createdAt: now.toISOString(),
|
|
|
+ updatedAt: now.toISOString(),
|
|
|
+ expiresAt: new Date(now.getTime() + this.config.jobRetentionMs).toISOString(),
|
|
|
+ });
|
|
|
+ console.info('[yuban-server] transcription job manually retried', {
|
|
|
+ jobId: job.id,
|
|
|
+ parentJobId: previous.id,
|
|
|
+ attempt: job.attempt + 1,
|
|
|
+ });
|
|
|
+ this.enqueue(job.id);
|
|
|
+ return job;
|
|
|
+ }
|
|
|
+
|
|
|
async drain() {
|
|
|
while (this.active < this.config.maxJobConcurrency && this.queue.length) {
|
|
|
const jobId = this.queue.shift();
|
|
|
@@ -74,8 +147,24 @@ export class TranscriptionWorker {
|
|
|
async run(jobId) {
|
|
|
let job = this.store.get(jobId);
|
|
|
if (!job || job.status === 'cancelled' || job.status === 'completed') return;
|
|
|
+ const retryAt = Date.parse(job.nextRetryAt || '');
|
|
|
+ if (job.status === 'retry_wait' && Number.isFinite(retryAt) && retryAt > Date.now()) {
|
|
|
+ this.enqueue(jobId);
|
|
|
+ return;
|
|
|
+ }
|
|
|
let workDirectory;
|
|
|
try {
|
|
|
+ job = await this.checkpoint(jobId, {
|
|
|
+ attempt: Math.max(0, Number(job.attempt || 0)) + 1,
|
|
|
+ runAttempt: Math.max(0, Number(job.runAttempt || 0)) + 1,
|
|
|
+ nextRetryAt: null,
|
|
|
+ error: undefined,
|
|
|
+ });
|
|
|
+ console.info('[yuban-server] transcription job started', {
|
|
|
+ jobId,
|
|
|
+ attempt: job.attempt,
|
|
|
+ resumedProviderOrder: Boolean(job.providerOrderId),
|
|
|
+ });
|
|
|
if (!job.providerOrderId) {
|
|
|
await mkdir(this.config.jobTempDir || tmpdir(), { recursive: true });
|
|
|
workDirectory = await mkdtemp(join(this.config.jobTempDir || tmpdir(), `${job.id}-`));
|
|
|
@@ -83,10 +172,10 @@ export class TranscriptionWorker {
|
|
|
const wavPath = join(workDirectory, 'recording.wav');
|
|
|
|
|
|
await this.checkpoint(jobId, { stage: 'downloading', progress: 10 });
|
|
|
- await downloadRemoteAudio(job.request.audioUrl, sourcePath, this.config);
|
|
|
+ await this.media.downloadRemoteAudio(job.request.audioUrl, sourcePath, this.config);
|
|
|
await this.checkpoint(jobId, { stage: 'transcoding', progress: 25 });
|
|
|
- await transcodeToPcmWav(sourcePath, wavPath, this.config);
|
|
|
- const durationMs = await probeDurationMs(wavPath, this.config);
|
|
|
+ await this.media.transcodeToPcmWav(sourcePath, wavPath, this.config);
|
|
|
+ const durationMs = await this.media.probeDurationMs(wavPath, this.config);
|
|
|
await this.checkpoint(jobId, { stage: 'submitting', progress: 40, durationMs });
|
|
|
const submitted = await this.provider.submit(wavPath, {
|
|
|
durationMs,
|
|
|
@@ -111,7 +200,14 @@ export class TranscriptionWorker {
|
|
|
transientFailures = 0;
|
|
|
if (response.status === 'completed') {
|
|
|
const text = normalizedTranscript(response.result.text);
|
|
|
+ job = await this.checkpoint(jobId, { stage: 'persisting', progress: 95 });
|
|
|
+ const persistence = await this.persistence.commit(job, {
|
|
|
+ text,
|
|
|
+ segments: response.result.segments,
|
|
|
+ });
|
|
|
const completedAt = new Date().toISOString();
|
|
|
+ const canonical = persistence?.canonical ||
|
|
|
+ buildCanonicalTranscript(text, response.result.segments, completedAt);
|
|
|
await this.store.update(jobId, {
|
|
|
status: 'completed',
|
|
|
stage: 'completed',
|
|
|
@@ -119,12 +215,21 @@ export class TranscriptionWorker {
|
|
|
heartbeatAt: completedAt,
|
|
|
completedAt,
|
|
|
result: {
|
|
|
- text,
|
|
|
+ text: canonical.transcription,
|
|
|
segments: response.result.segments,
|
|
|
- charCount: text.length,
|
|
|
- sha256: createHash('sha256').update(text, 'utf8').digest('hex'),
|
|
|
+ charCount: canonical.transcriptionCharCount,
|
|
|
+ sha256: canonical.transcriptionHash,
|
|
|
+ canonicalPersisted: Boolean(persistence?.persisted),
|
|
|
},
|
|
|
request: undefined,
|
|
|
+ nextRetryAt: null,
|
|
|
+ error: undefined,
|
|
|
+ });
|
|
|
+ console.info('[yuban-server] transcription job completed', {
|
|
|
+ jobId,
|
|
|
+ attempt: job.attempt,
|
|
|
+ durationMs: job.durationMs || null,
|
|
|
+ charCount: canonical.transcriptionCharCount,
|
|
|
});
|
|
|
return;
|
|
|
}
|
|
|
@@ -138,7 +243,7 @@ export class TranscriptionWorker {
|
|
|
transientFailures += 1;
|
|
|
}
|
|
|
const elapsed = Date.now() - pollStartedAt;
|
|
|
- await sleep(elapsed < 60_000 ? 5_000 : elapsed < 10 * 60_000 ? 15_000 : 30_000);
|
|
|
+ await this.sleep(elapsed < 60_000 ? 5_000 : elapsed < 10 * 60_000 ? 15_000 : 30_000);
|
|
|
}
|
|
|
throw new AppError('长录音转写仍在处理中,请重新查询或稍后重试', {
|
|
|
code: 'PROVIDER_TIMEOUT',
|
|
|
@@ -149,16 +254,47 @@ export class TranscriptionWorker {
|
|
|
const latest = this.store.get(jobId);
|
|
|
if (latest && latest.status !== 'cancelled') {
|
|
|
const appError = error instanceof AppError ? error : new AppError('任务执行失败');
|
|
|
- await this.store.update(jobId, {
|
|
|
- status: 'failed',
|
|
|
- stage: 'failed',
|
|
|
- heartbeatAt: new Date().toISOString(),
|
|
|
- error: {
|
|
|
+ const errorView = {
|
|
|
+ code: appError.code,
|
|
|
+ message: appError.message,
|
|
|
+ retryable: appError.retryable,
|
|
|
+ };
|
|
|
+ const runAttempt = Math.max(1, Number(latest.runAttempt || 1));
|
|
|
+ if (appError.retryable && runAttempt < this.config.maxJobAttempts) {
|
|
|
+ const delayMs = Math.min(
|
|
|
+ 15 * 60_000,
|
|
|
+ this.config.retryBaseDelayMs * 2 ** Math.max(0, runAttempt - 1),
|
|
|
+ );
|
|
|
+ const nextRetryAt = new Date(Date.now() + delayMs).toISOString();
|
|
|
+ await this.store.update(jobId, {
|
|
|
+ status: 'retry_wait',
|
|
|
+ stage: 'retry_wait',
|
|
|
+ heartbeatAt: new Date().toISOString(),
|
|
|
+ nextRetryAt,
|
|
|
+ error: errorView,
|
|
|
+ });
|
|
|
+ console.warn('[yuban-server] transcription job scheduled for retry', {
|
|
|
+ jobId,
|
|
|
+ attempt: latest.attempt,
|
|
|
+ code: appError.code,
|
|
|
+ nextRetryAt,
|
|
|
+ });
|
|
|
+ this.enqueue(jobId);
|
|
|
+ } else {
|
|
|
+ await this.store.update(jobId, {
|
|
|
+ status: 'failed',
|
|
|
+ stage: 'failed',
|
|
|
+ heartbeatAt: new Date().toISOString(),
|
|
|
+ nextRetryAt: null,
|
|
|
+ error: errorView,
|
|
|
+ });
|
|
|
+ console.error('[yuban-server] transcription job failed', {
|
|
|
+ jobId,
|
|
|
+ attempt: latest.attempt,
|
|
|
code: appError.code,
|
|
|
- message: appError.message,
|
|
|
retryable: appError.retryable,
|
|
|
- },
|
|
|
- });
|
|
|
+ });
|
|
|
+ }
|
|
|
}
|
|
|
} finally {
|
|
|
if (workDirectory) await rm(workDirectory, { recursive: true, force: true });
|