Skip to content

Commit 7c60a78

Browse files
committed
base future on throwspromise
1 parent 2fa5d59 commit 7c60a78

4 files changed

Lines changed: 9 additions & 10 deletions

File tree

agents/src/utils.ts

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -122,17 +122,17 @@ export class Queue<T> {
122122
}
123123

124124
/** @internal */
125-
export class Future<T = void> {
126-
#await: Promise<T>;
125+
export class Future<T = void, E extends Error = Error> {
126+
#await: ThrowsPromise<T, E>;
127127
#resolvePromise!: (value: T) => void;
128-
#rejectPromise!: (error: Error) => void;
128+
#rejectPromise!: (error: E) => void;
129129
#done: boolean = false;
130130
#rejected: boolean = false;
131131
#result: T | undefined = undefined;
132132
#error: Error | undefined = undefined;
133133

134134
constructor() {
135-
this.#await = new ThrowsPromise<T, Error>((resolve, reject) => {
135+
this.#await = new ThrowsPromise<T, E>((resolve, reject) => {
136136
this.#resolvePromise = resolve;
137137
this.#rejectPromise = reject;
138138
});
@@ -169,7 +169,7 @@ export class Future<T = void> {
169169
this.#resolvePromise(value);
170170
}
171171

172-
reject(error: Error) {
172+
reject(error: E) {
173173
this.#done = true;
174174
this.#rejected = true;
175175
this.#error = error;

agents/src/voice/agent_activity.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -169,7 +169,7 @@ export class AgentActivity implements RecognitionHooks {
169169
private _drainBlockedTasks: Task<any>[] = [];
170170
private _currentSpeech?: SpeechHandle;
171171
private speechQueue: Heap<[number, number, SpeechHandle]>; // [priority, timestamp, speechHandle]
172-
private q_updated: Future;
172+
private q_updated: Future<void, never>;
173173
private speechTasks: Set<Task<void>> = new Set();
174174
private lock = new Mutex();
175175
private audioStream = new MultiInputStream<AudioFrame>();
@@ -1316,7 +1316,7 @@ export class AgentActivity implements RecognitionHooks {
13161316
}
13171317

13181318
private async mainTask(signal: AbortSignal): Promise<void> {
1319-
const abortFuture = new Future();
1319+
const abortFuture = new Future<void, never>();
13201320
const abortHandler = () => {
13211321
abortFuture.resolve();
13221322
signal.removeEventListener('abort', abortHandler);

agents/src/voice/generation.ts

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -83,7 +83,7 @@ export interface _TTSGenerationData {
8383
/**
8484
* Future that resolves to a stream of timed transcripts, or null if TTS doesn't support it.
8585
*/
86-
timedTextsFut: Future<ReadableStream<TimedString> | null>;
86+
timedTextsFut: Future<ReadableStream<TimedString> | null, never>;
8787
/** Time to first byte (set when first audio frame is received) */
8888
ttfb?: number;
8989
}
@@ -566,7 +566,7 @@ export function performTTSInference(
566566
const outputWriter = audioStream.writable.getWriter();
567567
const audioOutputStream = audioStream.readable;
568568

569-
const timedTextsFut = new Future<ReadableStream<TimedString> | null>();
569+
const timedTextsFut = new Future<ReadableStream<TimedString> | null, never>();
570570
const timedTextsStream = new IdentityTransform<TimedString>();
571571
const timedTextsWriter = timedTextsStream.writable.getWriter();
572572

agents/src/voice/room_io/_output.ts

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,6 @@ import {
1515
TrackPublishOptions,
1616
TrackSource,
1717
} from '@livekit/rtc-node';
18-
import { ThrowsPromise } from '@livekit/throws-transformer/throws';
1918
import {
2019
ATTRIBUTE_TRANSCRIPTION_FINAL,
2120
ATTRIBUTE_TRANSCRIPTION_SEGMENT_ID,

0 commit comments

Comments
 (0)