mirror of
https://github.com/openclaw/openclaw.git
synced 2026-10-03 01:29:56 +00:00
* fix(voice): keep audio flowing while the Gateway is busy Move Discord voice transport, packet pacing, codecs and capture deadlines into a session-owned worker. Run GPT Live WebRTC and WebSocket media in workers and connect continuous playback with a bounded direct audio port. Preserve authorization, selected-agent work, recording receipts and transcript ownership on main, with generation fencing and ordered teardown. Validated affected tests, the capture-finalization red/green regression, source and compiled worker lifecycle, root and plugin builds, type/lint/format and line-cap checks. Co-authored-by: steipete <58493+steipete@users.noreply.github.com> * fix(voice): keep audio flowing while the Gateway is busy Worked on by: - @steipete Co-authored-by: steipete <58493+steipete@users.noreply.github.com> OpenClaw-Publication: 58c87505-0064-4391-994e-57587279bcca * fix(voice): use explicit Node message transfer lists Replace six Node postMessage lint suppressions with explicit empty transfer lists. Preserve the Node messaging contract and the unchanged production suppression allowlist. Reproduced the original CI shard plan before the fix; the same plan, focused worker messaging tests, and type-aware lint pass afterward. Co-authored-by: steipete <58493+steipete@users.noreply.github.com> * fix(openai): keep worker runtime out of cold voice catalogs Load the default media-socket factory only through abort-aware connection admission. Preserve injected factories and WebRTC routing, and centralize unchanged request-ID construction in the wire owner. Keep the cold-catalog contract and line cap intact. Co-authored-by: steipete <58493+steipete@users.noreply.github.com> * refactor(openai): keep socket contracts below connection admission Move unchanged socket interfaces into the existing shared leaf and migrate all private consumers, preserving lazy worker loading without a type import cycle. Co-authored-by: steipete <58493+steipete@users.noreply.github.com> * fix(discord): preserve direct speech ownership through playback Track exact speech with an epoch-scoped shared playback witness. Order completion behind bounded direct-port flush receipts, including no-audio responses, and prevent stale idle or flush events from retiring newer speech. Cover original failures and the pending-flush race with deterministic production-owner regressions. Co-authored-by: steipete <58493+steipete@users.noreply.github.com> * fix(discord): preserve audio ownership across worker retirement Worked on by: - @steipete Co-authored-by: steipete <58493+steipete@users.noreply.github.com> OpenClaw-Publication: 6294ffa7-0c2f-456b-9baf-3dcdab60efc8 * refactor(discord): keep the worker status slot private Worked on by: - @steipete Co-authored-by: steipete <58493+steipete@users.noreply.github.com> OpenClaw-Publication: dde5cff1-21c6-4bfa-825e-47c2f1ec2393 --------- Co-authored-by: steipete <58493+steipete@users.noreply.github.com>
296 lines
13 KiB
TypeScript
296 lines
13 KiB
TypeScript
import fs from "node:fs/promises";
|
|
import { PassThrough } from "node:stream";
|
|
import { setImmediate } from "node:timers/promises";
|
|
import { createDeferred } from "openclaw/plugin-sdk/extension-shared";
|
|
import { transcribeAudioFile } from "openclaw/plugin-sdk/media-understanding-runtime";
|
|
import { vi } from "vitest";
|
|
import { loadDiscordVoiceTestHarness } from "../extensions/discord/test-api.js";
|
|
import type { MediaUnderstandingModelConfig } from "../src/config/types.tools.js";
|
|
import { createPluginMetadataSnapshotFixture } from "../src/plugins/plugin-metadata.test-support.js";
|
|
import { createEmptyPluginRegistry } from "../src/plugins/registry-empty.js";
|
|
import { withPluginRuntimeGenerationScope } from "../src/plugins/runtime/generation-scope.js";
|
|
|
|
const { defineDiscordVoiceTests } = await loadDiscordVoiceTestHarness();
|
|
const cli = vi.hoisted(() => vi.fn());
|
|
vi.mock("../src/process/exec.js", () => ({ runExec: cli }));
|
|
|
|
defineDiscordVoiceTests(
|
|
({
|
|
expect,
|
|
it,
|
|
createManager,
|
|
makeVoiceConfig,
|
|
getSessionEntry,
|
|
getSessionConnection,
|
|
handleSpeakingStart,
|
|
startTranscripts,
|
|
decodeOpusStreamChunksMock,
|
|
transcribeAudioFileMock,
|
|
agentCommandMock,
|
|
controlRealtimeVoiceAgentRunMock,
|
|
loggerWarnMock,
|
|
receiveRecordedSpeech,
|
|
}) => {
|
|
it.each(
|
|
["conversation", "control"].flatMap((dispatch) =>
|
|
["oversized", "empty CLI", "fallback"].map((opening) => ({ dispatch, opening })),
|
|
),
|
|
)(
|
|
"preserves whole-utterance $dispatch across capture with $opening opening audio",
|
|
async ({ dispatch, opening }) => {
|
|
const suffix = dispatch === "control" ? "stop that" : "send the update";
|
|
const oversized = opening === "oversized";
|
|
const prefix = oversized ? "Do not" : dispatch === "control" ? "please" : "Opening context";
|
|
const openingFrames = oversized ? 300 : 50;
|
|
cli.mockReset().mockImplementation(async (_command: string, [filePath]: string[]) => {
|
|
const wav = await fs.readFile(filePath!);
|
|
return { stdout: wav[44] === 1 ? "" : suffix, stderr: "" };
|
|
});
|
|
const provider = vi.fn(async ({ buffer, model }: { buffer: Buffer; model?: string }) => {
|
|
if (model === "unavailable-stt") {
|
|
throw new Error("synthetic model unavailable");
|
|
}
|
|
return { text: buffer[44] === 1 ? prefix : suffix };
|
|
});
|
|
const registry = createEmptyPluginRegistry();
|
|
registry.mediaUnderstandingProviders.push({
|
|
pluginId: "synthetic-audio",
|
|
source: "test/transcripts-discord-audio.integration.test.ts",
|
|
provider: { id: "synthetic-audio", capabilities: ["audio"], transcribeAudio: provider },
|
|
});
|
|
// Keep the synthetic runtime and discovery inventory in one generation;
|
|
// transcription must not materialize unrelated bundled plugin owners.
|
|
const metadataSnapshot = createPluginMetadataSnapshotFixture({
|
|
plugins: [
|
|
{
|
|
id: "synthetic-audio",
|
|
contracts: { mediaUnderstandingProviders: ["synthetic-audio"] },
|
|
},
|
|
],
|
|
});
|
|
const apiModel: MediaUnderstandingModelConfig = {
|
|
provider: "synthetic-audio",
|
|
model: "synthetic-stt",
|
|
capabilities: ["audio"],
|
|
};
|
|
const models: MediaUnderstandingModelConfig[] =
|
|
opening === "empty CLI"
|
|
? [
|
|
{
|
|
type: "cli",
|
|
command: "synthetic-stt",
|
|
// The process succeeds with empty stdout; its input is still valid.
|
|
args: ["{{AttachmentPath}}"],
|
|
capabilities: ["audio"],
|
|
},
|
|
]
|
|
: [
|
|
...(opening === "fallback" ? [{ ...apiModel, model: "unavailable-stt" }] : []),
|
|
apiModel,
|
|
];
|
|
const cfg = {
|
|
tools: {
|
|
media: {
|
|
audio: { maxBytes: 1_048_576 },
|
|
models,
|
|
},
|
|
},
|
|
models: {
|
|
providers: {
|
|
"synthetic-audio": {
|
|
baseUrl: "https://unused.invalid",
|
|
apiKey: "synthetic-fixture-key",
|
|
models: [],
|
|
},
|
|
},
|
|
},
|
|
};
|
|
const manager = createManager(
|
|
makeVoiceConfig({}, { groupPolicy: "open", allowFrom: ["discord:u-owner"] }),
|
|
undefined,
|
|
cfg,
|
|
);
|
|
await manager.join({ guildId: "g1", channelId: "1001" });
|
|
const entry = getSessionEntry(manager);
|
|
const conversations = vi.spyOn(entry.conversations, "enqueue");
|
|
if (dispatch === "control") {
|
|
controlRealtimeVoiceAgentRunMock.mockImplementation(async () => {
|
|
return {
|
|
ok: true,
|
|
mode: "cancel",
|
|
sessionKey: entry.route.sessionKey,
|
|
active: true,
|
|
aborted: true,
|
|
message: "Cancelled",
|
|
speak: false,
|
|
show: true,
|
|
suppress: true,
|
|
};
|
|
});
|
|
}
|
|
const stream = new PassThrough({ objectMode: true });
|
|
getSessionConnection(entry).receiver.subscribe.mockReturnValueOnce(stream);
|
|
const openingDecoded = createDeferred<void>();
|
|
decodeOpusStreamChunksMock.mockImplementation(async (input, options) => {
|
|
let frames = 0;
|
|
for await (const packet of input) {
|
|
await options.onChunk(packet, packet);
|
|
if (++frames === openingFrames) {
|
|
openingDecoded.resolve();
|
|
}
|
|
}
|
|
});
|
|
const wavSizes: number[] = [];
|
|
const wavPaths: string[] = [];
|
|
const results: Awaited<ReturnType<typeof transcribeAudioFile>>[] = [];
|
|
const suffixTranscribed = createDeferred<void>();
|
|
transcribeAudioFileMock.mockImplementation(async (params) => {
|
|
wavPaths.push(params.filePath);
|
|
const wav = await fs.readFile(params.filePath);
|
|
wavSizes.push(wav.length);
|
|
const isOpening = wav[44] === 1;
|
|
// Uncaptured conversation and captured audio may finish out of order.
|
|
if (opening === "empty CLI" && isOpening) {
|
|
await suffixTranscribed.promise;
|
|
}
|
|
const result = await withPluginRuntimeGenerationScope(
|
|
{ metadataSnapshot, pluginRegistry: registry },
|
|
() => transcribeAudioFile(params),
|
|
);
|
|
results[oversized || isOpening ? 0 : 1] = result;
|
|
if (!isOpening) {
|
|
suffixTranscribed.resolve();
|
|
}
|
|
return result;
|
|
});
|
|
const recorded = createDeferred<void>();
|
|
const sink = vi.fn(() => recorded.resolve());
|
|
const receiving = handleSpeakingStart(manager, entry, "u-owner");
|
|
const receivingOutcome = oversized
|
|
? expect(receiving).rejects.toThrow("speak a shorter segment")
|
|
: receiving;
|
|
const recordingStream = oversized ? new PassThrough({ objectMode: true }) : stream;
|
|
try {
|
|
// Main rejects oversized uncaptured input before writing a WAV. A later
|
|
// capture scan may record its suffix, but cannot admit it as a new command.
|
|
// Successful-empty and fallback siblings use a one-second opening chunk.
|
|
for (let frame = 0; frame < openingFrames; frame++) {
|
|
if (stream.destroyed) {
|
|
break;
|
|
}
|
|
stream.write(Buffer.alloc(3_840, 1));
|
|
await setImmediate();
|
|
}
|
|
if (oversized) {
|
|
await receivingOutcome;
|
|
expect(stream.destroyed).toBe(true);
|
|
expect(transcribeAudioFileMock).not.toHaveBeenCalled();
|
|
getSessionConnection(entry).receiver.subscribe.mockReturnValueOnce(recordingStream);
|
|
entry.audio.speakingUsers.add("u-owner");
|
|
} else {
|
|
await openingDecoded.promise;
|
|
}
|
|
expect(await startTranscripts(manager, sink)).toMatchObject({ ok: true });
|
|
for (let frame = 0; frame < 50; frame++) {
|
|
recordingStream.write(Buffer.alloc(3_840, 2));
|
|
}
|
|
recordingStream.end();
|
|
await receivingOutcome;
|
|
if (oversized) {
|
|
// The capture scan owns a new receive operation after the rejected opening.
|
|
await recorded.promise;
|
|
}
|
|
await entry.processingQueue;
|
|
// Join owned work even when transcription fails and produces no sink/dispatch call.
|
|
await Promise.all(conversations.mock.results.map((result) => result.value));
|
|
expect(loggerWarnMock).not.toHaveBeenCalledWith(expect.stringContaining("cli-missing-"));
|
|
expect(wavSizes).toEqual(oversized ? [192_044] : [192_044, 192_044]);
|
|
expect(results[0]).toMatchObject({
|
|
text: oversized ? suffix : opening === "fallback" ? prefix : undefined,
|
|
decision: {
|
|
outcome: opening === "empty CLI" ? "skipped" : "success",
|
|
attachmentDispositions: {
|
|
0: { kind: opening === "empty CLI" ? "failed" : "handled" },
|
|
},
|
|
attachmentProcessing: { 0: "completed" },
|
|
},
|
|
});
|
|
expect(provider).toHaveBeenCalledTimes(opening === "empty CLI" ? 0 : oversized ? 1 : 4);
|
|
expect(cli).toHaveBeenCalledTimes(opening === "empty CLI" ? 2 : 0);
|
|
if (oversized) {
|
|
expect(provider.mock.calls[0]![0].buffer).toHaveLength(192_044);
|
|
expect(provider.mock.calls[0]![0].buffer[44]).toBe(2);
|
|
}
|
|
expect(sink).toHaveBeenCalledExactlyOnceWith(expect.objectContaining({ text: suffix }));
|
|
for (const wavPath of wavPaths) {
|
|
await expect(fs.stat(wavPath)).rejects.toMatchObject({ code: "ENOENT" });
|
|
}
|
|
const completeText = opening === "fallback" ? `${prefix}\n${suffix}` : suffix;
|
|
if (oversized) {
|
|
expect(agentCommandMock).not.toHaveBeenCalled();
|
|
expect(controlRealtimeVoiceAgentRunMock).not.toHaveBeenCalled();
|
|
} else if (dispatch === "control") {
|
|
expect(controlRealtimeVoiceAgentRunMock).toHaveBeenCalledExactlyOnceWith({
|
|
getToolAuthorityOverlay: expect.any(Function),
|
|
sessionKey: entry.route.sessionKey,
|
|
text: completeText,
|
|
});
|
|
expect(agentCommandMock).not.toHaveBeenCalled();
|
|
} else {
|
|
expect(agentCommandMock).toHaveBeenCalledOnce();
|
|
expect(agentCommandMock.mock.calls[0]?.[0]).toMatchObject({
|
|
message: expect.stringContaining(completeText),
|
|
});
|
|
expect(controlRealtimeVoiceAgentRunMock).not.toHaveBeenCalled();
|
|
}
|
|
} finally {
|
|
stream.end();
|
|
recordingStream.end();
|
|
await receivingOutcome;
|
|
await entry.processingQueue;
|
|
conversations.mockRestore();
|
|
await manager.destroy();
|
|
}
|
|
},
|
|
);
|
|
|
|
it.each([
|
|
{ command: undefined, args: ["{{AttachmentPath}}"], reason: "cli-missing-command" },
|
|
{ command: "synthetic-stt", args: undefined, reason: "cli-missing-attachment-arg" },
|
|
])("finishes voice capture with a warning for $reason", async ({ command, args, reason }) => {
|
|
cli.mockReset();
|
|
const manager = createManager(
|
|
makeVoiceConfig({}, { groupPolicy: "open", allowFrom: ["discord:u-owner"] }),
|
|
undefined,
|
|
{ tools: { media: { models: [{ type: "cli", command, args, capabilities: ["audio"] }] } } },
|
|
);
|
|
await manager.join({ guildId: "g1", channelId: "1001" });
|
|
const sink = vi.fn();
|
|
expect(await startTranscripts(manager, sink)).toMatchObject({ ok: true });
|
|
const wavPaths: string[] = [];
|
|
transcribeAudioFileMock.mockImplementation(async (params) => {
|
|
wavPaths.push(params.filePath);
|
|
return await withPluginRuntimeGenerationScope(
|
|
{
|
|
metadataSnapshot: createPluginMetadataSnapshotFixture({ plugins: [] }),
|
|
pluginRegistry: createEmptyPluginRegistry(),
|
|
},
|
|
() => transcribeAudioFile(params),
|
|
);
|
|
});
|
|
// This joins receive, recording, and conversation completion, including refusal.
|
|
await receiveRecordedSpeech(manager);
|
|
expect(transcribeAudioFileMock).toHaveBeenCalledOnce();
|
|
expect(cli).not.toHaveBeenCalled();
|
|
expect(loggerWarnMock).toHaveBeenCalledWith(
|
|
expect.stringContaining(`discord voice: recording failed: ${reason}; Set `),
|
|
);
|
|
expect(sink).not.toHaveBeenCalled();
|
|
expect(agentCommandMock).not.toHaveBeenCalled();
|
|
expect(controlRealtimeVoiceAgentRunMock).not.toHaveBeenCalled();
|
|
for (const wavPath of wavPaths) {
|
|
await expect(fs.stat(wavPath)).rejects.toMatchObject({ code: "ENOENT" });
|
|
}
|
|
});
|
|
},
|
|
);
|