Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
154 changes: 93 additions & 61 deletions packages/client/src/Call.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ import {
EndCallResponse,
GetCallReportResponse,
GetCallResponse,
GetCallRingStateResponse,
GetCallSessionParticipantStatsDetailsResponse,
GetOrCreateCallRequest,
GetOrCreateCallResponse,
Expand Down Expand Up @@ -152,14 +153,15 @@ import {
StatsReporter,
Tracer,
} from './stats';
import type { ClientEventReporter, JoinReason } from './reporting';
import type { ClientEventReporter, JoinReason, JoinSource } from './reporting';
import { AudioBindingsWatchdog } from './helpers/AudioBindingsWatchdog';
import { BlockedAudioTracker } from './helpers/BlockedAudioTracker';
import { TrackSubscriptionManager } from './helpers/TrackSubscriptionManager';
import { DynascaleManager } from './helpers/DynascaleManager';
import { createFirstVideoFrameDetector } from './helpers/firstVideoFrame';
import { ViewportTracker } from './helpers/ViewportTracker';
import { PermissionsContext } from './permissions';
import { RingStatePoller, RingTimeout, resolveOwnRingOutcome } from './ringing';
import { CallTypes } from './CallType';
import { StreamClient } from './coordinator/connection/client';
import { retryInterval, sleep } from './coordinator/connection/utils';
Expand Down Expand Up @@ -302,7 +304,8 @@ export class Call {
private statsReporter?: StatsReporter;
private sfuStatsReporter?: SfuStatsReporter;
private lastStatsOptions?: StatsOptions;
private dropTimeout: ReturnType<typeof setTimeout> | undefined;
private ringTimeout: RingTimeout | undefined;
private ringStatePoller: RingStatePoller | undefined;

private readonly clientStore: StreamVideoWriteableStateStore;
public readonly streamClient: StreamClient;
Expand Down Expand Up @@ -537,33 +540,20 @@ export class Call {
createSubscription(this.state.session$, (session) => {
if (!this.ringing) return;

const receiverId = this.clientStore.connectedUser?.id;
if (!receiverId) return;

const isAcceptedByMe = Boolean(session?.accepted_by[receiverId]);
const isRejectedByMe = Boolean(session?.rejected_by[receiverId]);

if (isAcceptedByMe || isRejectedByMe) {
this.cancelAutoDrop();
}

const isAcceptedElsewhere =
isAcceptedByMe && this.state.callingState === CallingState.RINGING;
const { settledByMe, leaveReason } = resolveOwnRingOutcome({
session,
currentUserId: this.currentUserId,
callingState: this.state.callingState,
});
if (settledByMe) this.cancelAutoDrop();
if (!leaveReason || hasPending(this.joinLeaveConcurrencyTag)) return;

if (
(isAcceptedElsewhere || isRejectedByMe) &&
!hasPending(this.joinLeaveConcurrencyTag)
) {
globalThis.streamRNVideoSDK?.callingX?.endCall(
this,
isAcceptedElsewhere ? 'answeredElsewhere' : 'rejected',
globalThis.streamRNVideoSDK?.callingX?.endCall(this, leaveReason);
this.leave().catch(() => {
this.logger.error(
'Could not leave a call that was accepted or rejected elsewhere',
);
this.leave().catch(() => {
this.logger.error(
'Could not leave a call that was accepted or rejected elsewhere',
);
});
}
});
}),
);
};
Expand Down Expand Up @@ -608,6 +598,7 @@ export class Call {
this.state.setCallingState(CallingState.RINGING);
}
this.scheduleAutoDrop();
this.scheduleRingStatePolling();
this.leaveCallHooks.add(registerRingingCallEventHandlers(this));
}
};
Expand Down Expand Up @@ -713,6 +704,13 @@ export class Call {
throw new Error('Cannot leave call that has already been left.');
}

// before the first await: the calling state stays RINGING well into the
// teardown, so pause both watchdogs before they can race this leave. They
// keep their deadlines, so a failed leave resumes them without handing the
// ring another window.
this.ringTimeout?.pause();
this.ringStatePoller?.pause();

await withoutConcurrency(this.joinLeaveConcurrencyTag, async () => {
const callingState = this.state.callingState;

Expand Down Expand Up @@ -821,6 +819,7 @@ export class Call {
this.unifiedSessionId = undefined;
this.ringingSubject.next(false);
this.cancelAutoDrop();
this.cancelRingStatePolling();
this.clientStore.unregisterCall(this);

globalThis.streamRNVideoSDK?.callManager.stop({
Expand Down Expand Up @@ -863,6 +862,17 @@ export class Call {
this.logger.warn('Failed to dispose media engine', err);
});
}
}).catch((err) => {
if (
!hasPending(this.joinLeaveConcurrencyTag) &&
this.state.callingState === CallingState.RINGING
) {
// resume, never re-arm: a fresh watchdog would restart the ring
// window this leave was already most of the way through
this.ringTimeout?.start();
this.ringStatePoller?.resume();
}
throw err;
});
};

Expand Down Expand Up @@ -1056,6 +1066,27 @@ export class Call {
);
};

/**
* Returns who accepted, rejected or missed the ring for a call session.
* Safe to poll: it performs no writes and emits no events.
*
* @param callSessionId the call session to read. Defaults to the current one.
* Pass it explicitly to read a session that has already ended, as ending a
* call clears its current session.
*/
getRingState = async (
callSessionId?: string,
): Promise<GetCallRingStateResponse> => {
const sessionId = callSessionId ?? this.state.session?.id;
if (!sessionId) {
throw new Error('Cannot read the ring state: the call has no session');
}
return this.streamClient.get<GetCallRingStateResponse>(
`${this.streamClientBasePath}/ring_state`,
{ call_session_id: sessionId },
);
};

/**
* A shortcut for {@link Call.get} with `notify` parameter set to `true`.
* Will send a `call.notification` event to the call members.
Expand Down Expand Up @@ -1112,12 +1143,14 @@ export class Call {
joinResponseTimeout,
rpcRequestTimeout,
allowOwnTracksLoopback = false,
joinSource,
...data
}: JoinCallData & {
maxJoinRetries?: number;
joinResponseTimeout?: number;
rpcRequestTimeout?: number;
allowOwnTracksLoopback?: boolean;
joinSource?: JoinSource;
} = {}): Promise<void> => {
const callingState = this.state.callingState;

Expand Down Expand Up @@ -1159,7 +1192,7 @@ export class Call {
try {
await this.clientEventReporter.withJoinLifecycle(
this.cid,
'first-attempt',
{ joinReason: 'first-attempt', joinSource },
async () => {
for (let attempt = 0; attempt < maxJoinRetries; attempt++) {
try {
Expand Down Expand Up @@ -2081,8 +2114,10 @@ export class Call {
this.reconnectReason === ReconnectReason.NETWORK_BACK_ONLINE
? 'network-available'
: 'full-rejoin';
await this.clientEventReporter.withJoinLifecycle(this.cid, joinReason, () =>
this.doJoin(this.joinCallData),
await this.clientEventReporter.withJoinLifecycle(
this.cid,
{ joinReason },
() => this.doJoin(this.joinCallData),
);
await this.restorePublishedTracks();
this.restoreSubscribedTracks();
Expand Down Expand Up @@ -2117,7 +2152,7 @@ export class Call {
const currentSfu = currentSfuClient.edgeName;
await this.clientEventReporter.withJoinLifecycle(
this.cid,
'migration',
{ joinReason: 'migration' },
() =>
this.doJoin({
...this.joinCallData,
Expand Down Expand Up @@ -3114,42 +3149,39 @@ export class Call {
*/
private scheduleAutoDrop = () => {
this.cancelAutoDrop();

const settings = this.state.settings;
if (!settings) return;
// ignore if the call is not ringing
if (this.state.callingState !== CallingState.RINGING) return;

const timeoutInMs = this.isCreatedByMe
? settings.ring.auto_cancel_timeout_ms
: settings.ring.incoming_call_timeout_ms;

// 0 means no auto-drop
if (timeoutInMs <= 0) return;
this.dropTimeout = setTimeout(() => {
// the call might have stopped ringing by this point,
// e.g. it was already accepted and joined
if (this.state.callingState !== CallingState.RINGING) return;
this.leave({
reject: true,
reason: 'timeout',
message: `ringing timeout - ${
this.isCreatedByMe
? 'no one accepted'
: `user didn't interact with incoming call screen`
}`,
}).catch((err) => {
this.logger.error('Failed to drop call', err);
});
}, timeoutInMs);
this.ringTimeout = new RingTimeout(this);
this.ringTimeout.start();
};

/**
* Cancels a scheduled auto-drop timeout.
*/
private cancelAutoDrop = () => {
clearTimeout(this.dropTimeout);
this.dropTimeout = undefined;
this.ringTimeout?.stop();
this.ringTimeout = undefined;
};

/**
* Starts polling for the ring outcome. Applicable only to ringing calls the
* current user created.
*/
private scheduleRingStatePolling = () => {
this.cancelRingStatePolling();

if (!this.isCreatedByMe) return;
const options = this.streamClient.options.ringStatePolling;
if (options === false) return;

this.ringStatePoller = new RingStatePoller(this, options);
this.ringStatePoller.start();
};

/**
* Cancels the ring state polling.
*/
private cancelRingStatePolling = () => {
this.ringStatePoller?.stop();
this.ringStatePoller = undefined;
};

/**
Expand Down
8 changes: 6 additions & 2 deletions packages/client/src/StreamSfuClient.ts
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,8 @@ import {
} from './helpers/promise';
import { withoutConcurrency } from './helpers/concurrency';
import { getTimers } from './timers';
import { Tracer, TraceSlice } from './stats';
import { getSdkInfo } from './helpers/client-details';
import { getStreamClientId, Tracer, TraceSlice } from './stats';
import { SfuJoinError, SfuTimeoutError } from './errors';

export type StreamSfuClientConstructor = {
Expand Down Expand Up @@ -271,7 +272,10 @@ export class StreamSfuClient {
this.rpc = createSignalClient({
baseUrl: server.url,
interceptors: [
withHeaders({ Authorization: `Bearer ${token}` }),
withHeaders({
Authorization: `Bearer ${token}`,
'X-Stream-Client': getStreamClientId(getSdkInfo()),
}),
this.tracer && withRequestTracer(this.tracer.trace),
this.logger.getLogLevel() === 'trace' && withRequestLogger(this.logger),
withTimeout(rpcRequestTimeout, this.tracer?.trace),
Expand Down
Loading
Loading