feat(speech): 支持同步Flash ASR模型和异步文件转录模型

- 新增同步Flash ASR模型请求流程,支持单音频文件识别
- 实现了对同步Flash模型不支持异步标志及参数的限制校验
- 异步文件转录模型支持单文件URL上传和语言参数细化
- 优化异步和同步鉴权域显示,丰富根帮助和分组帮助提示
- speech recognize增加dry-run测试覆盖多种识别场景
- free-tier自动停用功能优化,改用统一轮询函数处理批量请求
- 统一轮询逻辑,支持console和telemetry接口的异步任务完成判定
- 规范输出格式和错误提示,增强用户调试体验
- 版本升级到1.14.3,更新示例参数和模型ID引用
This commit is contained in:
zeyu.fz
2026-08-14 11:07:21 +08:00
64 changed files with 2776 additions and 1451 deletions
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "bailian-cli-runtime",
"version": "1.14.2",
"version": "1.14.3",
"description": "Runtime framework for bailian-cli (createCli, registry, args, output, pipeline). See https://www.npmjs.com/package/bailian-cli for usage.",
"homepage": "https://bailian.console.aliyun.com/cli",
"bugs": {
+180 -14
View File
@@ -10,6 +10,11 @@ import {
taskPath,
speechSynthesizePath,
speechRecognizePath,
resolveAsrApi,
buildAsrFlashRequest,
buildAsyncAsrLanguageFields,
collectAsrTranscriptionItems,
extractAsrFlashText,
stripUndefined,
resolveBooleanFlag,
resolveWatermark,
@@ -23,6 +28,7 @@ import {
type DashScopeTTSRequest,
type DashScopeTTSResponse,
type DashScopeASRRequest,
type DashScopeASRTaskResult,
type ChatMessageContent,
isLocalFile,
} from "bailian-cli-core";
@@ -573,27 +579,103 @@ export async function speechRecognize(
});
}
const model = input.model || "fun-asr";
const route = resolveAsrApi(model);
if (route.kind === "unsupported") {
throw new PipelineError(
"invalid_input",
route.unsupportedReason ?? `Unsupported ASR model: ${model}`,
{
step: "speech/recognize",
},
);
}
if (route.kind === "sync-flash") {
if (rawUrls.length !== 1) {
throw new PipelineError(
"invalid_input",
`Model "${model}" is a sync Flash ASR model and accepts exactly one url (got ${rawUrls.length})`,
{ step: "speech/recognize" },
);
}
const unsupportedFlags: string[] = [];
if (input.diarization) unsupportedFlags.push("diarization");
if (input["speaker-count"] !== undefined) unsupportedFlags.push("speaker-count");
// input-audio Flash supports vocabulary_id; qwen3 sync Flash does not
if (route.flashFamily === "qwen3" && input["vocabulary-id"] !== undefined) {
unsupportedFlags.push("vocabulary-id");
}
if (input["channel-id"] !== undefined) unsupportedFlags.push("channel-id");
if (unsupportedFlags.length > 0) {
throw new PipelineError(
"invalid_input",
`Model "${model}" uses sync Flash ASR and does not support: ${unsupportedFlags.join(", ")}`,
{ step: "speech/recognize" },
);
}
}
if (
route.kind === "async-filetrans" &&
route.asyncInputStyle === "file_url" &&
rawUrls.length !== 1
) {
throw new PipelineError(
"invalid_input",
`Model "${model}" accepts exactly one url (got ${rawUrls.length})`,
{ step: "speech/recognize" },
);
}
// Resolve local files to upload URLs
const fileUrls: string[] = [];
for (const u of rawUrls) {
if (isLocalFile(u)) {
for (const audioUrl of rawUrls) {
if (isLocalFile(audioUrl)) {
fileUrls.push(
await env.client.uploadFile(u, input.model || "fun-asr", {
await env.client.uploadFile(audioUrl, model, {
signal: ctx.signal,
}),
);
} else {
fileUrls.push(u);
fileUrls.push(audioUrl);
}
}
const model = input.model || "fun-asr";
if (route.kind === "sync-flash") {
const flashFamily = route.flashFamily!;
const body = buildAsrFlashRequest({
model,
audioUrl: fileUrls[0]!,
language: input.language,
vocabularyId: input["vocabulary-id"],
flashFamily,
});
const response = await env.client.requestJson<Record<string, unknown>>({
path: route.path,
method: "POST",
headers: { "X-DashScope-SSE": "disable" },
body,
signal: ctx.signal,
});
return {
text: extractAsrFlashText(response, flashFamily),
model,
mode: "sync",
raw: response,
};
}
const languageFields = buildAsyncAsrLanguageFields(
route.asyncLanguageStyle ?? "language_hints",
input.language,
);
const body: DashScopeASRRequest = {
model,
input: { file_urls: fileUrls },
input:
route.asyncInputStyle === "file_url" ? { file_url: fileUrls[0]! } : { file_urls: fileUrls },
parameters: {
channel_id: input["channel-id"] !== undefined ? [input["channel-id"]] : undefined,
language_hints: input.language ? [input.language] : undefined,
...languageFields,
diarization_enabled: input.diarization,
speaker_count: input["speaker-count"],
vocabulary_id: input["vocabulary-id"],
@@ -601,9 +683,8 @@ export async function speechRecognize(
};
stripUndefined(body.parameters as Record<string, unknown>);
const url = speechRecognizePath();
const asyncResp = await env.client.requestJson<DashScopeAsyncResponse>({
path: url,
path: speechRecognizePath(),
method: "POST",
body,
async: true,
@@ -614,7 +695,65 @@ export async function speechRecognize(
const pollIntervalMs = (input["poll-interval"] ?? 2) * 1000;
const timeoutMs = (ctx.timeoutSeconds ?? 300) * 1000;
return await pollTaskWithOptions(env, taskId, pollIntervalMs, timeoutMs, ctx);
// ASR polling reads original output, avoids generic flatten (avoids transcription_url polluting media urls)
const asrTask = await pollAsrTaskWithOptions(env, taskId, pollIntervalMs, timeoutMs, ctx);
const transcriptionItems = collectAsrTranscriptionItems(asrTask.output);
const base: Record<string, unknown> = {
task_id: asrTask.output.task_id,
task_status: asrTask.output.task_status,
request_id: asrTask.request_id,
mode: "async",
model,
};
if (asrTask.output.results) base.results = asrTask.output.results;
if (asrTask.output.result) {
base.result = asrTask.output.result;
if (typeof asrTask.output.result.transcription_url === "string") {
base.transcription_url = asrTask.output.result.transcription_url;
}
}
if (asrTask.output.task_metrics) base.task_metrics = asrTask.output.task_metrics;
if (asrTask.usage) base.usage = asrTask.usage;
if (transcriptionItems.length === 0) {
return base;
}
const texts: string[] = [];
const transcripts: Record<string, unknown>[] = [];
for (const item of transcriptionItems) {
if (!item.transcription_url) continue;
const transRes = await fetch(item.transcription_url, { signal: ctx.signal });
if (!transRes.ok) {
throw new PipelineError(
"async_task_failed",
`Failed to download transcription: HTTP ${transRes.status}`,
{ step: "speech/recognize", details: { taskId, url: item.transcription_url } },
);
}
const transData = (await transRes.json()) as Record<string, unknown>;
transcripts.push(transData);
const transcriptList = transData.transcripts as
| Array<{ text?: string; sentences?: Array<{ text?: string }> }>
| undefined;
if (!transcriptList?.length) continue;
for (const transcript of transcriptList) {
if (transcript.sentences?.length) {
for (const sentence of transcript.sentences) {
if (sentence.text) texts.push(sentence.text);
}
} else if (transcript.text) {
texts.push(transcript.text);
}
}
}
return {
...base,
text: texts.join("\n"),
transcripts,
};
}
// --- Shared: task polling ---
@@ -640,7 +779,7 @@ function flattenTaskResponse(resp: DashScopeTaskResponse): Record<string, unknow
if (urls.length > 0) flat.urls = urls;
}
if (output.results) {
const urls = output.results.map((r) => r.url).filter(Boolean);
const urls = output.results.map((item) => item.url).filter(Boolean);
if (urls.length > 0 && !flat.urls) flat.urls = urls;
}
if (output.task_metrics) flat.task_metrics = output.task_metrics;
@@ -658,13 +797,13 @@ async function pollTask(
return await pollTaskWithOptions(env, taskId, pollIntervalMs, timeoutMs, ctx);
}
async function pollTaskWithOptions(
async function pollUntilSucceeded(
env: PipelineEnv,
taskId: string,
pollIntervalMs: number,
timeoutMs: number,
ctx?: StepContext,
): Promise<Record<string, unknown>> {
): Promise<DashScopeTaskResponse> {
const started = Date.now();
let attempt = 0;
@@ -687,7 +826,7 @@ async function pollTaskWithOptions(
const status = result.output.task_status;
if (status === "SUCCEEDED") {
return flattenTaskResponse(result);
return result;
}
if (status === "FAILED") {
@@ -718,6 +857,33 @@ async function pollTaskWithOptions(
}
}
async function pollTaskWithOptions(
env: PipelineEnv,
taskId: string,
pollIntervalMs: number,
timeoutMs: number,
ctx?: StepContext,
): Promise<Record<string, unknown>> {
return flattenTaskResponse(await pollUntilSucceeded(env, taskId, pollIntervalMs, timeoutMs, ctx));
}
/** ASR task polling: preserve original output (includes results[] / result.transcription_url). */
async function pollAsrTaskWithOptions(
env: PipelineEnv,
taskId: string,
pollIntervalMs: number,
timeoutMs: number,
ctx?: StepContext,
): Promise<DashScopeASRTaskResult> {
return (await pollUntilSucceeded(
env,
taskId,
pollIntervalMs,
timeoutMs,
ctx,
)) as DashScopeASRTaskResult;
}
function delay(ms: number, signal?: AbortSignal): Promise<void> {
if (!signal) return new Promise((resolve) => setTimeout(resolve, ms));
return new Promise((resolve, reject) => {
+51 -14
View File
@@ -25,6 +25,13 @@ interface CommandNode {
children: Map<string, CommandNode>;
}
const AUTH_LABELS = {
apiKey: "API Key",
console: "Console",
openapi: "AK/SK",
none: "No Auth",
} satisfies Record<AuthRequirement, string>;
/**
* What a command path resolves to in the registry. The single judgement that
* feeds `resolve()` — no scattered `isGroupPath` + throwing `resolve`.
@@ -157,14 +164,37 @@ export class CommandRegistry {
};
}
private buildResourceLines(a: (s: string) => string, d: (s: string) => string): string {
const entries: Array<{ path: string; desc: string }> = [];
private buildCommandLines(
entries: Array<{ path: string; auth: AuthRequirement; desc: string }>,
accent: (text: string) => string,
dim: (text: string) => string,
): string {
const maxPathLength = Math.max(...entries.map((entry) => entry.path.length));
const maxAuthLength = Math.max(
...entries.map((entry) => `[${AUTH_LABELS[entry.auth]}]`.length),
);
const rows = entries.map((entry) => {
const authLabel = `[${AUTH_LABELS[entry.auth]}]`;
return ` ${accent(entry.path.padEnd(maxPathLength + 2))} ${accent(authLabel.padEnd(maxAuthLength + 2))} ${dim(entry.desc)}`;
});
return rows.join("\n");
}
private buildResourceLines(
accent: (text: string) => string,
dim: (text: string) => string,
): string {
const entries: Array<{ path: string; auth: AuthRequirement; desc: string }> = [];
const collect = (node: CommandNode, prefix: string) => {
for (const [name, child] of node.children) {
const fullPath = prefix ? `${prefix} ${name}` : name;
if (child.command) {
entries.push({ path: fullPath, desc: child.command.description });
entries.push({
path: fullPath,
auth: child.command.auth,
desc: child.command.description,
});
}
if (child.children.size > 0) {
collect(child, fullPath);
@@ -173,8 +203,7 @@ export class CommandRegistry {
};
collect(this.root, "");
const maxLen = Math.max(...entries.map((e) => e.path.length));
return entries.map((e) => ` ${a(e.path.padEnd(maxLen + 2))} ${d(e.desc)}`).join("\n");
return this.buildCommandLines(entries, accent, dim);
}
private buildFlagLines(
@@ -341,6 +370,7 @@ ${authFlagSections ? `${authFlagSections}\n\n` : ""}${b("Getting Help:")}
out.write(`\n${cmd.description}\n`);
out.write(`${b("Usage:")} ${prefix}${cmd.usageArgs ? ` ${cmd.usageArgs}` : ""}\n`);
out.write(`${b("Authentication:")} ${a(AUTH_LABELS[cmd.auth])}\n`);
const flagEntries = [
...Object.entries(cmd.flags ?? {}),
...Object.entries(credentialFlagDefs(cmd)),
@@ -373,18 +403,25 @@ ${authFlagSections ? `${authFlagSections}\n\n` : ""}${b("Getting Help:")}
}
private printChildren(node: CommandNode, prefix: string, out: NodeJS.WriteStream): void {
const entries: Array<{ fullName: string; description: string }> = [];
const collect = (n: CommandNode, p: string) => {
for (const [name, child] of n.children) {
const entries: Array<{ path: string; auth: AuthRequirement; desc: string }> = [];
const collect = (currentNode: CommandNode, currentPath: string) => {
for (const [name, child] of currentNode.children) {
if (child.command)
entries.push({ fullName: `${p} ${name}`, description: child.command.description });
if (child.children.size > 0) collect(child, `${p} ${name}`);
entries.push({
path: `${currentPath} ${name}`,
auth: child.command.auth,
desc: child.command.description,
});
if (child.children.size > 0) collect(child, `${currentPath} ${name}`);
}
};
collect(node, prefix);
const maxLen = Math.max(...entries.map((e) => e.fullName.length));
for (const { fullName, description } of entries) {
out.write(` ${this.accent(fullName.padEnd(maxLen), out)} ${this.dim(description, out)}\n`);
}
out.write(
this.buildCommandLines(
entries,
(text) => this.accent(text, out),
(text) => this.dim(text, out),
) + "\n",
);
}
}
@@ -0,0 +1,184 @@
import { expect, test } from "vite-plus/test";
import type { Client } from "bailian-cli-core";
import { PipelineError } from "../src/pipeline/errors.ts";
import type { PipelineEnv } from "../src/pipeline/bl-config.ts";
import { speechRecognize } from "../src/pipeline/steps/bl-api.ts";
import type { StepContext } from "../src/pipeline/types.ts";
type CapturedRequest = {
path?: string;
method?: string;
headers?: Record<string, string>;
body?: Record<string, unknown>;
async?: boolean;
};
function makeEnv(requestJsonImpl?: (opts: CapturedRequest) => Promise<unknown>): {
env: PipelineEnv;
captured: CapturedRequest[];
} {
const captured: CapturedRequest[] = [];
const client = {
uploadFile: async (source: string) => source,
requestJson: async (opts: CapturedRequest) => {
captured.push(opts);
if (requestJsonImpl) return requestJsonImpl(opts);
return { output: { text: "ok" } };
},
} as unknown as Client;
return {
env: {
client,
settings: { quiet: true, output: "json" } as PipelineEnv["settings"],
},
captured,
};
}
function makeCtx(): StepContext {
return { dryRun: false, signal: new AbortController().signal };
}
test("pipeline speechRecognize routes input-audio flash to sync multimodal endpoint", async () => {
const { env, captured } = makeEnv();
const result = (await speechRecognize(
env,
{
url: "https://example.com/a.wav",
model: "qwen-audio-3.0-asr-flash",
language: "en",
"vocabulary-id": "vocab-1",
},
makeCtx(),
)) as { mode?: string; text?: string };
expect(result.mode).toBe("sync");
expect(result.text).toBe("ok");
expect(captured).toHaveLength(1);
expect(captured[0]?.path).toBe("/api/v1/services/aigc/multimodal-generation/generation");
expect(captured[0]?.headers?.["X-DashScope-SSE"]).toBe("disable");
expect(captured[0]?.body).toMatchObject({
model: "qwen-audio-3.0-asr-flash",
parameters: {
format: "wav",
language_hints: ["en"],
vocabulary_id: "vocab-1",
},
});
});
test("pipeline speechRecognize maps qwen3-filetrans language to parameters.language", async () => {
const { env, captured } = makeEnv(async (opts) => {
if (opts.async || opts.method === "POST") {
return { output: { task_id: "task-1", task_status: "PENDING" } };
}
return {
output: { task_id: "task-1", task_status: "SUCCEEDED", results: [] },
request_id: "r1",
};
});
await speechRecognize(
env,
{
url: "https://example.com/a.wav",
model: "qwen3-asr-flash-filetrans",
language: "zh",
"poll-interval": 0,
},
makeCtx(),
);
expect(captured[0]?.path).toBe("/api/v1/services/audio/asr/transcription");
expect(captured[0]?.async).toBe(true);
expect(captured[0]?.body).toMatchObject({
model: "qwen3-asr-flash-filetrans",
input: { file_url: "https://example.com/a.wav" },
parameters: { language: "zh" },
});
expect(
(captured[0]?.body?.parameters as Record<string, unknown> | undefined)?.language_hints,
).toBeUndefined();
});
test("pipeline speechRecognize rejects realtime models before requesting", async () => {
const { env, captured } = makeEnv();
await expect(
speechRecognize(
env,
{ url: "https://example.com/a.wav", model: "qwen3-asr-flash-realtime" },
makeCtx(),
),
).rejects.toBeInstanceOf(PipelineError);
expect(captured).toHaveLength(0);
});
test("pipeline speechRecognize rejects multiple urls for sync flash", async () => {
const { env, captured } = makeEnv();
await expect(
speechRecognize(
env,
{
url: ["https://example.com/a.wav", "https://example.com/b.wav"],
model: "fun-asr-flash-2026-06-15",
},
makeCtx(),
),
).rejects.toBeInstanceOf(PipelineError);
expect(captured).toHaveLength(0);
});
test("pipeline speechRecognize downloads qwen3 singular result.transcription_url", async () => {
const originalFetch = globalThis.fetch;
const transcriptionUrl = "https://example.com/transcription.json";
globalThis.fetch = (async (input: RequestInfo | URL) => {
const url = typeof input === "string" ? input : input instanceof URL ? input.href : input.url;
expect(url).toBe(transcriptionUrl);
return new Response(
JSON.stringify({
transcripts: [{ text: "pipeline hello", sentences: [{ text: "pipeline hello" }] }],
}),
{ status: 200, headers: { "Content-Type": "application/json" } },
);
}) as typeof fetch;
try {
const { env, captured } = makeEnv(async (opts) => {
if (opts.async || opts.method === "POST") {
return { output: { task_id: "task-1", task_status: "PENDING" } };
}
return {
output: {
task_id: "task-1",
task_status: "SUCCEEDED",
result: { transcription_url: transcriptionUrl },
},
request_id: "r1",
};
});
const result = (await speechRecognize(
env,
{
url: "https://example.com/a.wav",
model: "qwen3-asr-flash-filetrans",
"poll-interval": 0,
},
makeCtx(),
)) as {
mode?: string;
text?: string;
transcription_url?: string;
result?: { transcription_url?: string };
};
expect(captured[0]?.async).toBe(true);
expect(result.mode).toBe("async");
expect(result.text).toBe("pipeline hello");
expect(result.transcription_url).toBe(transcriptionUrl);
expect(result.result?.transcription_url).toBe(transcriptionUrl);
} finally {
globalThis.fetch = originalFetch;
}
});