FN-8571: add Parakeet STT transcription backend
Add a configurable sherpa-onnx Parakeet v3 speech-to-text service and dashboard API. - Add model download, validation, caching, and lifecycle management for the bundled STT runtime. - Register multipart-safe voice transcription routes with request-size handling and service-backed responses. - Expose STT settings, documentation, release metadata, and focused integration coverage. Files changed: .changeset/fn-8571-voice-stt-backend.md | 7 + docs/settings-reference.md | 12 + .../core/src/__tests__/settings-parity.test.ts | 10 + packages/core/src/index.ts | 2 +- packages/core/src/settings-schema.ts | 2 + packages/core/src/types.ts | 2 + packages/core/src/types/settings-scope.ts | 19 ++ packages/dashboard/package.json | 3 + .../voice-body-parser-integration.test.ts | 67 ++++++ packages/dashboard/src/routes.ts | 2 + packages/dashboard/src/routes/README.md | 74 +++--- .../routes/__tests__/register-voice-routes.test.ts | 159 ++++++++++++ .../src/routes/create-api-routes-mount-sequence.ts | 2 +- .../dashboard/src/routes/register-voice-routes.ts | 92 +++++++ packages/dashboard/src/server.ts | 21 +- .../src/stt/__tests__/model-manager.test.ts | 171 +++++++++++++ .../dashboard/src/stt/__tests__/voice-stt.test.ts | 48 ++++ packages/dashboard/src/stt/model-manager.ts | 268 +++++++++++++++++++++ packages/dashboard/src/stt/parakeet-service.ts | 77 ++++++ packages/dashboard/src/stt/types.ts | 14 ++ pnpm-lock.yaml | 65 +++++ 21 files changed, 1086 insertions(+), 31 deletions(-) Fusion-Task-Id: FN-8571 Fusion-Task-Lineage: b11e0cea-9c9c-41d4-a141-c47e5ed56dba Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8571-voice-stt-backend.md
Normal file
7
.changeset/fn-8571-voice-stt-backend.md
Normal file
@@ -0,0 +1,7 @@
|
||||
---
|
||||
"@runfusion/fusion": minor
|
||||
---
|
||||
|
||||
summary: Add opt-in voice transcription model lifecycle and API support.
|
||||
category: feature
|
||||
dev: Adds voiceInput settings, optional lazy sherpa runtime, checksum-gated shared model cache, and project-bound PCM voice endpoints.
|
||||
@@ -45,6 +45,18 @@ See [Signals Connectors](./signals-connectors.md) for setup, signing, payload, a
|
||||
|
||||
---
|
||||
|
||||
## Voice input
|
||||
|
||||
`VoiceInputSettings` is exported from `@fusion/core`. Both global and project settings accept
|
||||
`voiceInput: { enabled?, model?, language? }`; values resolve per request with project precedence.
|
||||
`enabled` defaults to false and gates dictation only, so model status/download/delete remain available
|
||||
while disabled. `model` defaults to registry identifier `"parakeet-v3"` and `language` to `"en"`;
|
||||
unsupported values are rejected and never become URLs or paths. The optional sherpa runtime and
|
||||
user-scoped cache degrade to unavailable safely. Downloads are on demand and require a pinned SHA-256;
|
||||
unpinned assets refuse download. Status polling reports `queued`/`downloading`; deleting fences an
|
||||
in-flight download. Voice chunks alone allow 2 MiB JSON, use 16 kHz mono PCM, and sessions retain
|
||||
closed tombstones for 60 seconds (size-cap responses stay 413 before eviction).
|
||||
|
||||
## Global Settings
|
||||
|
||||
Defaults from `DEFAULT_GLOBAL_SETTINGS`; key scope from `GLOBAL_SETTINGS_KEYS`.
|
||||
|
||||
@@ -1,3 +1,4 @@
|
||||
import type { VoiceInputSettings, ProjectSettings } from "../types.js";
|
||||
import { describe, expect, it } from "vitest";
|
||||
import {
|
||||
DEFAULT_GLOBAL_SETTINGS,
|
||||
@@ -542,6 +543,7 @@ describe("settings key parity", () => {
|
||||
// GLOBAL_SETTINGS_KEYS order.
|
||||
expect(overlap).toEqual([
|
||||
"testMode",
|
||||
"voiceInput",
|
||||
"mergeRequestContractShadowEnabled",
|
||||
"taskTokenBudget",
|
||||
"githubTrackingDefaultRepo",
|
||||
@@ -717,6 +719,14 @@ describe("model lane key parity regression (FN-1729)", () => {
|
||||
expect(isGlobalSettingsKey("testMode")).toBe(true);
|
||||
});
|
||||
|
||||
it("exports voiceInput from the public type barrel in both scopes", () => {
|
||||
const sample: VoiceInputSettings = { enabled: true, model: "parakeet-v3", language: "en" };
|
||||
const projectValue: ProjectSettings["voiceInput"] = sample;
|
||||
expect(projectValue).toEqual(sample);
|
||||
expect(isProjectSettingsKey("voiceInput")).toBe(true);
|
||||
expect(isGlobalSettingsKey("voiceInput")).toBe(true);
|
||||
});
|
||||
|
||||
it("no model lane provider exists without its corresponding modelId key", () => {
|
||||
const allKeys = new Set([...GLOBAL_SETTINGS_KEYS, ...PROJECT_SETTINGS_KEYS]);
|
||||
|
||||
|
||||
File diff suppressed because one or more lines are too long
@@ -143,6 +143,7 @@ export const DEFAULT_GLOBAL_SETTINGS = {
|
||||
defaultProvider: undefined,
|
||||
defaultModelId: undefined,
|
||||
testMode: undefined,
|
||||
voiceInput: undefined,
|
||||
modelPricingOverrides: undefined,
|
||||
modelPricingFetchedAt: undefined,
|
||||
modelPricingSource: undefined,
|
||||
@@ -436,6 +437,7 @@ export const DEFAULT_PROJECT_SETTINGS = {
|
||||
// with this on but auto-merge off, review threads are resolved but the PR is not merged.
|
||||
autoResolveReviewComments: true,
|
||||
testMode: undefined,
|
||||
voiceInput: undefined,
|
||||
mergeRequestContractShadowEnabled: false,
|
||||
mergeStrategy: "direct",
|
||||
directMergeCommitStrategy: "always-squash",
|
||||
|
||||
@@ -737,6 +737,7 @@ import type {
|
||||
RemoteAccessLifecycleConfig,
|
||||
RemoteAccessProjectSettings,
|
||||
SecretsEnvSettings,
|
||||
VoiceInputSettings,
|
||||
ProjectSettings,
|
||||
Settings,
|
||||
} from "./types/settings-scope.js";
|
||||
@@ -787,6 +788,7 @@ export type {
|
||||
RemoteAccessLifecycleConfig,
|
||||
RemoteAccessProjectSettings,
|
||||
SecretsEnvSettings,
|
||||
VoiceInputSettings,
|
||||
ProjectSettings,
|
||||
Settings,
|
||||
};
|
||||
|
||||
@@ -289,6 +289,23 @@ export interface BackupSettingsMigrationConflict {
|
||||
recordedAt: string;
|
||||
}
|
||||
|
||||
/**
|
||||
* FNXC:VoiceInput 2026-07-21-12:00:
|
||||
* Voice dictation is opt-in and resolved per request with project values overriding global
|
||||
* values. The flag gates dictation only: operators may manage the on-demand model while
|
||||
* disabled, and missing optional runtime support reports unavailable rather than failing boot.
|
||||
* `model` is a registry identifier (default `parakeet-v3`), never a path or URL; `language`
|
||||
* defaults to `en` and unsupported values are rejected when a session is created.
|
||||
*/
|
||||
export interface VoiceInputSettings {
|
||||
/** Default false; gates dictation/transcription only. */
|
||||
enabled?: boolean;
|
||||
/** Registry model identifier; default `parakeet-v3`. */
|
||||
model?: string;
|
||||
/** Recognition language code; default `en`. */
|
||||
language?: string;
|
||||
}
|
||||
|
||||
export interface GlobalSettings {
|
||||
/** Maximum PostgreSQL server connections for Fusion's embedded database. Applied on the next Fusion restart. */
|
||||
embeddedPostgresMaxConnections?: number;
|
||||
@@ -352,6 +369,7 @@ export interface GlobalSettings {
|
||||
* of per-task or per-lane overrides. No network calls, zero token cost.
|
||||
* Project `testMode` takes precedence over the global value. */
|
||||
testMode?: boolean;
|
||||
voiceInput?: VoiceInputSettings;
|
||||
/**
|
||||
* User-edited or one-click-fetched pricing entries keyed by lowercased `provider:model`.
|
||||
*
|
||||
@@ -1117,6 +1135,7 @@ export interface ProjectSettings {
|
||||
/** When true, force every AI lane onto the deterministic mock provider regardless
|
||||
* of per-task or per-lane overrides. No network calls, zero token cost. */
|
||||
testMode?: boolean;
|
||||
voiceInput?: VoiceInputSettings;
|
||||
/** Phase-1 FN-5741 write-only shadow seam toggle.
|
||||
* Overrides global `mergeRequestContractShadowEnabled` when defined.
|
||||
* Default: false. */
|
||||
|
||||
@@ -155,6 +155,9 @@
|
||||
"ws": "^8.18.0",
|
||||
"zod": "^3.25.76"
|
||||
},
|
||||
"optionalDependencies": {
|
||||
"sherpa-onnx-node": "^1.13.4"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@testing-library/jest-dom": "^6.9.1",
|
||||
"@testing-library/react": "^16.3.2",
|
||||
|
||||
@@ -0,0 +1,67 @@
|
||||
// @vitest-environment node
|
||||
|
||||
import { EventEmitter } from "node:events";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import type { Settings, TaskStore } from "@fusion/core";
|
||||
import { createServer } from "../server.js";
|
||||
import { request } from "../test-request.js";
|
||||
|
||||
/**
|
||||
* FNXC:VoiceInput 2026-07-21-18:20:
|
||||
* This intentionally boots server.ts rather than a registrar-only router. The global JSON
|
||||
* middleware is installed before registrars, so only this path proves the voice exclusion is
|
||||
* load-bearing while all other API routes retain the existing 100 KiB budget.
|
||||
*/
|
||||
class VoiceParserStore extends EventEmitter {
|
||||
constructor(private readonly voiceEnabled = false) { super(); }
|
||||
getRootDir() { return process.cwd(); }
|
||||
getFusionDir() { return `${process.cwd()}/.fusion`; }
|
||||
getSettings = vi.fn(async (): Promise<Settings> => ({ voiceInput: { enabled: this.voiceEnabled } } as Settings));
|
||||
getSettingsFast = this.getSettings;
|
||||
getGlobalSettingsStore = () => ({ getSettings: async () => ({}) });
|
||||
getAsyncLayer = vi.fn(() => ({ db: { update: vi.fn(() => ({ set: vi.fn(() => ({ where: vi.fn(() => ({ returning: vi.fn(async () => []) })) })) })) } }));
|
||||
getProjectScopedPluginMcpServers = vi.fn().mockResolvedValue([]);
|
||||
getTaskWorkflowSelection = vi.fn();
|
||||
getWorkflowDefinition = vi.fn(async () => undefined);
|
||||
getWorkflowSettingValues = vi.fn(() => ({}));
|
||||
getWorkflowSettingsProjectId = vi.fn(() => "default");
|
||||
}
|
||||
|
||||
const app = (voiceEnabled = false) => createServer(new VoiceParserStore(voiceEnabled) as unknown as TaskStore, { noAuth: true });
|
||||
|
||||
describe("voice body parser middleware boundary", () => {
|
||||
it("lets a voice JSON body above the global 100 KiB limit reach the chunk handler intact", async () => {
|
||||
// Enable dictation so the handler must read sessionId after its 2 MiB parser, rather than
|
||||
// rejecting on the guard before consuming the parsed request body.
|
||||
const server = app(true);
|
||||
const body = JSON.stringify({ sessionId: "never", audio: "A".repeat(110 * 1024), sequence: 0, final: false });
|
||||
const response = await request(server, "POST", "/api/voice/transcribe", body, { "content-type": "application/json" });
|
||||
// The route proves it received the complete parsed session id; without the server.ts
|
||||
// exclusion Express rejects this same request at the global 100 KiB parser with 413.
|
||||
expect(response.status).toBe(404);
|
||||
expect(response.body).toEqual({ error: "unknown-session" });
|
||||
});
|
||||
|
||||
it("treats the trailing-slash spelling as the same 2 MiB voice endpoint", async () => {
|
||||
const server = app(true);
|
||||
const body = JSON.stringify({ sessionId: "never", audio: "A".repeat(110 * 1024), sequence: 0, final: false });
|
||||
const response = await request(server, "POST", "/api/voice/transcribe/", body, { "content-type": "application/json" });
|
||||
expect(response.status).toBe(404);
|
||||
expect(response.body).toEqual({ error: "unknown-session" });
|
||||
});
|
||||
|
||||
it("keeps the global 100 KiB JSON budget outside the one voice exclusion", async () => {
|
||||
const overGlobalLimit = await request(app(), "POST", "/api/no-such-json-route", JSON.stringify({ payload: "x".repeat(110 * 1024) }), { "content-type": "application/json" });
|
||||
expect(overGlobalLimit.status).toBe(413);
|
||||
expect(overGlobalLimit.body).toEqual({ error: "payload-too-large" });
|
||||
});
|
||||
|
||||
it("maps voice parser size and JSON syntax errors to the public JSON contract", async () => {
|
||||
const tooLarge = await request(app(), "POST", "/api/voice/transcribe", JSON.stringify({ audio: "A".repeat(2 * 1024 * 1024 + 1) }), { "content-type": "application/json" });
|
||||
expect(tooLarge.status).toBe(413);
|
||||
expect(tooLarge.body).toEqual({ error: "payload-too-large", limitBytes: 2 * 1024 * 1024 });
|
||||
const malformed = await request(app(), "POST", "/api/voice/transcribe", "{", { "content-type": "application/json" });
|
||||
expect(malformed.status).toBe(400);
|
||||
expect(malformed.body).toEqual({ error: "invalid-request" });
|
||||
});
|
||||
});
|
||||
@@ -84,6 +84,7 @@ import { registerAuthRoutes } from "./routes/register-auth-routes.js";
|
||||
import { registerRuntimeProviderRoutes } from "./routes/register-runtime-provider-routes.js";
|
||||
import { registerFnBinaryRoutes } from "./routes/register-fn-binary-routes.js";
|
||||
import { registerUpdateCheckRoutes } from "./routes/register-update-check-routes.js";
|
||||
import { registerVoiceRoutes } from "./routes/register-voice-routes.js";
|
||||
import { registerDiagnosticsRoutes } from "./routes/register-diagnostics-routes.js";
|
||||
import { registerSystemRoutes } from "./routes/register-system-routes.js";
|
||||
import { registerCliAgentHooksRoute } from "./routes/cli-agent-hooks.js";
|
||||
@@ -1252,6 +1253,7 @@ export function createApiRoutes(store: TaskStore, options?: ServerOptions): Rout
|
||||
// by opening storm-guarded fix tasks back in triage.
|
||||
registrarMounter.mount("registerMonitorRoutes", () => registerMonitorRoutes(routeContext));
|
||||
registrarMounter.mount("registerUpdateCheckRoutes", () => registerUpdateCheckRoutes(routeContext));
|
||||
registrarMounter.mount("registerVoiceRoutes", () => registerVoiceRoutes(routeContext));
|
||||
registrarMounter.mount("registerDiagnosticsRoutes", () => registerDiagnosticsRoutes(routeContext));
|
||||
// CLI Agent Executor hook ingestion (U17) — per-session token auth, exempt from
|
||||
// the daemon bearer-token middleware (hook scripts only hold the session token).
|
||||
|
||||
@@ -40,6 +40,7 @@ The following is the complete top-level registrar map currently imported by `rou
|
||||
- `registerSignalRoutes` — domain registrar mounted by `createApiRoutes`.
|
||||
- `registerMonitorRoutes` — domain registrar mounted by `createApiRoutes`.
|
||||
- `registerUpdateCheckRoutes` — domain registrar mounted by `createApiRoutes`.
|
||||
- `registerVoiceRoutes` — opt-in voice model lifecycle and project-bound PCM transcription endpoints.
|
||||
- `registerDiagnosticsRoutes` — domain registrar mounted by `createApiRoutes`.
|
||||
- `registerCliAgentHooksRoute` — domain registrar mounted by `createApiRoutes`.
|
||||
- `registerCliAgentSettingsRoutes` — domain registrar mounted by `createApiRoutes`.
|
||||
@@ -105,33 +106,34 @@ Express matches in registration order. `create-api-routes-mount-sequence.ts` is
|
||||
28. `registerSignalRoutes`
|
||||
29. `registerMonitorRoutes`
|
||||
30. `registerUpdateCheckRoutes`
|
||||
31. `registerDiagnosticsRoutes`
|
||||
32. `registerCliAgentHooksRoute`
|
||||
33. `registerCliAgentSettingsRoutes`
|
||||
34. `registerActivityLogRoutes`
|
||||
35. `registerAgentCoreListCreateRoutes`
|
||||
36. `registerAgentImportExportRoutes`
|
||||
37. `registerOrgPortabilityRoutes`
|
||||
38. `registerAgentCoreRoutes`
|
||||
39. `registerAgentRuntimeRoutes`
|
||||
40. `registerSystemRoutes`
|
||||
41. `registerAgentReflectionRatingRoutes`
|
||||
42. `registerAgentGenerationRoutes`
|
||||
43. `registerIntegratedRouters`
|
||||
44. `registerProjectRoutes`
|
||||
45. `registerNodeRoutes`
|
||||
46. `registerDockerNodeRoutes`
|
||||
47. `registerDockerProvisioningRoutes`
|
||||
48. `registerSettingsSyncRoutes`
|
||||
49. `registerSecretsSyncRoutes`
|
||||
50. `registerMeshRoutes`
|
||||
51. `registerDiscoveryRoutes`
|
||||
52. `registerSettingsSyncInboundRoutes`
|
||||
53. `registerSecretsSyncInboundRoutes`
|
||||
54. `registerSetupActivityRoutes`
|
||||
55. `registerIntegratedDevServerRouter`
|
||||
56. `registerAgentSkillsRoutes`
|
||||
57. `registerProxyRoutes`
|
||||
31. `registerVoiceRoutes`
|
||||
32. `registerDiagnosticsRoutes`
|
||||
33. `registerCliAgentHooksRoute`
|
||||
34. `registerCliAgentSettingsRoutes`
|
||||
35. `registerActivityLogRoutes`
|
||||
36. `registerAgentCoreListCreateRoutes`
|
||||
37. `registerAgentImportExportRoutes`
|
||||
38. `registerOrgPortabilityRoutes`
|
||||
39. `registerAgentCoreRoutes`
|
||||
40. `registerAgentRuntimeRoutes`
|
||||
41. `registerSystemRoutes`
|
||||
42. `registerAgentReflectionRatingRoutes`
|
||||
43. `registerAgentGenerationRoutes`
|
||||
44. `registerIntegratedRouters`
|
||||
45. `registerProjectRoutes`
|
||||
46. `registerNodeRoutes`
|
||||
47. `registerDockerNodeRoutes`
|
||||
48. `registerDockerProvisioningRoutes`
|
||||
49. `registerSettingsSyncRoutes`
|
||||
50. `registerSecretsSyncRoutes`
|
||||
51. `registerMeshRoutes`
|
||||
52. `registerDiscoveryRoutes`
|
||||
53. `registerSettingsSyncInboundRoutes`
|
||||
54. `registerSecretsSyncInboundRoutes`
|
||||
55. `registerSetupActivityRoutes`
|
||||
56. `registerIntegratedDevServerRouter`
|
||||
57. `registerAgentSkillsRoutes`
|
||||
58. `registerProxyRoutes`
|
||||
<!-- mount-sequence:end -->
|
||||
|
||||
## Ordering rules
|
||||
@@ -153,3 +155,21 @@ Residual inline handlers in `routes.ts` are grandfathered only. `pnpm check:rout
|
||||
pnpm --filter @fusion/dashboard typecheck
|
||||
pnpm --filter @fusion/dashboard exec vitest run src/routes/__tests__/create-api-routes-mount-order.test.ts --silent=passed-only --reporter=dot
|
||||
```
|
||||
|
||||
## Voice transcription
|
||||
|
||||
`registerVoiceRoutes` exposes `GET /voice/status`, `POST`/`DELETE /voice/model`, and dictation
|
||||
`POST /voice/session`, `POST /voice/transcribe`, and `DELETE /voice/session/:id`. Settings are
|
||||
resolved per request through `getScopedStore(req)` with project-over-global precedence. Voice is
|
||||
opt-in: only dictation endpoints require `voiceInput.enabled`; model inspection, download, and
|
||||
delete remain available while off because the user-scoped model cache is shared by projects.
|
||||
|
||||
Audio chunks are base64 raw 16 kHz mono signed-16-bit little-endian PCM. Chunks are ordered,
|
||||
limited to 1 MiB (2 MiB JSON body), and sessions are project-bound. Active sessions become
|
||||
60-second closed tombstones on completion, delete, expiry, model removal, or the 16 MiB cap; cap
|
||||
tombstones return 413 while other closed sessions return 409, then all evict to 404. Repeated
|
||||
DELETE during the tombstone returns `{ closed:true, alreadyClosed:true }`; unknown and foreign IDs
|
||||
return 404. Download returns 202 with queued/downloading state; poll status for progress.
|
||||
|
||||
`server.ts` excludes only `/api/voice/transcribe` from its global 100 KiB JSON parser so the
|
||||
route's 2 MiB parser can return JSON 413/400 errors; other routes retain raw-body HMAC capture.
|
||||
|
||||
@@ -0,0 +1,159 @@
|
||||
import express from "express";
|
||||
import type { AddressInfo } from "node:net";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { createRegisterVoiceRoutes } from "../register-voice-routes.js";
|
||||
import type { ApiRoutesContext } from "../types.js";
|
||||
|
||||
const servers: Array<ReturnType<ReturnType<typeof express>["listen"]>> = [];
|
||||
async function harness(enabled: boolean, ready = false, projectId = "project-a") {
|
||||
const app = express(); const router = express.Router(); app.use(router);
|
||||
const scoped = { getSettings: async () => ({ voiceInput: { enabled } }), getGlobalSettingsStore: () => ({ getSettings: async () => ({}) }) };
|
||||
const manager = { getState: async () => ({ status: ready ? "installed" as const : "not-installed" as const, installedPath: ready ? "/model" : undefined }), peekState: () => ({ status: "not-installed" as const }), scheduleDownload: () => ({ accepted: false, state: { status: "error" as const, errorReason: "checksum-unpinned" as const } }), remove: async () => {}, download: async () => ({ status: "not-installed" as const }), subscribe: () => () => {} };
|
||||
const service = ready ? { getRuntimeStatus: async () => ({ status: "available" as const }), createSession: async () => ({ acceptChunk: () => ({ partial: "ok" }), finish: () => ({ text: "ok" }), close: () => {} }) } : undefined;
|
||||
createRegisterVoiceRoutes({ manager, ...(service ? { service } : {}) })({ router, getScopedStore: async () => scoped, getProjectIdFromRequest: () => projectId } as unknown as ApiRoutesContext);
|
||||
const server = app.listen(0); servers.push(server);
|
||||
await new Promise<void>((resolve) => server.once("listening", resolve));
|
||||
const port = (server.address() as AddressInfo).port;
|
||||
return (path: string, init?: RequestInit) => fetch(`http://127.0.0.1:${port}${path}`, init);
|
||||
}
|
||||
afterEach(async () => { await Promise.all(servers.splice(0).map((server) => new Promise<void>((resolve) => server.close(() => resolve())))); });
|
||||
|
||||
describe("voice route authorization split", () => {
|
||||
it("allows lifecycle inspection while dictation is disabled", async () => {
|
||||
const request = await harness(false);
|
||||
expect((await request("/voice/status")).status).toBe(200);
|
||||
expect((await request("/voice/session", { method: "POST" })).status).toBe(409);
|
||||
expect(await (await request("/voice/model/download", { method: "POST" })).json()).toMatchObject({ error: "checksum-unpinned" });
|
||||
expect((await request("/voice/model", { method: "DELETE" })).status).toBe(200);
|
||||
});
|
||||
|
||||
it("fences a pending session creation when model deletion starts", async () => {
|
||||
const app = express();
|
||||
const router = express.Router();
|
||||
app.use(router);
|
||||
const close = vi.fn();
|
||||
let resolveCreation!: (session: { acceptChunk: () => { partial: string }; finish: () => { text: string }; close: () => void }) => void;
|
||||
let resolveRemoval!: () => void;
|
||||
let signalCreation!: () => void;
|
||||
let signalRemoval!: () => void;
|
||||
const creationStarted = new Promise<void>((resolve) => { signalCreation = resolve; });
|
||||
const removalStarted = new Promise<void>((resolve) => { signalRemoval = resolve; });
|
||||
const manager = {
|
||||
getState: async () => ({ status: "installed" as const, installedPath: "/model" }),
|
||||
peekState: () => ({ status: "installed" as const, installedPath: "/model" }),
|
||||
scheduleDownload: () => ({ accepted: true as const, state: { status: "installed" as const } }),
|
||||
remove: async () => { signalRemoval(); await new Promise<void>((resolve) => { resolveRemoval = resolve; }); },
|
||||
download: async () => ({ status: "installed" as const }),
|
||||
subscribe: () => () => {},
|
||||
};
|
||||
const service = {
|
||||
getRuntimeStatus: async () => ({ status: "available" as const }),
|
||||
createSession: async () => {
|
||||
signalCreation();
|
||||
return new Promise<{ acceptChunk: () => { partial: string }; finish: () => { text: string }; close: () => void }>((resolve) => { resolveCreation = resolve; });
|
||||
},
|
||||
};
|
||||
createRegisterVoiceRoutes({ manager, service })({
|
||||
router,
|
||||
getScopedStore: async () => ({ getSettings: async () => ({ voiceInput: { enabled: true } }), getGlobalSettingsStore: () => ({ getSettings: async () => ({}) }) }),
|
||||
getProjectIdFromRequest: () => "epoch-project",
|
||||
} as unknown as ApiRoutesContext);
|
||||
const server = app.listen(0); servers.push(server);
|
||||
await new Promise<void>((resolve) => server.once("listening", resolve));
|
||||
const port = (server.address() as AddressInfo).port;
|
||||
const request = (path: string, init?: RequestInit) => fetch(`http://127.0.0.1:${port}${path}`, init);
|
||||
|
||||
const creating = request("/voice/session", { method: "POST" });
|
||||
await creationStarted;
|
||||
const deleting = request("/voice/model", { method: "DELETE" });
|
||||
await removalStarted;
|
||||
resolveCreation({ acceptChunk: () => ({ partial: "" }), finish: () => ({ text: "" }), close });
|
||||
const rejected = await creating;
|
||||
expect(rejected.status).toBe(409);
|
||||
expect(await rejected.json()).toEqual({ error: "unavailable", reason: "model-removed" });
|
||||
expect(close).toHaveBeenCalledOnce();
|
||||
resolveRemoval();
|
||||
expect((await deleting).status).toBe(200);
|
||||
});
|
||||
|
||||
it("rejects non-canonical base64 before it reaches the recognizer", async () => {
|
||||
const request = await harness(true, true);
|
||||
const created = await request("/voice/session", { method: "POST" });
|
||||
const { sessionId } = await created.json() as { sessionId: string };
|
||||
const response = await request("/voice/transcribe", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ sessionId, audio: "AAAAAAAAA=", sequence: 0, final: false }) });
|
||||
expect(response.status).toBe(400);
|
||||
expect(await response.json()).toEqual({ error: "invalid-audio-payload" });
|
||||
});
|
||||
|
||||
it("maps malformed and oversized voice JSON to API errors", async () => {
|
||||
const request = await harness(true);
|
||||
const malformed = await request("/voice/transcribe", { method: "POST", headers: { "content-type": "application/json" }, body: "{" });
|
||||
expect([400, 409]).toContain(malformed.status);
|
||||
const huge = await request("/voice/transcribe", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ audio: "x".repeat(2 * 1024 * 1024 + 1) }) });
|
||||
expect(huge.status).toBe(413);
|
||||
expect(await huge.json()).toMatchObject({ error: "payload-too-large", limitBytes: 2 * 1024 * 1024 });
|
||||
});
|
||||
|
||||
it("keeps closed sessions as project-bound DELETE tombstones", async () => {
|
||||
const owner = await harness(true, true, "owner");
|
||||
const foreign = await harness(true, true, "foreign");
|
||||
const { sessionId } = await (await owner("/voice/session", { method: "POST" })).json() as { sessionId: string };
|
||||
expect(await (await owner(`/voice/session/${sessionId}`, { method: "DELETE" })).json()).toEqual({ sessionId, closed: true, alreadyClosed: false });
|
||||
expect(await (await owner(`/voice/session/${sessionId}`, { method: "DELETE" })).json()).toEqual({ sessionId, closed: true, alreadyClosed: true });
|
||||
const foreignResponse = await foreign(`/voice/session/${sessionId}`, { method: "DELETE" });
|
||||
const unknownResponse = await foreign("/voice/session/never-existed", { method: "DELETE" });
|
||||
expect(foreignResponse.status).toBe(404);
|
||||
expect(await foreignResponse.text()).toBe(await unknownResponse.text());
|
||||
});
|
||||
|
||||
it("enforces the active-session cap per project and frees capacity at tombstone close", async () => {
|
||||
const request = await harness(true, true, "capacity-project");
|
||||
const created = await Promise.all(Array.from({ length: 8 }, () => request("/voice/session", { method: "POST" })));
|
||||
expect(created.every((response) => response.status === 201)).toBe(true);
|
||||
expect((await request("/voice/session", { method: "POST" })).status).toBe(429);
|
||||
const { sessionId } = await created[0].json() as { sessionId: string };
|
||||
expect((await request(`/voice/session/${sessionId}`, { method: "DELETE" })).status).toBe(200);
|
||||
expect((await request("/voice/session", { method: "POST" })).status).toBe(201);
|
||||
// A different scoped project has an independent eight-session budget.
|
||||
const other = await harness(true, true, "capacity-other");
|
||||
expect((await other("/voice/session", { method: "POST" })).status).toBe(201);
|
||||
});
|
||||
|
||||
it("retains TTL and size-cap closures as tombstones before evicting them", async () => {
|
||||
const request = await harness(true, true, "expiry-project");
|
||||
const created = await request("/voice/session", { method: "POST" });
|
||||
const { sessionId } = await created.json() as { sessionId: string };
|
||||
const base = Date.now();
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(base + 60_001);
|
||||
const expired = await request("/voice/transcribe", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ sessionId, audio: "AAA=", sequence: 0, final: false }) });
|
||||
expect(await expired.json()).toEqual({ error: "session-closed", reason: "ttl-expired" });
|
||||
vi.setSystemTime(base + 120_002);
|
||||
const evicted = await request("/voice/transcribe", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ sessionId, audio: "AAA=", sequence: 0, final: false }) });
|
||||
expect(await evicted.json()).toEqual({ error: "unknown-session" });
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("sweeps expired sessions before model deletion preserves their TTL tombstone", async () => {
|
||||
const request = await harness(true, true, "delete-sweep-project");
|
||||
const { sessionId } = await (await request("/voice/session", { method: "POST" })).json() as { sessionId: string };
|
||||
const base = Date.now();
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(base + 60_001);
|
||||
expect((await request("/voice/model", { method: "DELETE" })).status).toBe(200);
|
||||
const expired = await request("/voice/transcribe", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ sessionId, audio: "AAA=", sequence: 0, final: false }) });
|
||||
expect(await expired.json()).toEqual({ error: "session-closed", reason: "ttl-expired" });
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it("keeps a size-cap tombstone at 413 until eviction", async () => {
|
||||
const request = await harness(true, true, "size-cap-project");
|
||||
const { sessionId } = await (await request("/voice/session", { method: "POST" })).json() as { sessionId: string };
|
||||
// The cap is checked before recognizer work; setting an already-capped session through
|
||||
// sixteen 1 MiB chunks would add slow, low-signal HTTP work to this focused route test.
|
||||
const oversized = await request("/voice/transcribe", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ sessionId, audio: Buffer.alloc(1024 * 1024 + 2).toString("base64"), sequence: 0, final: false }) });
|
||||
expect(await oversized.json()).toEqual({ error: "payload-too-large", limitBytes: 1024 * 1024 });
|
||||
const repeated = await request("/voice/transcribe", { method: "POST", headers: { "content-type": "application/json" }, body: JSON.stringify({ sessionId, audio: "AAA=", sequence: 0, final: false }) });
|
||||
expect(await repeated.json()).toEqual({ error: "payload-too-large", limitBytes: 16 * 1024 * 1024 });
|
||||
});
|
||||
});
|
||||
@@ -12,7 +12,7 @@ export const CREATE_API_ROUTES_REGISTRAR_MOUNT_SEQUENCE = [
|
||||
"registerPluginsAutomationRoutes", "registerApprovalRoutes", "registerWorktrunkRoutes", "registerConfigMcpPiSettingsRoutes", "registerSystemMaintenanceRoutes", "registerModelRoutes",
|
||||
"registerCustomProviderRoutes", "registerAuthRoutes", "registerRuntimeProviderRoutes", "registerFnBinaryRoutes",
|
||||
"registerAiTextAssistantRoutes", "registerUsageRoutes", "registerCommandCenterRoutes", "registerKnowledgeRoutes", "registerReportRoutes",
|
||||
"registerSignalRoutes", "registerMonitorRoutes", "registerUpdateCheckRoutes", "registerDiagnosticsRoutes",
|
||||
"registerSignalRoutes", "registerMonitorRoutes", "registerUpdateCheckRoutes", "registerVoiceRoutes", "registerDiagnosticsRoutes",
|
||||
"registerCliAgentHooksRoute", "registerCliAgentSettingsRoutes", "registerActivityLogRoutes", "registerAgentCoreListCreateRoutes", "registerAgentImportExportRoutes",
|
||||
"registerOrgPortabilityRoutes", "registerAgentCoreRoutes", "registerAgentRuntimeRoutes", "registerSystemRoutes",
|
||||
"registerAgentReflectionRatingRoutes", "registerAgentGenerationRoutes", "registerIntegratedRouters", "registerProjectRoutes",
|
||||
|
||||
92
packages/dashboard/src/routes/register-voice-routes.ts
Normal file
92
packages/dashboard/src/routes/register-voice-routes.ts
Normal file
@@ -0,0 +1,92 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import express from "express";
|
||||
import type { VoiceInputSettings } from "@fusion/core";
|
||||
import { createVoiceModelManager } from "../stt/model-manager.js";
|
||||
import { createParakeetService, VoiceInputError } from "../stt/parakeet-service.js";
|
||||
import { DEFAULT_VOICE_LANGUAGE, DEFAULT_VOICE_MODEL_ID, resolveVoiceLanguage, resolveVoiceModelId } from "../stt/types.js";
|
||||
import type { ApiRouteRegistrar } from "./types.js";
|
||||
|
||||
const defaultManager = createVoiceModelManager();
|
||||
const defaultService = createParakeetService({ manager: defaultManager });
|
||||
type Session = { projectId?: string; recognizer: Awaited<ReturnType<typeof defaultService.createSession>>; next: number; bytes: number; last: number; expires: number; closed?: string; tombstone?: number };
|
||||
const sessions = new Map<string, Session>();
|
||||
// Session creation awaits native initialization; reservations close the await-window race.
|
||||
const pendingSessionReservations = new Map<string | undefined, number>();
|
||||
|
||||
/**
|
||||
* FNXC:VoiceInput 2026-07-21-21:10:
|
||||
* Model deletion fences pending native session creation immediately. A creation that began
|
||||
* before shared-cache removal must close its just-created recognizer and remain unavailable;
|
||||
* it must never insert an active session backed by a removed model.
|
||||
*/
|
||||
let modelEpoch = 0;
|
||||
const LIMIT = 1024 * 1024; const TOTAL_LIMIT = 16 * LIMIT;
|
||||
|
||||
/** FNXC:VoiceInput 2026-07-21-12:00: voice mode is opt-in, but model lifecycle remains usable while off so operators can install from settings first. */
|
||||
async function settingsFor(ctx: Parameters<ApiRouteRegistrar>[0], req: express.Request) {
|
||||
const scoped = await ctx.getScopedStore(req); const project = await scoped.getSettings() as { voiceInput?: VoiceInputSettings };
|
||||
const global = await scoped.getGlobalSettingsStore().getSettings() as { voiceInput?: VoiceInputSettings };
|
||||
const voice = { ...global.voiceInput, ...project.voiceInput }; return { voice, projectId: ctx.getProjectIdFromRequest(req) };
|
||||
}
|
||||
function sweep() { const now = Date.now(); for (const [id, session] of sessions) { if (!session.closed && (now - session.last > 60_000 || now > session.expires)) { session.recognizer.close(); session.closed = "ttl-expired"; session.tombstone = now + 60_000; } if (session.closed && (session.tombstone ?? 0) <= now) sessions.delete(id); } }
|
||||
function error(res: express.Response, code: number, body: object) { res.status(code).json(body); }
|
||||
|
||||
/**
|
||||
* FNXC:VoiceInput 2026-07-21-19:10:
|
||||
* Audio chunks use canonical standard base64. Node's permissive decoder accepts malformed
|
||||
* padding, so validate structure and round-trip equality before audio reaches a recognizer.
|
||||
*/
|
||||
function decodeCanonicalBase64(value: string): Buffer | undefined {
|
||||
if (!value || value.length % 4 !== 0 || !/^(?:[A-Za-z0-9+/]{4})*(?:[A-Za-z0-9+/]{2}==|[A-Za-z0-9+/]{3}=)?$/.test(value)) return undefined;
|
||||
const decoded = Buffer.from(value, "base64");
|
||||
return decoded.toString("base64") === value ? decoded : undefined;
|
||||
}
|
||||
|
||||
export function createRegisterVoiceRoutes(deps: { manager?: typeof defaultManager; service?: typeof defaultService } = {}): ApiRouteRegistrar {
|
||||
const manager = deps.manager ?? defaultManager;
|
||||
const service = deps.service ?? defaultService;
|
||||
return (ctx) => {
|
||||
const { router } = ctx;
|
||||
router.get("/voice/status", async (req, res) => { const { voice } = await settingsFor(ctx, req); const model = resolveVoiceModelId(voice.model); const language = resolveVoiceLanguage(voice.language); const modelState = manager.peekState().status === "queued" || manager.peekState().status === "downloading" ? manager.peekState() : await manager.getState(); res.json({ enabled: voice.enabled === true, modelId: "id" in model ? model.id : undefined, language: "language" in language ? language.language : undefined, unsupportedModel: "unsupported" in model ? model.unsupported : undefined, unsupportedLanguage: "unsupported" in language ? language.unsupported : undefined, model: modelState, runtime: await service.getRuntimeStatus() }); });
|
||||
router.post("/voice/model/download", async (req, res) => { const { voice } = await settingsFor(ctx, req); const model = resolveVoiceModelId(voice.model); if ("unsupported" in model) return error(res, 400, { error: "unsupported-model", value: model.unsupported, supported: [DEFAULT_VOICE_MODEL_ID] }); const scheduled = manager.scheduleDownload(); if (!scheduled.accepted) return error(res, 409, { error: scheduled.state.errorReason }); res.status(202).json({ state: scheduled.state }); });
|
||||
router.delete("/voice/model", async (_req, res) => {
|
||||
// Increment before awaiting cleanup so pending createSession() calls are fenced immediately.
|
||||
modelEpoch++;
|
||||
sweep();
|
||||
await manager.remove();
|
||||
const now = Date.now();
|
||||
for (const session of sessions.values()) if (!session.closed) {
|
||||
session.recognizer.close();
|
||||
session.closed = "model-removed";
|
||||
session.tombstone = now + 60_000;
|
||||
}
|
||||
res.json({ state: manager.peekState() });
|
||||
});
|
||||
router.post("/voice/session", async (req, res) => { sweep(); const { voice, projectId } = await settingsFor(ctx, req); if (voice.enabled !== true) return error(res, 409, { error: "disabled" }); const model = resolveVoiceModelId(voice.model); const language = resolveVoiceLanguage(voice.language); if ("unsupported" in model) return error(res, 400, { error: "unsupported-model", value: model.unsupported, supported: [DEFAULT_VOICE_MODEL_ID] }); if ("unsupported" in language) return error(res, 400, { error: "unsupported-language", value: language.unsupported, supported: [DEFAULT_VOICE_LANGUAGE] }); const active = [...sessions.values()].filter((s) => !s.closed && s.projectId === projectId).length; const reserved = pendingSessionReservations.get(projectId) ?? 0; if (active + reserved >= 8) return error(res, 429, { error: "too-many-sessions" });
|
||||
// FNXC:VoiceInput 2026-07-21-20:30: Reserve before native session creation awaits so
|
||||
// concurrent requests cannot all pass the eight-active-session capacity check.
|
||||
pendingSessionReservations.set(projectId, reserved + 1);
|
||||
const creationEpoch = modelEpoch;
|
||||
try {
|
||||
const recognizer = await service.createSession({ modelId: model.id, language: language.language });
|
||||
if (creationEpoch !== modelEpoch) {
|
||||
recognizer.close();
|
||||
return error(res, 409, { error: "unavailable", reason: "model-removed" });
|
||||
}
|
||||
const id = randomUUID();
|
||||
const now = Date.now();
|
||||
sessions.set(id, { projectId, recognizer, next: 0, bytes: 0, last: now, expires: now + 300_000 });
|
||||
res.status(201).json({ sessionId: id, expiresAt: new Date(now + 300_000).toISOString(), modelId: model.id, language: language.language });
|
||||
} catch (e) {
|
||||
return error(res, 409, { error: "unavailable", reason: e instanceof Error ? e.message : "unavailable" });
|
||||
} finally {
|
||||
const remaining = (pendingSessionReservations.get(projectId) ?? 1) - 1;
|
||||
if (remaining > 0) pendingSessionReservations.set(projectId, remaining); else pendingSessionReservations.delete(projectId);
|
||||
}
|
||||
});
|
||||
const parserError: express.ErrorRequestHandler = (err, _req, res, next) => { if (err?.type === "entity.too.large") return error(res, 413, { error: "payload-too-large", limitBytes: 2 * LIMIT }); if (err instanceof SyntaxError) return error(res, 400, { error: "invalid-request" }); next(err); };
|
||||
router.post("/voice/transcribe", express.json({ limit: "2mb" }), parserError, async (req: express.Request, res: express.Response) => { sweep(); const { voice, projectId } = await settingsFor(ctx, req); if (voice.enabled !== true) return error(res, 409, { error: "disabled" }); const body = req.body as { sessionId?: unknown; audio?: unknown; sequence?: unknown; final?: unknown; sampleRate?: unknown; channels?: unknown; encoding?: unknown; }; if (!body || typeof body.sessionId !== "string" || typeof body.audio !== "string" || typeof body.sequence !== "number" || typeof body.final !== "boolean") return error(res, 400, { error: "invalid-request" }); const session = sessions.get(body.sessionId); if (!session || session.projectId !== projectId) return error(res, 404, { error: "unknown-session" }); if (session.closed) return session.closed === "size-cap-exceeded" ? error(res, 413, { error: "payload-too-large", limitBytes: TOTAL_LIMIT }) : error(res, 409, { error: "session-closed", reason: session.closed }); if (body.sequence !== session.next) return error(res, 400, { error: "out-of-order-chunk", expected: session.next }); if ((body.sampleRate !== undefined && body.sampleRate !== 16000) || (body.channels !== undefined && body.channels !== 1) || (body.encoding !== undefined && body.encoding !== "pcm_s16le")) return error(res, 400, { error: "unsupported-audio-format", expected: { encoding: "pcm_s16le", sampleRate: 16000, channels: 1 } }); const audio = decodeCanonicalBase64(body.audio); if (!audio || !audio.length || audio.length % 2) return error(res, 400, { error: "invalid-audio-payload" }); if (audio.length > LIMIT || session.bytes + audio.length > TOTAL_LIMIT) { session.recognizer.close(); session.closed = "size-cap-exceeded"; session.tombstone = Date.now() + 60_000; return error(res, 413, { error: "payload-too-large", limitBytes: audio.length > LIMIT ? LIMIT : TOTAL_LIMIT }); } try { const result = session.recognizer.acceptChunk(audio, { final: body.final }); session.next++; session.bytes += audio.length; session.last = Date.now(); if (body.final) { session.recognizer.close(); session.closed = "completed"; session.tombstone = Date.now() + 60_000; return res.json({ sessionId: body.sessionId, sequence: body.sequence, text: result.text ?? "", final: true }); } return res.json({ sessionId: body.sessionId, sequence: body.sequence, partial: result.partial ?? "", final: false }); } catch (e) { return error(res, e instanceof VoiceInputError ? 400 : 409, { error: e instanceof VoiceInputError ? e.code : "unavailable" }); } });
|
||||
router.delete("/voice/session/:id", async (req, res) => { sweep(); const { voice, projectId } = await settingsFor(ctx, req); if (voice.enabled !== true) return error(res, 409, { error: "disabled" }); const session = sessions.get(req.params.id); if (!session || session.projectId !== projectId) return error(res, 404, { error: "unknown-session" }); if (session.closed) return res.json({ sessionId: req.params.id, closed: true, alreadyClosed: true }); session.recognizer.close(); session.closed = "deleted"; session.tombstone = Date.now() + 60_000; res.json({ sessionId: req.params.id, closed: true, alreadyClosed: false }); }); };
|
||||
}
|
||||
|
||||
export const registerVoiceRoutes = createRegisterVoiceRoutes();
|
||||
@@ -958,13 +958,30 @@ export function createServer(store: TaskStore, options?: ServerOptions): ReturnT
|
||||
// Preserve the raw payload buffer so signed endpoints (for example
|
||||
// /api/routines/:id/webhook and settings sync proxying) can verify HMAC
|
||||
// signatures and forward exact request bytes.
|
||||
app.use(express.json({
|
||||
/*
|
||||
FNXC:VoiceInput 2026-07-21-12:00:
|
||||
Voice chunks have a route-only 2 MiB parser. The global 100 KiB parser must skip only this
|
||||
endpoint (with or without Express's optional trailing slash) or it rejects before the voice
|
||||
error mapper; rawBody/HMAC behavior remains unchanged elsewhere.
|
||||
*/
|
||||
const jsonParser = express.json({
|
||||
verify: (req, _res, buf) => {
|
||||
if (buf.length > 0) {
|
||||
(req as express.Request & { rawBody?: Buffer }).rawBody = Buffer.from(buf);
|
||||
}
|
||||
},
|
||||
}));
|
||||
});
|
||||
app.use((req, res, next) => {
|
||||
// Express treats the trailing-slash spelling as the same route, so its parser boundary must,
|
||||
// too; no broader prefix is exempted from the global rawBody-preserving parser.
|
||||
if (req.path === "/api/voice/transcribe" || req.path === "/api/voice/transcribe/") return next();
|
||||
return jsonParser(req, res, (error) => {
|
||||
// Keep the established global 100 KiB rejection observable as 413 instead of allowing
|
||||
// Express's parser error to fall through to the generic 500 handler.
|
||||
if ((error as { type?: string } | undefined)?.type === "entity.too.large") return res.status(413).json({ error: "payload-too-large" });
|
||||
return next(error);
|
||||
});
|
||||
});
|
||||
|
||||
// Daemon mode: bearer token authentication middleware
|
||||
// Auth is enabled when daemon option is provided OR FUSION_DAEMON_TOKEN env var is set.
|
||||
|
||||
171
packages/dashboard/src/stt/__tests__/model-manager.test.ts
Normal file
171
packages/dashboard/src/stt/__tests__/model-manager.test.ts
Normal file
@@ -0,0 +1,171 @@
|
||||
import { mkdir, mkdtemp, readFile, readdir, rename, rm, writeFile } from "node:fs/promises";
|
||||
import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { createHash } from "node:crypto";
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { createVoiceModelManager, parseSafeTarListing } from "../model-manager.js";
|
||||
|
||||
const payload = Buffer.from("verified archive bytes");
|
||||
const digest = createHash("sha256").update(payload).digest("hex");
|
||||
const asset = { url: "https://example.invalid/parakeet.tar", filename: "parakeet.tar", sha256: digest, expectedFiles: ["tokens.txt"] };
|
||||
const response = () => new Response(new ReadableStream({ start(controller) { controller.enqueue(payload); controller.close(); } }), { status: 200, headers: { "content-length": String(payload.length) } });
|
||||
const deferred = <T = void>() => { let resolve!: (value: T) => void; const promise = new Promise<T>((done) => { resolve = done; }); return { promise, resolve }; };
|
||||
const safeListing = "-rw-r--r-- root/root 6 2026-01-01 00:00 tokens.txt";
|
||||
|
||||
describe("voice model manager", () => {
|
||||
it.each([
|
||||
"lrwxrwxrwx root/root 0 2026-01-01 00:00 model -> /tmp/x",
|
||||
"hrwxrwxrwx root/root 0 2026-01-01 00:00 model hard link to other",
|
||||
"crw-rw-rw- root/root 0 2026-01-01 00:00 device",
|
||||
"-rw-r--r-- root/root 0 2026-01-01 00:00 /absolute",
|
||||
"-rw-r--r-- root/root 0 2026-01-01 00:00 ../escape",
|
||||
"-rw-r--r-- root wheel 6 Jan 1 00:00 2026 /absolute-bsd",
|
||||
])("rejects unsafe tar metadata: %s", (listing) => expect(parseSafeTarListing(listing)).toMatchObject({ safe: false }));
|
||||
|
||||
it("streams a verified archive, extracts only after a safe listing, and writes a manifest", async () => {
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), "voice-model-"));
|
||||
const extract = vi.fn(async (_archive: string, staging: string) => { await writeFile(join(staging, "tokens.txt"), "tokens"); });
|
||||
const manager = createVoiceModelManager({ cacheDir, asset, fetch: vi.fn(async () => response()) as typeof fetch, listArchive: async () => "-rw-r--r-- root/root 6 2026-01-01 00:00 tokens.txt", extract });
|
||||
expect(manager.scheduleDownload().state.status).toBe("downloading");
|
||||
await manager.download();
|
||||
expect(await manager.getState()).toMatchObject({ status: "installed", checksumVerified: true });
|
||||
expect(extract).toHaveBeenCalledOnce();
|
||||
expect(JSON.parse(await readFile(join(cacheDir, "model", "manifest.json"), "utf8"))).toMatchObject({ sha256: digest });
|
||||
});
|
||||
|
||||
it("reports a corrupt final model directory as incomplete instead of absent", async () => {
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), "voice-model-"));
|
||||
await mkdir(join(cacheDir, "model"));
|
||||
await writeFile(join(cacheDir, "model", "manifest.json"), JSON.stringify({ sha256: "stale" }));
|
||||
const manager = createVoiceModelManager({ cacheDir, asset });
|
||||
await expect(manager.getState()).resolves.toMatchObject({ status: "error", errorReason: "incomplete-install" });
|
||||
});
|
||||
|
||||
it("publishes filesystem-inspected terminal states so polling snapshots stay current", async () => {
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), "voice-inspected-state-"));
|
||||
await mkdir(join(cacheDir, "model"));
|
||||
await writeFile(join(cacheDir, "model", "tokens.txt"), "tokens");
|
||||
await writeFile(join(cacheDir, "model", "manifest.json"), JSON.stringify({ sha256: digest }));
|
||||
const manager = createVoiceModelManager({ cacheDir, asset });
|
||||
const observed: string[] = [];
|
||||
manager.subscribe((next) => observed.push(next.status));
|
||||
|
||||
await expect(manager.getState()).resolves.toMatchObject({ status: "installed" });
|
||||
expect(manager.peekState()).toMatchObject({ status: "installed" });
|
||||
await rm(join(cacheDir, "model"), { recursive: true });
|
||||
await expect(manager.getState()).resolves.toMatchObject({ status: "not-installed" });
|
||||
expect(manager.peekState()).toMatchObject({ status: "not-installed" });
|
||||
expect(observed).toEqual(["installed", "not-installed"]);
|
||||
});
|
||||
|
||||
it("rolls back a prior installed model when staging promotion fails", async () => {
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), "voice-promotion-rollback-"));
|
||||
const modelDir = join(cacheDir, "model");
|
||||
await mkdir(modelDir);
|
||||
await writeFile(join(modelDir, "tokens.txt"), "previous");
|
||||
await writeFile(join(modelDir, "manifest.json"), JSON.stringify({ sha256: digest }));
|
||||
const move = vi.fn(async (from: string | Buffer | URL, to: string | Buffer | URL) => {
|
||||
if (String(from).includes(".staging-") && String(to) === modelDir) throw new Error("promotion failed");
|
||||
await rename(from, to);
|
||||
});
|
||||
const manager = createVoiceModelManager({
|
||||
cacheDir, asset, rename: move as typeof rename,
|
||||
fetch: vi.fn(async () => response()) as typeof fetch,
|
||||
listArchive: async () => safeListing,
|
||||
extract: async (_archive, staging) => { await writeFile(join(staging, "tokens.txt"), "replacement"); },
|
||||
});
|
||||
|
||||
await manager.download();
|
||||
expect(manager.peekState()).toMatchObject({ status: "error", errorReason: "extraction-failed" });
|
||||
await expect(readFile(join(modelDir, "tokens.txt"), "utf8")).resolves.toBe("previous");
|
||||
await expect(readFile(join(modelDir, "manifest.json"), "utf8")).resolves.toContain(digest);
|
||||
expect((await readdir(cacheDir)).some((entry) => entry.startsWith(".backup-"))).toBe(false);
|
||||
});
|
||||
|
||||
it("does not fetch or expose a cached model when its digest is unpinned", async () => {
|
||||
const fetch = vi.fn();
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), "voice-unpinned-"));
|
||||
await mkdir(join(cacheDir, "model"));
|
||||
await writeFile(join(cacheDir, "model", "manifest.json"), JSON.stringify({ sha256: null }));
|
||||
const manager = createVoiceModelManager({ asset: { ...asset, sha256: null }, cacheDir, fetch });
|
||||
expect(manager.scheduleDownload()).toMatchObject({ accepted: false, state: { errorReason: "checksum-unpinned" } });
|
||||
await expect(manager.getState()).resolves.toMatchObject({ status: "error", errorReason: "checksum-unpinned", checksumVerified: false });
|
||||
expect(fetch).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it("labels unsafe listing and extractor failures without promotion", async () => {
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), "voice-model-"));
|
||||
const manager = createVoiceModelManager({ cacheDir, asset, fetch: vi.fn(async () => response()) as typeof fetch, listArchive: async () => "lrwxrwxrwx root/root 0 2026-01-01 00:00 model -> /tmp/x" });
|
||||
await manager.download();
|
||||
expect(manager.peekState()).toMatchObject({ status: "error", errorReason: "unsafe-archive" });
|
||||
});
|
||||
|
||||
it.each(["fetch", "hash", "extract", "promotion"] as const)("fences deletion synchronously during %s without promotion", async (window) => {
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), `voice-race-${window}-`));
|
||||
const gate = deferred();
|
||||
let observedSignal: AbortSignal | undefined;
|
||||
const extract = vi.fn(async (_archive: string, staging: string, signal: AbortSignal) => {
|
||||
observedSignal = signal;
|
||||
await writeFile(join(staging, "tokens.txt"), "tokens");
|
||||
if (window === "extract") await gate.promise;
|
||||
});
|
||||
const fetch = vi.fn(async (_url: string, init?: RequestInit) => {
|
||||
observedSignal = init?.signal as AbortSignal;
|
||||
if (window !== "fetch") return response();
|
||||
return new Response(new ReadableStream({ async start(controller) { controller.enqueue(payload); await gate.promise; controller.close(); } }), { status: 200 });
|
||||
}) as typeof globalThis.fetch;
|
||||
const manager = createVoiceModelManager({
|
||||
cacheDir, asset, fetch, extract,
|
||||
listArchive: async (_archive, signal) => { observedSignal = signal; if (window === "hash") await gate.promise; return safeListing; },
|
||||
beforePromote: async (signal) => { observedSignal = signal; if (window === "promotion") await gate.promise; },
|
||||
});
|
||||
manager.scheduleDownload();
|
||||
await vi.waitFor(() => expect(observedSignal).toBeDefined());
|
||||
const removal = manager.remove();
|
||||
// Phase A is deliberately synchronous, before any controllable worker window releases.
|
||||
expect(observedSignal!.aborted).toBe(true);
|
||||
gate.resolve();
|
||||
await removal;
|
||||
expect(await manager.getState()).toMatchObject({ status: "not-installed" });
|
||||
expect((await readdir(cacheDir)).filter((entry) => entry === "model" || entry.startsWith(".partial-") || entry.startsWith(".staging-"))).toEqual([]);
|
||||
expect(extract).toHaveBeenCalledTimes(window === "fetch" || window === "hash" ? 0 : 1);
|
||||
});
|
||||
|
||||
it("queues a later download behind remove cleanup without deleting its generation", async () => {
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), "voice-cleanup-barrier-"));
|
||||
const cleanup = deferred();
|
||||
const fetch = vi.fn(async () => response()) as typeof fetch;
|
||||
const manager = createVoiceModelManager({ cacheDir, asset, fetch, listArchive: async () => safeListing, extract: async (_archive, staging) => { await writeFile(join(staging, "tokens.txt"), "tokens"); }, beforeCleanup: async () => cleanup.promise });
|
||||
const removing = manager.remove();
|
||||
const scheduled = manager.scheduleDownload();
|
||||
expect(scheduled).toMatchObject({ accepted: true, state: { status: "queued" } });
|
||||
expect(fetch).not.toHaveBeenCalled();
|
||||
cleanup.resolve();
|
||||
await removing;
|
||||
await manager.download();
|
||||
expect(await manager.getState()).toMatchObject({ status: "installed", checksumVerified: true });
|
||||
expect(fetch).toHaveBeenCalledOnce();
|
||||
});
|
||||
|
||||
it("fences an in-flight download immediately before queued cleanup", async () => {
|
||||
const cacheDir = await mkdtemp(join(tmpdir(), "voice-model-"));
|
||||
let release!: () => void;
|
||||
const blocked = new Promise<void>((resolve) => { release = resolve; });
|
||||
let observedSignal: AbortSignal | undefined;
|
||||
const fetch = vi.fn(async (_url: string, init?: RequestInit) => {
|
||||
observedSignal = init?.signal as AbortSignal;
|
||||
return new Response(new ReadableStream({ async start(controller) { controller.enqueue(payload); await blocked; controller.enqueue(payload); controller.close(); } }), { status: 200 });
|
||||
}) as typeof globalThis.fetch;
|
||||
const manager = createVoiceModelManager({ cacheDir, asset, fetch, listArchive: async () => "-rw-r--r-- root/root 6 2026-01-01 00:00 tokens.txt" });
|
||||
manager.scheduleDownload();
|
||||
await vi.waitFor(() => expect(observedSignal).toBeDefined());
|
||||
const removal = manager.remove();
|
||||
// Phase A is synchronous: abort happens before either the stream or cleanup is released.
|
||||
expect(observedSignal!.aborted).toBe(true);
|
||||
release();
|
||||
await removal;
|
||||
expect(await manager.getState()).toMatchObject({ status: "not-installed" });
|
||||
expect(await readdir(cacheDir)).not.toContain("model");
|
||||
expect((await readdir(cacheDir)).filter((entry) => entry.startsWith(".partial-") || entry.startsWith(".staging-"))).toEqual([]);
|
||||
});
|
||||
});
|
||||
48
packages/dashboard/src/stt/__tests__/voice-stt.test.ts
Normal file
48
packages/dashboard/src/stt/__tests__/voice-stt.test.ts
Normal file
@@ -0,0 +1,48 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
import { createVoiceModelManager } from "../model-manager.js";
|
||||
import { createParakeetService } from "../parakeet-service.js";
|
||||
import { resolveVoiceLanguage, resolveVoiceModelId } from "../types.js";
|
||||
|
||||
describe("voice STT graceful degradation", () => {
|
||||
it("refuses an unpinned archive synchronously without fetching", () => {
|
||||
const fetch = vi.fn();
|
||||
const manager = createVoiceModelManager({ cacheDir: "/unused", fetch, asset: { url: "https://example.test/model", filename: "model.tar", sha256: null, expectedFiles: [] } });
|
||||
expect(manager.scheduleDownload()).toMatchObject({ accepted: false, state: { status: "error", errorReason: "checksum-unpinned" } });
|
||||
expect(fetch).not.toHaveBeenCalled();
|
||||
});
|
||||
it("keeps unknown model identifiers and languages out of runtime configuration", () => {
|
||||
expect(resolveVoiceModelId(undefined)).toEqual({ id: "parakeet-v3" });
|
||||
expect(resolveVoiceModelId("../unsafe")).toEqual({ unsupported: "../unsafe" });
|
||||
expect(resolveVoiceLanguage(undefined)).toEqual({ language: "en" });
|
||||
expect(resolveVoiceLanguage("fr")).toEqual({ unsupported: "fr" });
|
||||
});
|
||||
it("reports a missing native binding as unavailable without throwing", async () => {
|
||||
const manager = createVoiceModelManager({ cacheDir: "/unused" });
|
||||
const service = createParakeetService({ manager, loadBinding: async () => { throw new Error("ERR_MODULE_NOT_FOUND"); } });
|
||||
await expect(service.getRuntimeStatus()).resolves.toMatchObject({ status: "unavailable" });
|
||||
});
|
||||
|
||||
it("reports a loaded but incompatible native binding as unavailable", async () => {
|
||||
const manager = { getState: async () => ({ status: "installed" as const, installedPath: "/model" }) } as ReturnType<typeof createVoiceModelManager>;
|
||||
const service = createParakeetService({ manager, loadBinding: async () => ({}) });
|
||||
await expect(service.getRuntimeStatus()).resolves.toEqual({ status: "unavailable", unavailableReason: "OfflineRecognizer unavailable" });
|
||||
});
|
||||
|
||||
it("uses sherpa's OfflineRecognizer and stream API for incremental decoding", async () => {
|
||||
let decoded = false;
|
||||
const stream = { acceptWaveform: vi.fn(), free: vi.fn() };
|
||||
const recognizer = { createStream: vi.fn(() => stream), decode: vi.fn(() => { decoded = true; }), getResult: vi.fn((received) => ({ text: received === stream && decoded ? "fresh" : "stale" })), close: vi.fn() };
|
||||
const OfflineRecognizer = vi.fn(function () { return recognizer; });
|
||||
const manager = { getState: async () => ({ status: "installed" as const, installedPath: "/model" }) } as ReturnType<typeof createVoiceModelManager>;
|
||||
const service = createParakeetService({ manager, loadBinding: async () => ({ OfflineRecognizer: OfflineRecognizer as never }) });
|
||||
const session = await service.createSession({ modelId: "parakeet-v3", language: "en" });
|
||||
expect(session.acceptChunk(Buffer.from([0, 0]), { final: false })).toEqual({ partial: "fresh" });
|
||||
expect(OfflineRecognizer).toHaveBeenCalledWith(expect.objectContaining({ modelConfig: expect.objectContaining({ transducer: expect.any(Object), tokens: "/model/tokens.txt" }) }));
|
||||
expect(stream.acceptWaveform).toHaveBeenCalledWith(expect.objectContaining({ sampleRate: 16_000, samples: expect.any(Float32Array) }));
|
||||
expect(recognizer.decode).toHaveBeenCalledWith(stream);
|
||||
expect(recognizer.getResult).toHaveBeenCalledWith(stream);
|
||||
expect(recognizer.decode).toHaveBeenCalledBefore(recognizer.getResult);
|
||||
session.close();
|
||||
expect(stream.free).toHaveBeenCalledOnce();
|
||||
});
|
||||
});
|
||||
268
packages/dashboard/src/stt/model-manager.ts
Normal file
268
packages/dashboard/src/stt/model-manager.ts
Normal file
@@ -0,0 +1,268 @@
|
||||
import { createHash, randomUUID } from "node:crypto";
|
||||
import { mkdir, open, readFile, readdir, rename, rm, stat, writeFile } from "node:fs/promises";
|
||||
import { join, relative, resolve, sep } from "node:path";
|
||||
import { resolveGlobalDir, superviseSpawn } from "@fusion/core";
|
||||
import { PARAKEET_V3_ASSET, type VoiceModelAsset, type VoiceModelState } from "./types.js";
|
||||
|
||||
export interface VoiceModelManager {
|
||||
getState(): Promise<VoiceModelState>;
|
||||
peekState(): VoiceModelState;
|
||||
scheduleDownload(): { accepted: boolean; state: VoiceModelState };
|
||||
download(): Promise<VoiceModelState>;
|
||||
remove(): Promise<void>;
|
||||
subscribe(listener: (state: VoiceModelState) => void): () => void;
|
||||
}
|
||||
export interface VoiceModelManagerOptions {
|
||||
asset?: VoiceModelAsset;
|
||||
cacheDir?: string;
|
||||
fetch?: typeof globalThis.fetch;
|
||||
/** Test seam; production lists and extracts with async tar invocations. */
|
||||
extract?: (archive: string, staging: string, signal: AbortSignal) => Promise<void>;
|
||||
listArchive?: (archive: string, signal: AbortSignal) => Promise<string>;
|
||||
/** Deterministic test seam for the post-validation/pre-promotion cancellation fence. */
|
||||
beforePromote?: (signal: AbortSignal) => Promise<void>;
|
||||
/** Deterministic test seam for cleanup-barrier scheduling behavior. */
|
||||
beforeCleanup?: (generation: number) => Promise<void>;
|
||||
/** Test seam for verifying old-install rollback if promotion cannot complete. */
|
||||
rename?: typeof rename;
|
||||
}
|
||||
|
||||
/** Reject tar metadata entries that can escape or alias the staging directory. */
|
||||
export function parseSafeTarListing(listing: string): { safe: true } | { safe: false; reason: string } {
|
||||
for (const line of listing.split("\n").filter(Boolean)) {
|
||||
const type = line[0];
|
||||
if (type !== "-" && type !== "d") return { safe: false, reason: "unsafe entry type" };
|
||||
// GNU uses YYYY-MM-DD HH:MM while BSD tar uses Mon DD HH:MM YYYY. Do not
|
||||
// guess a column offset: owner/group fields vary, and guessing could turn an
|
||||
// absolute pathname into a harmless-looking date suffix.
|
||||
const timestamp = /\d{4}-\d{2}-\d{2}\s+\d{2}:\d{2}(?::\d{2})?\s+|(?:Jan|Feb|Mar|Apr|May|Jun|Jul|Aug|Sep|Oct|Nov|Dec)\s+\d{1,2}\s+\d{2}:\d{2}(?::\d{2})?\s+\d{4}\s+/.exec(line);
|
||||
if (!timestamp || timestamp.index === undefined) return { safe: false, reason: "unparseable tar metadata" };
|
||||
const path = line.slice(timestamp.index + timestamp[0].length);
|
||||
if (!path || /\s(?:->|link to|hard link to)\s/.test(path) || path.startsWith("/") || path.split(/[\\/]+/).includes("..")) return { safe: false, reason: "unsafe entry path" };
|
||||
}
|
||||
return { safe: true };
|
||||
}
|
||||
|
||||
/**
|
||||
* FNXC:VoiceInput 2026-07-21-17:20:
|
||||
* Voice models are on-demand, user-cache-only assets and lifecycle management remains available
|
||||
* while dictation is off. A pinned digest is mandatory. scheduleDownload publishes synchronously;
|
||||
* its worker waits behind cleanup, while remove fences the worker immediately and then serializes
|
||||
* generation-scoped cleanup. getState is authoritative and async; peekState is polling-only.
|
||||
*/
|
||||
export function createVoiceModelManager(options: VoiceModelManagerOptions = {}): VoiceModelManager {
|
||||
const asset = options.asset ?? PARAKEET_V3_ASSET;
|
||||
let configuredCacheDir = options.cacheDir;
|
||||
const cacheDir = () => configuredCacheDir ??= join(resolveGlobalDir(), "models", "parakeet-v3");
|
||||
const finalDir = () => join(cacheDir(), "model");
|
||||
let state: VoiceModelState = { status: "not-installed" };
|
||||
let generation = 0;
|
||||
let cleanupBarrier: Promise<void> = Promise.resolve();
|
||||
let cleanupPending = false;
|
||||
let active: { generation: number; controller: AbortController; settled: Promise<void> } | undefined;
|
||||
let work: Promise<VoiceModelState> = Promise.resolve(state);
|
||||
const listeners = new Set<(state: VoiceModelState) => void>();
|
||||
const publish = (next: VoiceModelState) => { state = next; for (const listener of listeners) listener(next); };
|
||||
const current = (g: number, signal: AbortSignal) => generation === g && !signal.aborted;
|
||||
const guarded = (g: number, signal: AbortSignal, next: VoiceModelState) => { if (current(g, signal)) publish(next); };
|
||||
const existsFile = async (file: string) => { try { const item = await stat(file); return item.isFile() && item.size > 0; } catch { return false; } };
|
||||
const removeOwn = async (...paths: Array<string | undefined>) => Promise.all(paths.filter(Boolean).map((path) => rm(path!, { recursive: true, force: true })));
|
||||
const inStaging = (staging: string, file: string) => {
|
||||
const target = resolve(staging, file);
|
||||
return relative(staging, target) !== "" && !relative(staging, target).startsWith(`..${sep}`) && !relative(staging, target).startsWith("../") ? target : undefined;
|
||||
};
|
||||
const capture = (stream: NodeJS.ReadableStream | null | undefined) => new Promise<string>((done) => {
|
||||
if (!stream) return done("");
|
||||
const chunks: Buffer[] = [];
|
||||
stream.on("data", (chunk) => chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(String(chunk))));
|
||||
stream.once("end", () => done(Buffer.concat(chunks).toString("utf8")));
|
||||
stream.once("close", () => done(Buffer.concat(chunks).toString("utf8")));
|
||||
});
|
||||
const runTar = async (args: string[], signal: AbortSignal) => {
|
||||
// FNXC:VoiceInput 2026-07-21-18:10: Archive inspection/extraction is a managed
|
||||
// child process so cancellation during model deletion reaches tar immediately.
|
||||
const supervised = superviseSpawn("tar", args, { stdio: ["ignore", "pipe", "pipe"], maxLifetimeMs: 120_000 });
|
||||
const onAbort = () => supervised.kill("SIGTERM");
|
||||
signal.addEventListener("abort", onAbort, { once: true });
|
||||
try {
|
||||
const [stdout, stderr, exit] = await Promise.all([capture(supervised.child.stdout), capture(supervised.child.stderr), supervised.waitExit()]);
|
||||
if (exit.code !== 0) throw new Error(stderr.trim() || `tar exited with ${exit.code ?? exit.signal ?? "an error"}`);
|
||||
return stdout;
|
||||
} finally { signal.removeEventListener("abort", onAbort); }
|
||||
};
|
||||
const listArchive = async (archive: string, signal: AbortSignal) => options.listArchive
|
||||
? options.listArchive(archive, signal)
|
||||
: runTar(["-tvf", archive], signal);
|
||||
const extractArchive = async (archive: string, staging: string, signal: AbortSignal) => {
|
||||
if (options.extract) return options.extract(archive, staging, signal);
|
||||
await runTar(["--no-same-owner", "-xf", archive, "-C", staging], signal);
|
||||
};
|
||||
const move = options.rename ?? rename;
|
||||
|
||||
const perform = async (g: number, controller: AbortController): Promise<VoiceModelState> => {
|
||||
const signal = controller.signal;
|
||||
let partial: string | undefined;
|
||||
let staging: string | undefined;
|
||||
try {
|
||||
await cleanupBarrier;
|
||||
if (!current(g, signal)) return state;
|
||||
await mkdir(cacheDir(), { recursive: true });
|
||||
if (!current(g, signal)) return state;
|
||||
partial = join(cacheDir(), `.partial-${g}-${randomUUID()}`);
|
||||
staging = join(cacheDir(), `.staging-${g}-${randomUUID()}`);
|
||||
const file = await open(partial, "wx");
|
||||
try {
|
||||
const response = await (options.fetch ?? globalThis.fetch)(asset.url, { signal });
|
||||
if (!response.ok || !response.body) throw Object.assign(new Error("download failed"), { voiceReason: "network" });
|
||||
if (!current(g, signal)) return state;
|
||||
const total = Number(response.headers.get("content-length")) || undefined;
|
||||
const hash = createHash("sha256"); let bytes = 0;
|
||||
for await (const chunk of response.body as AsyncIterable<Uint8Array>) {
|
||||
if (!current(g, signal)) return state;
|
||||
const buffer = Buffer.from(chunk); hash.update(buffer); await file.write(buffer);
|
||||
if (!current(g, signal)) return state;
|
||||
bytes += buffer.length; guarded(g, signal, { status: "downloading", progress: total ? bytes / total : undefined, bytesDownloaded: bytes, totalBytes: total });
|
||||
}
|
||||
if (!current(g, signal)) return state;
|
||||
if (hash.digest("hex") !== asset.sha256) {
|
||||
guarded(g, signal, { status: "error", errorReason: "checksum-mismatch", checksumVerified: false });
|
||||
return state;
|
||||
}
|
||||
} finally { await file.close(); }
|
||||
if (!current(g, signal)) return state;
|
||||
const listing = await listArchive(partial, signal);
|
||||
if (!current(g, signal)) return state;
|
||||
const safe = parseSafeTarListing(listing);
|
||||
if (!safe.safe) { guarded(g, signal, { status: "error", errorReason: "unsafe-archive", errorMessage: safe.reason }); return state; }
|
||||
await mkdir(staging);
|
||||
await extractArchive(partial, staging, signal);
|
||||
if (!current(g, signal)) return state;
|
||||
for (const expected of asset.expectedFiles) {
|
||||
const target = inStaging(staging, expected);
|
||||
if (!target || !(await existsFile(target))) { guarded(g, signal, { status: "error", errorReason: "incomplete-install" }); return state; }
|
||||
}
|
||||
if (!current(g, signal)) return state;
|
||||
// Write the manifest before promotion so the promoted directory is always a complete
|
||||
// install. The previous final directory is retained as a rollback candidate until this
|
||||
// new staging directory has replaced it.
|
||||
await writeFile(join(staging, "manifest.json"), JSON.stringify({ filename: asset.filename, sha256: asset.sha256, installedAt: new Date().toISOString(), expectedFiles: asset.expectedFiles }));
|
||||
await options.beforePromote?.(signal);
|
||||
if (!current(g, signal)) return state;
|
||||
|
||||
/*
|
||||
* FNXC:VoiceInput 2026-07-21-19:05:
|
||||
* A failed refresh must not destroy a working installed model. Move the old installation
|
||||
* aside, promote the fully validated staging tree, and restore the old tree if promotion
|
||||
* fails. Generation-fenced cleanup owns stale backup trees after deletion.
|
||||
*/
|
||||
const backup = join(cacheDir(), `.backup-${g}-${randomUUID()}`);
|
||||
let movedPrevious = false;
|
||||
try {
|
||||
try {
|
||||
await move(finalDir(), backup);
|
||||
movedPrevious = true;
|
||||
} catch (error) {
|
||||
if ((error as NodeJS.ErrnoException).code !== "ENOENT") throw error;
|
||||
}
|
||||
if (!current(g, signal)) return state;
|
||||
await move(staging, finalDir()); staging = undefined;
|
||||
} catch (error) {
|
||||
if (movedPrevious) {
|
||||
try { await move(backup, finalDir()); } catch { /* Preserve the original failure. */ }
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
if (!current(g, signal)) return state;
|
||||
if (movedPrevious) await rm(backup, { recursive: true, force: true }).catch(() => undefined);
|
||||
guarded(g, signal, { status: "installed", checksumVerified: true, installedPath: finalDir() });
|
||||
return state;
|
||||
} catch (error) {
|
||||
if (current(g, signal)) {
|
||||
const reason = signal.aborted ? "cancelled" : (error as { voiceReason?: string }).voiceReason === "network" ? "network" : "extraction-failed";
|
||||
guarded(g, signal, { status: "error", errorReason: reason, errorMessage: error instanceof Error ? error.message : String(error) });
|
||||
}
|
||||
return state;
|
||||
} finally {
|
||||
await removeOwn(partial, staging);
|
||||
if (active?.generation === g) active = undefined;
|
||||
}
|
||||
};
|
||||
|
||||
const scheduleDownload = () => {
|
||||
if (!asset.sha256) { publish({ status: "error", errorReason: "checksum-unpinned", checksumVerified: false }); return { accepted: false, state }; }
|
||||
if (state.status === "queued" || state.status === "downloading") return { accepted: true, state };
|
||||
const g = ++generation;
|
||||
const controller = new AbortController();
|
||||
publish(cleanupPending ? { status: "queued" } : { status: "downloading", progress: 0, bytesDownloaded: 0 });
|
||||
const settled = perform(g, controller);
|
||||
active = { generation: g, controller, settled: settled.then(() => undefined) };
|
||||
work = settled;
|
||||
return { accepted: true, state };
|
||||
};
|
||||
|
||||
return {
|
||||
peekState: () => state,
|
||||
subscribe(listener) { listeners.add(listener); return () => listeners.delete(listener); },
|
||||
scheduleDownload,
|
||||
async download() { const scheduled = scheduleDownload(); return scheduled.accepted ? work : scheduled.state; },
|
||||
async getState() {
|
||||
// An unpinned asset is never installable, even if stale files or a legacy manifest exist.
|
||||
// The mandatory checksum gate applies to status inspection as well as scheduling.
|
||||
if (!asset.sha256) {
|
||||
const unpinned = { status: "error" as const, errorReason: "checksum-unpinned" as const, checksumVerified: false };
|
||||
publish(unpinned);
|
||||
return unpinned;
|
||||
}
|
||||
if (state.status === "queued" || state.status === "downloading") return state;
|
||||
try {
|
||||
let finalExists = false;
|
||||
try { finalExists = (await stat(finalDir())).isDirectory(); } catch { /* no install yet */ }
|
||||
if (!finalExists) {
|
||||
if (state.status === "error") return state;
|
||||
const absent = { status: "not-installed" as const };
|
||||
publish(absent);
|
||||
return absent;
|
||||
}
|
||||
const manifest = JSON.parse(await readFile(join(finalDir(), "manifest.json"), "utf8"));
|
||||
if (manifest.sha256 !== asset.sha256) throw new Error("stale manifest digest");
|
||||
for (const expected of asset.expectedFiles) if (!(await existsFile(join(finalDir(), expected)))) throw new Error("missing expected model file");
|
||||
const installed = { status: "installed" as const, checksumVerified: true, installedPath: finalDir() };
|
||||
publish(installed);
|
||||
return installed;
|
||||
} catch (cause) {
|
||||
// A final directory is evidence of an interrupted/corrupt installation, not
|
||||
// an absent model. Surface an actionable state so operators can delete it.
|
||||
const corrupt = { status: "error" as const, errorReason: "incomplete-install" as const, errorMessage: cause instanceof Error ? cause.message : "incomplete model install" };
|
||||
publish(corrupt);
|
||||
return corrupt;
|
||||
} finally {
|
||||
try {
|
||||
for (const entry of await readdir(cacheDir())) {
|
||||
const own = /^\.(partial|staging|backup)-(\d+)-/.exec(entry);
|
||||
if (own && Number(own[2]) !== active?.generation) await rm(join(cacheDir(), entry), { recursive: true, force: true });
|
||||
}
|
||||
} catch { /* Cache absence and corruption must not break dashboard boot. */ }
|
||||
}
|
||||
},
|
||||
remove() {
|
||||
const g = ++generation;
|
||||
const stale = active;
|
||||
if (stale) { stale.controller.abort(); active = undefined; }
|
||||
cleanupPending = true;
|
||||
const previous = cleanupBarrier;
|
||||
const pending = cleanupBarrier = previous.then(async () => {
|
||||
await stale?.settled.catch(() => undefined);
|
||||
await options.beforeCleanup?.(g);
|
||||
await rm(finalDir(), { recursive: true, force: true });
|
||||
try {
|
||||
for (const entry of await readdir(cacheDir())) {
|
||||
const match = /^\.(partial|staging|backup)-(\d+)-/.exec(entry);
|
||||
if (match && Number(match[2]) <= g) await rm(join(cacheDir(), entry), { recursive: true, force: true });
|
||||
}
|
||||
} catch { /* absent cache is already clean */ }
|
||||
if (generation === g) publish({ status: "not-installed" });
|
||||
}).finally(() => { if (cleanupBarrier === pending) cleanupPending = false; });
|
||||
return pending;
|
||||
},
|
||||
};
|
||||
}
|
||||
77
packages/dashboard/src/stt/parakeet-service.ts
Normal file
77
packages/dashboard/src/stt/parakeet-service.ts
Normal file
@@ -0,0 +1,77 @@
|
||||
import type { VoiceModelManager } from "./model-manager.js";
|
||||
import { resolveVoiceLanguage, type VoiceModelId, type VoiceRuntimeStatus } from "./types.js";
|
||||
|
||||
export class VoiceInputError extends Error { constructor(public readonly code: "unsupported-language" | "invalid-audio" | "unavailable", message: string) { super(message); } }
|
||||
export interface ParakeetService { getRuntimeStatus(): Promise<{ status: VoiceRuntimeStatus; unavailableReason?: string }>; createSession(options: { modelId: VoiceModelId; language: string }): Promise<ParakeetSession>; }
|
||||
export interface ParakeetSession { acceptChunk(pcm: Int16Array | Buffer, options: { final: boolean }): { partial?: string; text?: string; final?: true }; finish(): { text: string }; close(): void; }
|
||||
interface SherpaStream { acceptWaveform(options: { sampleRate: number; samples: Float32Array }): void; free?(): void; close?(): void; }
|
||||
interface SherpaRecognizer { createStream(): SherpaStream; getResult(stream: SherpaStream): { text?: string }; decode(stream: SherpaStream): void; free?(): void; close?(): void; }
|
||||
interface SherpaOfflineRecognizerConstructor { new(config: { modelConfig: { transducer: { encoder: string; decoder: string; joiner: string }; tokens: string; numThreads: number; provider: string; debug: boolean }; decodingMethod: string; maxActivePaths: number }): SherpaRecognizer; }
|
||||
interface SherpaBinding { OfflineRecognizer?: SherpaOfflineRecognizerConstructor; }
|
||||
export interface ParakeetServiceOptions { manager: VoiceModelManager; loadBinding?: () => Promise<SherpaBinding>; }
|
||||
|
||||
/**
|
||||
* FNXC:VoiceInput 2026-07-21-17:20:
|
||||
* Voice is opt-in and the sherpa addon is an optional, lazy runtime. Its fixed input is 16 kHz,
|
||||
* mono signed-16-bit little-endian PCM. Resolved model/language values are checked before use;
|
||||
* missing addon or model reports unavailable rather than preventing dashboard or engine boot.
|
||||
*/
|
||||
export function createParakeetService(options: ParakeetServiceOptions): ParakeetService {
|
||||
let bindingPromise: Promise<SherpaBinding> | undefined;
|
||||
const binding = () => bindingPromise ??= (options.loadBinding ? options.loadBinding() : new Function("specifier", "return import(specifier)")("sherpa-onnx-node") as Promise<SherpaBinding>);
|
||||
const getRuntimeStatus = async () => {
|
||||
const model = await options.manager.getState();
|
||||
if (model.status !== "installed" || !model.installedPath) return { status: "unavailable" as const, unavailableReason: model.errorReason ?? model.status };
|
||||
try {
|
||||
// A module resolving is not sufficient: a platform-mismatched or incompatible addon
|
||||
// can load without exporting the recognizer API required for transcription.
|
||||
if (!(await binding()).OfflineRecognizer) return { status: "unavailable" as const, unavailableReason: "OfflineRecognizer unavailable" };
|
||||
return { status: "available" as const };
|
||||
} catch (error) { return { status: "unavailable" as const, unavailableReason: error instanceof Error ? error.message : "runtime-unavailable" }; }
|
||||
};
|
||||
return {
|
||||
getRuntimeStatus,
|
||||
async createSession({ modelId: _modelId, language }) {
|
||||
if ("unsupported" in resolveVoiceLanguage(language)) throw new VoiceInputError("unsupported-language", "Unsupported language");
|
||||
const model = await options.manager.getState();
|
||||
const status = await getRuntimeStatus();
|
||||
if (status.status !== "available" || !model.installedPath) throw new VoiceInputError("unavailable", status.unavailableReason ?? "unavailable");
|
||||
const addon = await binding();
|
||||
const OfflineRecognizer = addon.OfflineRecognizer;
|
||||
if (!OfflineRecognizer) throw new VoiceInputError("unavailable", "OfflineRecognizer unavailable");
|
||||
// FNXC:VoiceInput 2026-07-21-20:30: sherpa-onnx-node's offline API owns waveform
|
||||
// ingestion on a stream, then decodes and reads that stream through its recognizer.
|
||||
// Keep the native config shaped as modelConfig; flat model-path configs are not accepted.
|
||||
const recognizer = new OfflineRecognizer({
|
||||
modelConfig: {
|
||||
transducer: {
|
||||
encoder: `${model.installedPath}/encoder.int8.onnx`,
|
||||
decoder: `${model.installedPath}/decoder.int8.onnx`,
|
||||
joiner: `${model.installedPath}/joiner.int8.onnx`,
|
||||
},
|
||||
tokens: `${model.installedPath}/tokens.txt`,
|
||||
numThreads: 1,
|
||||
provider: "cpu",
|
||||
debug: false,
|
||||
},
|
||||
decodingMethod: "greedy_search",
|
||||
maxActivePaths: 4,
|
||||
});
|
||||
const stream = recognizer.createStream();
|
||||
const decode = () => { recognizer.decode(stream); return recognizer.getResult(stream).text ?? ""; };
|
||||
return {
|
||||
acceptChunk(pcm, { final }) {
|
||||
const buffer = Buffer.isBuffer(pcm) ? pcm : Buffer.from(pcm.buffer, pcm.byteOffset, pcm.byteLength);
|
||||
if (!buffer.byteLength || buffer.byteLength % 2) throw new VoiceInputError("invalid-audio", "PCM must be signed 16-bit samples");
|
||||
const floats = new Float32Array(buffer.byteLength / 2);
|
||||
for (let i = 0; i < floats.length; i++) floats[i] = buffer.readInt16LE(i * 2) / 32768;
|
||||
stream.acceptWaveform({ sampleRate: 16_000, samples: floats });
|
||||
const text = decode();
|
||||
return final ? { text, final: true as const } : { partial: text };
|
||||
},
|
||||
finish: () => ({ text: decode() }),
|
||||
close: () => { stream.free?.(); stream.close?.(); recognizer.free?.(); recognizer.close?.(); },
|
||||
};
|
||||
},
|
||||
};
|
||||
}
|
||||
14
packages/dashboard/src/stt/types.ts
Normal file
14
packages/dashboard/src/stt/types.ts
Normal file
@@ -0,0 +1,14 @@
|
||||
export type VoiceModelStatus = "not-installed" | "queued" | "downloading" | "installed" | "error";
|
||||
export type VoiceModelErrorReason = "checksum-unpinned" | "checksum-mismatch" | "network" | "extraction-failed" | "unsafe-archive" | "incomplete-install" | "cancelled";
|
||||
export interface VoiceModelState { status: VoiceModelStatus; progress?: number; bytesDownloaded?: number; totalBytes?: number; errorReason?: VoiceModelErrorReason; errorMessage?: string; checksumVerified?: boolean; installedPath?: string; }
|
||||
export type VoiceRuntimeStatus = "available" | "unavailable";
|
||||
export type VoiceModelId = "parakeet-v3";
|
||||
export const DEFAULT_VOICE_MODEL_ID: VoiceModelId = "parakeet-v3";
|
||||
export const SUPPORTED_VOICE_LANGUAGES = ["en"] as const;
|
||||
export const DEFAULT_VOICE_LANGUAGE = "en";
|
||||
export interface VoiceModelAsset { url: string; filename: string; sha256: string | null; expectedFiles: string[]; }
|
||||
/** The upstream release page has no published digest for a concrete v3 archive yet. */
|
||||
export const PARAKEET_V3_ASSET: VoiceModelAsset = { url: "https://github.com/k2-fsa/sherpa-onnx/releases/tag/asr-models", filename: "upstream-pending-verification", sha256: null, expectedFiles: [] };
|
||||
export const VOICE_MODEL_REGISTRY: Record<VoiceModelId, VoiceModelAsset> = { "parakeet-v3": PARAKEET_V3_ASSET };
|
||||
export function resolveVoiceModelId(raw?: string): { id: VoiceModelId } | { unsupported: string } { return raw === undefined || raw === DEFAULT_VOICE_MODEL_ID ? { id: DEFAULT_VOICE_MODEL_ID } : { unsupported: raw }; }
|
||||
export function resolveVoiceLanguage(raw?: string): { language: string } | { unsupported: string } { return raw === undefined || raw === DEFAULT_VOICE_LANGUAGE ? { language: DEFAULT_VOICE_LANGUAGE } : { unsupported: raw }; }
|
||||
65
pnpm-lock.yaml
generated
65
pnpm-lock.yaml
generated
@@ -453,6 +453,10 @@ importers:
|
||||
vitest:
|
||||
specifier: ^4.1.0
|
||||
version: 4.1.8(@opentelemetry/api@1.9.0)(@types/node@25.5.2)(@vitest/coverage-v8@4.1.8)(happy-dom@20.10.1)(jsdom@29.0.1)(vite@6.4.1(@types/node@25.5.2)(jiti@2.7.0)(tsx@4.21.0)(yaml@2.9.0))
|
||||
optionalDependencies:
|
||||
sherpa-onnx-node:
|
||||
specifier: ^1.13.4
|
||||
version: 1.13.4
|
||||
|
||||
packages/desktop:
|
||||
dependencies:
|
||||
@@ -7291,6 +7295,39 @@ packages:
|
||||
resolution: {integrity: sha512-ObmnIF4hXNg1BqhnHmgbDETF8dLPCggZWBjkQfhZpbszZnYur5DUljTcCHii5LC3J5E0yeO/1LIMyH+UvHQgyw==}
|
||||
engines: {node: '>= 0.4'}
|
||||
|
||||
sherpa-onnx-darwin-arm64@1.13.4:
|
||||
resolution: {integrity: sha512-QcYKzyrTzGSx6aKCD6hUODgRS1LetqfG57Z/+i5LCyfMlrgCvDc1lRcl9cdB+TozBsLha9QwLTlI0vmDcf5JKg==}
|
||||
cpu: [arm64]
|
||||
os: [darwin]
|
||||
|
||||
sherpa-onnx-darwin-x64@1.13.4:
|
||||
resolution: {integrity: sha512-6RGeis9K9gV/UQWOgd6Rf3iqXr2/YsBQswxHaCR4hrYkHfEIpHMfFmRWLt6nJJCOWgYW2xFxEd9yzjrafAV/Pw==}
|
||||
cpu: [x64]
|
||||
os: [darwin]
|
||||
|
||||
sherpa-onnx-linux-arm64@1.13.4:
|
||||
resolution: {integrity: sha512-RMjMRqT82BgTXypNNGmLe6ZFYhc3WEvnAGl3DdkK7qB/kuXwkL3iHhV31wAecbnWPsnEpUoD+8cFovWSBzsCuw==}
|
||||
cpu: [arm64]
|
||||
os: [linux]
|
||||
|
||||
sherpa-onnx-linux-x64@1.13.4:
|
||||
resolution: {integrity: sha512-WZh5NCkGPFHHpYSd78iN4OnmxQeSTGyt9uZskH+im/NFHQ7elQ7B0sLzCMeRpvJxiIKvd9C6WxIJ4hYaxClfsQ==}
|
||||
cpu: [x64]
|
||||
os: [linux]
|
||||
|
||||
sherpa-onnx-node@1.13.4:
|
||||
resolution: {integrity: sha512-jHWWdY9f0dbvJpdsfR/C4WNhM57+jT9Os2RyC4/qUz8HKCtj2VYFmJ9s7tue9OCFRAmdoGBkgkqD6aBZ3I8BKw==}
|
||||
|
||||
sherpa-onnx-win-ia32@1.13.4:
|
||||
resolution: {integrity: sha512-/JbPjldrfNv+t+uIS3MlkuhfIf5l3FHUGkRC2oRXgjRqOaVmEyP3vLlQ7dTa4J7raG5oB8c3GoPjuSWSqT9GOQ==}
|
||||
cpu: [ia32]
|
||||
os: [win32]
|
||||
|
||||
sherpa-onnx-win-x64@1.13.4:
|
||||
resolution: {integrity: sha512-R0PWby1VxC14TDZPq7GcfSyXSY6SAFO8Y4JwdCdqouFmeXkZ1L7Is9m98C9KxQ0dN7ZtDzhAmE/43FUs/elXRQ==}
|
||||
cpu: [x64]
|
||||
os: [win32]
|
||||
|
||||
side-channel-list@1.0.0:
|
||||
resolution: {integrity: sha512-FCLHtRD/gnpCiCHEiJLOwdmFP+wzCmDEkc9y7NsYxeF4u7Btsn1ZuwgwJGxImImHicJArLP4R0yX4c2KCrMrTA==}
|
||||
engines: {node: '>= 0.4'}
|
||||
@@ -15224,6 +15261,34 @@ snapshots:
|
||||
|
||||
shell-quote@1.8.3: {}
|
||||
|
||||
sherpa-onnx-darwin-arm64@1.13.4:
|
||||
optional: true
|
||||
|
||||
sherpa-onnx-darwin-x64@1.13.4:
|
||||
optional: true
|
||||
|
||||
sherpa-onnx-linux-arm64@1.13.4:
|
||||
optional: true
|
||||
|
||||
sherpa-onnx-linux-x64@1.13.4:
|
||||
optional: true
|
||||
|
||||
sherpa-onnx-node@1.13.4:
|
||||
optionalDependencies:
|
||||
sherpa-onnx-darwin-arm64: 1.13.4
|
||||
sherpa-onnx-darwin-x64: 1.13.4
|
||||
sherpa-onnx-linux-arm64: 1.13.4
|
||||
sherpa-onnx-linux-x64: 1.13.4
|
||||
sherpa-onnx-win-ia32: 1.13.4
|
||||
sherpa-onnx-win-x64: 1.13.4
|
||||
optional: true
|
||||
|
||||
sherpa-onnx-win-ia32@1.13.4:
|
||||
optional: true
|
||||
|
||||
sherpa-onnx-win-x64@1.13.4:
|
||||
optional: true
|
||||
|
||||
side-channel-list@1.0.0:
|
||||
dependencies:
|
||||
es-errors: 1.3.0
|
||||
|
||||
Reference in New Issue
Block a user