refactor(electron): encoding recording on the fly (#11457)

fix AF-2460, AF-2463

When recording is started, we start polling the pending raw buffers that are waiting for encoding. The buffers are determined by the cursor of the original raw buffer file. When recording is stopped, we will flush the pending buffers and wrap the encoded chunks into WebM.

```mermaid
sequenceDiagram
    participant App as App/UI
    participant RecordingFeature as Recording Feature
    participant StateMachine as State Machine
    participant FileSystem as File System
    participant StreamEncoder as Stream Encoder
    participant OpusEncoder as Opus Encoder
    participant WebM as WebM Muxer

    Note over App,WebM: Recording Start Flow
    App->>RecordingFeature: startRecording()
    RecordingFeature->>StateMachine: dispatch(START_RECORDING)
    StateMachine-->>RecordingFeature: status: 'recording'
    RecordingFeature->>StreamEncoder: createStreamEncoder(id, {sampleRate, channels})

    Note over App,WebM: Streaming Flow
    loop Audio Data Streaming
        RecordingFeature->>FileSystem: Write raw audio chunks to .raw file
        StreamEncoder->>FileSystem: Poll raw audio data
        FileSystem-->>StreamEncoder: Raw audio chunks
        StreamEncoder->>OpusEncoder: Encode chunks
        OpusEncoder-->>StreamEncoder: Encoded Opus frames
    end

    Note over App,WebM: Recording Stop Flow
    App->>RecordingFeature: stopRecording()
    RecordingFeature->>StateMachine: dispatch(STOP_RECORDING)
    StateMachine-->>RecordingFeature: status: 'stopped'
    StreamEncoder->>OpusEncoder: flush()
    StreamEncoder->>WebM: muxToWebM(encodedChunks)
    WebM-->>RecordingFeature: WebM buffer
    RecordingFeature->>FileSystem: Save as .opus file
    RecordingFeature->>StateMachine: dispatch(SAVE_RECORDING)
```
This commit is contained in:
pengx17
2025-04-03 15:56:53 +00:00
parent 8ce10e6d0a
commit 133be72ac2
7 changed files with 260 additions and 76 deletions

View File

@@ -1,7 +1,11 @@
import { Button } from '@affine/component'; import { Button } from '@affine/component';
import { useAsyncCallback } from '@affine/core/components/hooks/affine-async-hooks'; import { useAsyncCallback } from '@affine/core/components/hooks/affine-async-hooks';
import { appIconMap } from '@affine/core/utils'; import { appIconMap } from '@affine/core/utils';
import { encodeRawBufferToOpus } from '@affine/core/utils/webm-encoding'; import {
createStreamEncoder,
encodeRawBufferToOpus,
type OpusStreamEncoder,
} from '@affine/core/utils/webm-encoding';
import { apis, events } from '@affine/electron-api'; import { apis, events } from '@affine/electron-api';
import { useI18n } from '@affine/i18n'; import { useI18n } from '@affine/i18n';
import track from '@affine/track'; import track from '@affine/track';
@@ -23,6 +27,8 @@ type Status = {
appGroupId?: number; appGroupId?: number;
icon?: Buffer; icon?: Buffer;
filepath?: string; filepath?: string;
sampleRate?: number;
numberOfChannels?: number;
}; };
export const useRecordingStatus = () => { export const useRecordingStatus = () => {
@@ -99,7 +105,8 @@ export function Recording() {
await apis?.recording?.stopRecording(status.id); await apis?.recording?.stopRecording(status.id);
}, [status]); }, [status]);
const handleProcessStoppedRecording = useAsyncCallback(async () => { const handleProcessStoppedRecording = useAsyncCallback(
async (currentStreamEncoder?: OpusStreamEncoder) => {
let id: number | undefined; let id: number | undefined;
try { try {
const result = await apis?.recording?.getCurrentRecording(); const result = await apis?.recording?.getCurrentRecording();
@@ -115,7 +122,9 @@ export function Recording() {
return; return;
} }
const [buffer] = await Promise.all([ const [buffer] = await Promise.all([
encodeRawBufferToOpus({ currentStreamEncoder
? currentStreamEncoder.finish()
: encodeRawBufferToOpus({
filepath, filepath,
sampleRate, sampleRate,
numberOfChannels, numberOfChannels,
@@ -134,21 +143,62 @@ export function Recording() {
await apis?.recording.removeRecording(id); await apis?.recording.removeRecording(id);
} }
} }
}, []); },
[]
);
useEffect(() => { useEffect(() => {
// allow processing stopped event in tray menu as well: let currentStreamEncoder: OpusStreamEncoder | undefined;
return events?.recording.onRecordingStatusChanged(status => {
apis?.recording
.getCurrentRecording()
.then(status => {
if (status) {
return handleRecordingStatusChanged(status);
}
return;
})
.catch(console.error);
const handleRecordingStatusChanged = async (status: Status) => {
if (status?.status === 'new') { if (status?.status === 'new') {
track.popup.$.recordingBar.toggleRecordingBar({ track.popup.$.recordingBar.toggleRecordingBar({
type: 'Meeting record', type: 'Meeting record',
appName: status.appName || 'System Audio', appName: status.appName || 'System Audio',
}); });
} }
if (
status?.status === 'recording' &&
status.sampleRate &&
status.numberOfChannels &&
(!currentStreamEncoder || currentStreamEncoder.id !== status.id)
) {
currentStreamEncoder?.close();
currentStreamEncoder = createStreamEncoder(status.id, {
sampleRate: status.sampleRate,
numberOfChannels: status.numberOfChannels,
});
currentStreamEncoder.poll().catch(console.error);
}
if (status?.status === 'stopped') { if (status?.status === 'stopped') {
handleProcessStoppedRecording(); handleProcessStoppedRecording(currentStreamEncoder);
currentStreamEncoder = undefined;
}
};
// allow processing stopped event in tray menu as well:
const unsubscribe = events?.recording.onRecordingStatusChanged(status => {
if (status) {
handleRecordingStatusChanged(status).catch(console.error);
} }
}); });
return () => {
unsubscribe?.();
currentStreamEncoder?.close();
};
}, [handleProcessStoppedRecording]); }, [handleProcessStoppedRecording]);
const handleStartRecording = useAsyncCallback(async () => { const handleStartRecording = useAsyncCallback(async () => {

View File

@@ -1,5 +1,6 @@
/* oxlint-disable no-var-requires */ /* oxlint-disable no-var-requires */
import { execSync } from 'node:child_process'; import { execSync } from 'node:child_process';
import fsp from 'node:fs/promises';
import path from 'node:path'; import path from 'node:path';
// Should not load @affine/native for unsupported platforms // Should not load @affine/native for unsupported platforms
@@ -240,7 +241,12 @@ function setupNewRunningAppGroup() {
); );
} }
function createRecording(status: RecordingStatus) { export function createRecording(status: RecordingStatus) {
let recording = recordings.get(status.id);
if (recording) {
return recording;
}
const bufferedFilePath = path.join( const bufferedFilePath = path.join(
SAVED_RECORDINGS_DIR, SAVED_RECORDINGS_DIR,
`${status.appGroup?.bundleIdentifier ?? 'unknown'}-${status.id}-${status.startTime}.raw` `${status.appGroup?.bundleIdentifier ?? 'unknown'}-${status.id}-${status.startTime}.raw`
@@ -275,7 +281,7 @@ function createRecording(status: RecordingStatus) {
? status.app.rawInstance.tapAudio(tapAudioSamples) ? status.app.rawInstance.tapAudio(tapAudioSamples)
: ShareableContent.tapGlobalAudio(null, tapAudioSamples); : ShareableContent.tapGlobalAudio(null, tapAudioSamples);
const recording: Recording = { recording = {
id: status.id, id: status.id,
startTime: status.startTime, startTime: status.startTime,
app: status.app, app: status.app,
@@ -284,6 +290,8 @@ function createRecording(status: RecordingStatus) {
stream, stream,
}; };
recordings.set(status.id, recording);
return recording; return recording;
} }
@@ -330,7 +338,6 @@ function setupRecordingListeners() {
// create a recording if not exists // create a recording if not exists
if (!recording) { if (!recording) {
recording = createRecording(status); recording = createRecording(status);
recordings.set(status.id, recording);
} }
} else if (status?.status === 'stopped') { } else if (status?.status === 'stopped') {
const recording = recordings.get(status.id); const recording = recordings.get(status.id);
@@ -518,6 +525,10 @@ export function startRecording(
appGroup: normalizeAppGroupInfo(appGroup), appGroup: normalizeAppGroupInfo(appGroup),
}); });
if (state?.status === 'recording') {
createRecording(state);
}
// set a timeout to stop the recording after MAX_DURATION_FOR_TRANSCRIPTION // set a timeout to stop the recording after MAX_DURATION_FOR_TRANSCRIPTION
setTimeout(() => { setTimeout(() => {
if ( if (
@@ -544,7 +555,7 @@ export function resumeRecording(id: number) {
export async function stopRecording(id: number) { export async function stopRecording(id: number) {
const recording = recordings.get(id); const recording = recordings.get(id);
if (!recording) { if (!recording) {
logger.error(`Recording ${id} not found`); logger.error(`stopRecording: Recording ${id} not found`);
return; return;
} }
@@ -590,9 +601,6 @@ export async function stopRecording(id: number) {
const recordingStatus = recordingStateMachine.dispatch({ const recordingStatus = recordingStateMachine.dispatch({
type: 'STOP_RECORDING', type: 'STOP_RECORDING',
id, id,
filepath: String(recording.file.path),
sampleRate: recording.stream.sampleRate,
numberOfChannels: recording.stream.channels,
}); });
if (!recordingStatus) { if (!recordingStatus) {
@@ -620,11 +628,35 @@ export async function stopRecording(id: number) {
} }
} }
export async function getRawAudioBuffers(
id: number,
cursor?: number
): Promise<{
buffer: Buffer;
nextCursor: number;
}> {
const recording = recordings.get(id);
if (!recording) {
throw new Error(`getRawAudioBuffers: Recording ${id} not found`);
}
const start = cursor ?? 0;
const file = await fsp.open(recording.file.path, 'r');
const stats = await file.stat();
const buffer = Buffer.alloc(stats.size - start);
const result = await file.read(buffer, 0, buffer.length, start);
await file.close();
return {
buffer,
nextCursor: start + result.bytesRead,
};
}
export async function readyRecording(id: number, buffer: Buffer) { export async function readyRecording(id: number, buffer: Buffer) {
const recordingStatus = recordingStatus$.value; const recordingStatus = recordingStatus$.value;
const recording = recordings.get(id); const recording = recordings.get(id);
if (!recordingStatus || recordingStatus.id !== id || !recording) { if (!recordingStatus || recordingStatus.id !== id || !recording) {
logger.error(`Recording ${id} not found`); logger.error(`readyRecording: Recording ${id} not found`);
return; return;
} }
@@ -635,6 +667,16 @@ export async function readyRecording(id: number, buffer: Buffer) {
await fs.writeFile(filepath, buffer); await fs.writeFile(filepath, buffer);
// can safely remove the raw file now
const rawFilePath = recording.file.path;
logger.info('remove raw file', rawFilePath);
if (rawFilePath) {
try {
await fs.unlink(rawFilePath);
} catch (err) {
logger.error('failed to remove raw file', err);
}
}
// Update the status through the state machine // Update the status through the state machine
recordingStateMachine.dispatch({ recordingStateMachine.dispatch({
type: 'SAVE_RECORDING', type: 'SAVE_RECORDING',
@@ -689,7 +731,8 @@ export interface SerializedRecordingStatus {
export function serializeRecordingStatus( export function serializeRecordingStatus(
status: RecordingStatus status: RecordingStatus
): SerializedRecordingStatus { ): SerializedRecordingStatus | null {
const recording = recordings.get(status.id);
return { return {
id: status.id, id: status.id,
status: status.status, status: status.status,
@@ -697,9 +740,10 @@ export function serializeRecordingStatus(
appGroupId: status.appGroup?.processGroupId, appGroupId: status.appGroup?.processGroupId,
icon: status.appGroup?.icon, icon: status.appGroup?.icon,
startTime: status.startTime, startTime: status.startTime,
filepath: status.filepath, filepath:
sampleRate: status.sampleRate, status.filepath ?? (recording ? String(recording.file.path) : undefined),
numberOfChannels: status.numberOfChannels, sampleRate: recording?.stream.sampleRate,
numberOfChannels: recording?.stream.channels,
}; };
} }

View File

@@ -12,6 +12,7 @@ import {
checkRecordingAvailable, checkRecordingAvailable,
checkScreenRecordingPermission, checkScreenRecordingPermission,
disableRecordingFeature, disableRecordingFeature,
getRawAudioBuffers,
getRecording, getRecording,
handleBlockCreationFailed, handleBlockCreationFailed,
handleBlockCreationSuccess, handleBlockCreationSuccess,
@@ -47,6 +48,9 @@ export const recordingHandlers = {
stopRecording: async (_, id: number) => { stopRecording: async (_, id: number) => {
return stopRecording(id); return stopRecording(id);
}, },
getRawAudioBuffers: async (_, id: number, cursor?: number) => {
return getRawAudioBuffers(id, cursor);
},
// save the encoded recording buffer to the file system // save the encoded recording buffer to the file system
readyRecording: async (_, id: number, buffer: Uint8Array) => { readyRecording: async (_, id: number, buffer: Uint8Array) => {
return readyRecording(id, Buffer.from(buffer)); return readyRecording(id, Buffer.from(buffer));

View File

@@ -9,15 +9,15 @@ import type { AppGroupInfo, RecordingStatus } from './types';
*/ */
export type RecordingEvent = export type RecordingEvent =
| { type: 'NEW_RECORDING'; appGroup?: AppGroupInfo } | { type: 'NEW_RECORDING'; appGroup?: AppGroupInfo }
| { type: 'START_RECORDING'; appGroup?: AppGroupInfo } | {
type: 'START_RECORDING';
appGroup?: AppGroupInfo;
}
| { type: 'PAUSE_RECORDING'; id: number } | { type: 'PAUSE_RECORDING'; id: number }
| { type: 'RESUME_RECORDING'; id: number } | { type: 'RESUME_RECORDING'; id: number }
| { | {
type: 'STOP_RECORDING'; type: 'STOP_RECORDING';
id: number; id: number;
filepath: string;
sampleRate: number;
numberOfChannels: number;
} }
| { | {
type: 'SAVE_RECORDING'; type: 'SAVE_RECORDING';
@@ -81,12 +81,7 @@ export class RecordingStateMachine {
newStatus = this.handleResumeRecording(); newStatus = this.handleResumeRecording();
break; break;
case 'STOP_RECORDING': case 'STOP_RECORDING':
newStatus = this.handleStopRecording( newStatus = this.handleStopRecording(event.id);
event.id,
event.filepath,
event.sampleRate,
event.numberOfChannels
);
break; break;
case 'SAVE_RECORDING': case 'SAVE_RECORDING':
newStatus = this.handleSaveRecording(event.id, event.filepath); newStatus = this.handleSaveRecording(event.id, event.filepath);
@@ -208,12 +203,7 @@ export class RecordingStateMachine {
/** /**
* Handle the STOP_RECORDING event * Handle the STOP_RECORDING event
*/ */
private handleStopRecording( private handleStopRecording(id: number): RecordingStatus | null {
id: number,
filepath: string,
sampleRate: number,
numberOfChannels: number
): RecordingStatus | null {
const currentStatus = this.recordingStatus$.value; const currentStatus = this.recordingStatus$.value;
if (!currentStatus || currentStatus.id !== id) { if (!currentStatus || currentStatus.id !== id) {
@@ -232,9 +222,6 @@ export class RecordingStateMachine {
return { return {
...currentStatus, ...currentStatus,
status: 'stopped', status: 'stopped',
filepath,
sampleRate,
numberOfChannels,
}; };
} }

View File

@@ -29,6 +29,7 @@ export interface Recording {
file: WriteStream; file: WriteStream;
stream: AudioTapStream; stream: AudioTapStream;
startTime: number; startTime: number;
filepath?: string; // the filepath of the recording (only available when status is ready)
} }
export interface RecordingStatus { export interface RecordingStatus {
@@ -52,7 +53,5 @@ export interface RecordingStatus {
app?: TappableAppInfo; app?: TappableAppInfo;
appGroup?: AppGroupInfo; appGroup?: AppGroupInfo;
startTime: number; // 0 means not started yet startTime: number; // 0 means not started yet
filepath?: string; // the filepath of the recording (only available when status is ready) filepath?: string; // encoded file path
sampleRate?: number;
numberOfChannels?: number;
} }

View File

@@ -1,7 +1,11 @@
import { join } from 'node:path'; import { join } from 'node:path';
import { setTimeout } from 'node:timers/promises'; import { setTimeout } from 'node:timers/promises';
import { BrowserWindow, type BrowserWindowConstructorOptions } from 'electron'; import {
app,
BrowserWindow,
type BrowserWindowConstructorOptions,
} from 'electron';
import { BehaviorSubject } from 'rxjs'; import { BehaviorSubject } from 'rxjs';
import { popupViewUrl } from '../constants'; import { popupViewUrl } from '../constants';
@@ -96,6 +100,9 @@ abstract class PopupWindow {
}, },
}); });
// it seems that the dock will disappear when popup windows are shown
await app.dock?.show();
// required to make the window transparent // required to make the window transparent
browserWindow.setBackgroundColor('#00000000'); browserWindow.setBackgroundColor('#00000000');
browserWindow.setVisibleOnAllWorkspaces(true, { browserWindow.setVisibleOnAllWorkspaces(true, {

View File

@@ -1,4 +1,5 @@
import { DebugLogger } from '@affine/debug'; import { DebugLogger } from '@affine/debug';
import { apis } from '@affine/electron-api';
import { ArrayBufferTarget, Muxer } from 'webm-muxer'; import { ArrayBufferTarget, Muxer } from 'webm-muxer';
interface AudioEncodingConfig { interface AudioEncodingConfig {
@@ -12,10 +13,10 @@ const logger = new DebugLogger('webm-encoding');
/** /**
* Creates and configures an Opus encoder with the given settings * Creates and configures an Opus encoder with the given settings
*/ */
async function createOpusEncoder(config: AudioEncodingConfig): Promise<{ export function createOpusEncoder(config: AudioEncodingConfig): {
encoder: AudioEncoder; encoder: AudioEncoder;
encodedChunks: EncodedAudioChunk[]; encodedChunks: EncodedAudioChunk[];
}> { } {
const encodedChunks: EncodedAudioChunk[] = []; const encodedChunks: EncodedAudioChunk[] = [];
const encoder = new AudioEncoder({ const encoder = new AudioEncoder({
output: chunk => { output: chunk => {
@@ -81,7 +82,7 @@ async function encodeAudioFrames({
/** /**
* Creates a WebM container with the encoded audio chunks * Creates a WebM container with the encoded audio chunks
*/ */
function muxToWebM( export function muxToWebM(
encodedChunks: EncodedAudioChunk[], encodedChunks: EncodedAudioChunk[],
config: AudioEncodingConfig config: AudioEncodingConfig
): Uint8Array { ): Uint8Array {
@@ -121,7 +122,7 @@ export async function encodeRawBufferToOpus({
throw new Error('Response body is null'); throw new Error('Response body is null');
} }
const { encoder, encodedChunks } = await createOpusEncoder({ const { encoder, encodedChunks } = createOpusEncoder({
sampleRate, sampleRate,
numberOfChannels, numberOfChannels,
}); });
@@ -193,7 +194,7 @@ export async function encodeAudioBlobToOpus(
bitrate: targetBitrate, bitrate: targetBitrate,
}; };
const { encoder, encodedChunks } = await createOpusEncoder(config); const { encoder, encodedChunks } = createOpusEncoder(config);
// Combine all channels into a single Float32Array // Combine all channels into a single Float32Array
const audioData = new Float32Array( const audioData = new Float32Array(
@@ -220,3 +221,95 @@ export async function encodeAudioBlobToOpus(
await audioContext.close(); await audioContext.close();
} }
} }
export const createStreamEncoder = (
recordingId: number,
codecs: {
sampleRate: number;
numberOfChannels: number;
targetBitrate?: number;
}
) => {
const { encoder, encodedChunks } = createOpusEncoder({
sampleRate: codecs.sampleRate,
numberOfChannels: codecs.numberOfChannels,
bitrate: codecs.targetBitrate,
});
const toAudioData = (buffer: Uint8Array) => {
// Each sample in f32 format is 4 bytes
const BYTES_PER_SAMPLE = 4;
return new AudioData({
format: 'f32',
sampleRate: codecs.sampleRate,
numberOfChannels: codecs.numberOfChannels,
numberOfFrames:
buffer.length / BYTES_PER_SAMPLE / codecs.numberOfChannels,
timestamp: 0,
data: buffer,
});
};
let cursor = 0;
let isClosed = false;
const next = async () => {
if (!apis || isClosed) {
throw new Error('Electron API is not available');
}
const { buffer, nextCursor } = await apis.recording.getRawAudioBuffers(
recordingId,
cursor
);
if (isClosed || cursor === nextCursor) {
return;
}
cursor = nextCursor;
logger.debug('Encoding next chunk', cursor, nextCursor);
encoder.encode(toAudioData(buffer));
};
const poll = async () => {
if (isClosed) {
return;
}
logger.debug('Polling next chunk');
await next();
await new Promise(resolve => setTimeout(resolve, 1000));
await poll();
};
const close = () => {
if (isClosed) {
return;
}
isClosed = true;
return encoder.close();
};
return {
id: recordingId,
next,
poll,
flush: () => {
return encoder.flush();
},
close,
finish: async () => {
logger.debug('Finishing encoding');
await next();
close();
const buffer = muxToWebM(encodedChunks, {
sampleRate: codecs.sampleRate,
numberOfChannels: codecs.numberOfChannels,
bitrate: codecs.targetBitrate,
});
return buffer;
},
[Symbol.dispose]: () => {
close();
},
};
};
export type OpusStreamEncoder = ReturnType<typeof createStreamEncoder>;