Skip to content
Open
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
4 changes: 3 additions & 1 deletion packages/firebase_ai/firebase_ai/lib/firebase_ai.dart
Original file line number Diff line number Diff line change
Expand Up @@ -75,14 +75,16 @@ export 'src/live_api.dart'
LiveServerContent,
LiveServerToolCall,
LiveServerToolCallCancellation,
LiveServerVoiceActivity,
LiveServerResponse,
RealtimeInputConfig,
Sensitivity,
SessionResumptionConfig,
SessionResumptionUpdate,
SlidingWindow,
Transcription,
TurnCoverage;
TurnCoverage,
VoiceActivityType;
export 'src/live_session.dart' show LiveSession;
export 'src/mime_types.dart' show FirebaseAIMimeTypes;
export 'src/schema.dart' show JSONSchema, Schema, SchemaType;
Expand Down
75 changes: 72 additions & 3 deletions packages/firebase_ai/firebase_ai/lib/src/live_api.dart
Original file line number Diff line number Diff line change
Expand Up @@ -379,6 +379,41 @@ class GoingAwayNotice implements LiveServerMessage {
final String? timeLeft;
}

/// Whether the server detected the start or end of user speech.
enum VoiceActivityType {
/// The user started speaking.
activityStart('ACTIVITY_START'),

/// The user stopped speaking.
activityEnd('ACTIVITY_END');

const VoiceActivityType(this.value);

/// The JSON wire string value.
final String value;
}

/// A server message indicating that voice activity was detected.
///
/// Gemini 3.x Live models send this when the user starts or stops speaking.
class LiveServerVoiceActivity implements LiveServerMessage {
/// Creates a [LiveServerVoiceActivity] instance.
///
/// [type] (optional): Whether speech started or stopped.
/// [audioOffset] (optional): When the activity was detected, as a duration
/// string such as `2.520s`.
const LiveServerVoiceActivity({this.type, this.audioOffset});

/// Whether speech started or stopped.
///
/// Null when the server sends a type this SDK does not recognize.
final VoiceActivityType? type;

/// The time the activity was detected, relative to the start of the audio
/// stream. A duration string such as `2.520s`.
final String? audioOffset;
}

/// An update of the session resumption state.
///
/// This message is only sent if [SessionResumptionConfig] was set in the
Expand Down Expand Up @@ -558,6 +593,7 @@ class LiveClientToolResponse {
/// - `toolCall` messages indicating function calls requested by the model.
/// - `toolCallCancellation` messages to cancel pending function calls.
/// - `setupComplete` messages signaling the completion of the server setup.
/// - `voiceActivity` messages signaling the start or end of user speech.
///
/// If the JSON object does not match any of the expected formats, an
/// [FirebaseAISdkException] is thrown.
Expand Down Expand Up @@ -590,11 +626,23 @@ class LiveClientToolResponse {
/// Returns:
/// - A [LiveServerResponse] object representing the parsed message.
LiveServerResponse parseServerResponse(Object jsonObject) {
LiveServerMessage message = _parseServerMessage(jsonObject);
return tryParseServerResponse(jsonObject) ??
(throw unhandledFormat('LiveServerMessage', jsonObject));
}

/// Parses a live server message.
///
/// Returns null when [jsonObject] has a top-level key this SDK does not
/// recognize. Error payloads still throw [FirebaseAIException].
LiveServerResponse? tryParseServerResponse(Object jsonObject) {
final LiveServerMessage? message = _parseServerMessage(jsonObject);
if (message == null) {
return null;
}
return LiveServerResponse(message: message);
}

LiveServerMessage _parseServerMessage(Object jsonObject) {
LiveServerMessage? _parseServerMessage(Object jsonObject) {
if (jsonObject case {'error': final Object error}) {
throw parseError(error);
}
Expand Down Expand Up @@ -668,7 +716,28 @@ LiveServerMessage _parseServerMessage(Object jsonObject) {
lastConsumedClientMessageIndex:
sessionResumptionUpdateJson['lastConsumedClientMessageIndex'] as int?,
);
} else if (json.containsKey('voiceActivity')) {
return _parseVoiceActivity(json['voiceActivity']);
} else {
throw unhandledFormat('LiveServerMessage', json);
return null;
}
}

LiveServerVoiceActivity _parseVoiceActivity(Object? value) {
if (value is! Map) {
return const LiveServerVoiceActivity();
}
final json = Map<String, dynamic>.from(value);
final type = json['type'];
final audioOffset = json['audioOffset'];
return LiveServerVoiceActivity(
type: _parseVoiceActivityType(type),
audioOffset: audioOffset is String ? audioOffset : null,
);
}

VoiceActivityType? _parseVoiceActivityType(Object? value) => switch (value) {
'ACTIVITY_START' => VoiceActivityType.activityStart,
'ACTIVITY_END' => VoiceActivityType.activityEnd,
_ => null,
};
11 changes: 8 additions & 3 deletions packages/firebase_ai/firebase_ai/lib/src/live_session.dart
Original file line number Diff line number Diff line change
Expand Up @@ -159,9 +159,14 @@ class LiveSession {
final String jsonString =
message is String ? message : utf8.decode(message as List<int>);
var response = json.decode(jsonString);

if (!_messageController.isClosed) {
_messageController.add(parseServerResponse(response));
final parsed = tryParseServerResponse(response);
if (parsed == null) {
log(
'live_session: Ignoring unrecognized LiveServerMessage',
error: response,
);
} else if (!_messageController.isClosed) {
_messageController.add(parsed);
}
} catch (e) {
if (!_messageController.isClosed && _messageController.hasListener) {
Expand Down
86 changes: 86 additions & 0 deletions packages/firebase_ai/firebase_ai/test/live_session_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -122,6 +122,92 @@ void main() {
fakeWs.close();
});

test('receive stays open across a voiceActivity frame', () async {
final fakeWs = FakeWebSocketChannel();
final session = LiveSession.forTesting(fakeWs);

final messages = <LiveServerMessage>[];
final errors = <Object>[];
var done = false;
final completer = Completer<void>();
final subscription = session.receive().listen(
(response) {
messages.add(response.message);
if (messages.length == 2 && !completer.isCompleted) {
completer.complete();
}
},
onError: errors.add,
onDone: () => done = true,
);

fakeWs.emit(jsonEncode({
'voiceActivity': {
'type': 'ACTIVITY_START',
'audioOffset': '2.520s',
},
}));
fakeWs.emit('{"setupComplete": {}}');

await completer.future.timeout(const Duration(seconds: 5));
expect(messages[0], isA<LiveServerVoiceActivity>());
expect(
(messages[0] as LiveServerVoiceActivity).type,
VoiceActivityType.activityStart,
);
expect(
(messages[0] as LiveServerVoiceActivity).audioOffset,
'2.520s',
);
expect(messages[1], isA<LiveServerSetupComplete>());
expect(errors, isEmpty);
expect(done, isFalse);

await subscription.cancel();
fakeWs.close();
});

test('receive stays open after an unrecognized frame', () async {
final fakeWs = FakeWebSocketChannel();
final session = LiveSession.forTesting(fakeWs);

final errors = <Object>[];
var done = false;
final completer = Completer<LiveServerMessage>();
final subscription = session.receive().listen(
(response) {
if (response.message is LiveServerSetupComplete &&
!completer.isCompleted) {
completer.complete(response.message);
}
},
onError: (Object error) {
errors.add(error);
if (!completer.isCompleted) {
completer.completeError(error);
}
},
onDone: () {
done = true;
if (!completer.isCompleted) {
completer.completeError(StateError('receive closed'));
}
},
);

fakeWs.emit('{"unknown": {}}');
fakeWs.emit('{"setupComplete": {}}');

final message =
await completer.future.timeout(const Duration(seconds: 5));
expect(message, isA<LiveServerSetupComplete>());
expect(errors, isEmpty);
expect(done, isFalse);

await subscription.cancel();
fakeWs.close();
});

test('sendStartActivityRealtime sends correct activity_start message',
() async {
final fakeWs = FakeWebSocketChannel();
Expand Down
36 changes: 36 additions & 0 deletions packages/firebase_ai/firebase_ai/test/live_test.dart
Original file line number Diff line number Diff line change
Expand Up @@ -314,6 +314,42 @@ void main() {
expect(goAwayMessage.timeLeft, '50s');
});

test('parseServerMessage parses voiceActivity message correctly', () {
final start = parseServerResponse({
'voiceActivity': {
'type': 'ACTIVITY_START',
'audioOffset': '2.520s',
},
});
expect(start.message, isA<LiveServerVoiceActivity>());
final startMessage = start.message as LiveServerVoiceActivity;
expect(startMessage.type, VoiceActivityType.activityStart);
expect(startMessage.audioOffset, '2.520s');

final end = parseServerResponse({
'voiceActivity': {
'type': 'ACTIVITY_END',
'audioOffset': '3.640s',
},
});
final endMessage = end.message as LiveServerVoiceActivity;
expect(endMessage.type, VoiceActivityType.activityEnd);
expect(endMessage.audioOffset, '3.640s');
});

test('parseServerMessage keeps voiceActivity with an unrecognized type',
() {
final response = parseServerResponse({
'voiceActivity': {
'type': 'ACTIVITY_UNKNOWN',
'audioOffset': '1s',
},
});
final message = response.message as LiveServerVoiceActivity;
expect(message.type, isNull);
expect(message.audioOffset, '1s');
});

test('parseServerMessage throws VertexAIException for error message', () {
final jsonObject = {'error': {}};
expect(() => parseServerResponse(jsonObject),
Expand Down
Loading