feat: fix voice join/leave, channel edit defaults, channel-update stream restart
Some checks failed
Build Linux Binaries / linux/amd64 (push) Has been cancelled
Build Linux Binaries / linux/arm64 (push) Has been cancelled

Three bugs fixed across the full stack (proto/server/core/ABI/Win/macOS/iOS):

1. Join/Leave Voice now truly subscribes/unsubscribes from the voice plane.
   Previously the button only toggled the local mic — receiving was always on
   (gated by channel membership alone). Added a protocol-level voice subscription
   concept: new SubscribeVoiceRequest/UnsubscribeVoiceRequest/VoiceSubscriptionResult
   proto messages, User.voice_subscribed field, vc_join_voice/vc_leave_voice C ABI
   functions, VC_EVENT_VOICE_STATE event, server-side voice_subscribed flag checked
   by the SFU relay recipient filter, and core-client gating of remote-stream
   decoder setup. All three clients rewired to subscribe+mic on Join / unsubscribe
   on Leave. Text chat works regardless of voice subscription.

2. Channel edit dialog now shows the channel's actual current settings. The read
   struct vc_channel was missing sort_order and audio fields — only the write
   struct vc_channel_info had them. Extended vc_channel with both (additive, no
   ABI break), updated the session model and list_channels marshaling to populate
   them, and updated all three clients' edit callers to use actual channel info
   instead of hardcoded defaults.

3. Channel parameter updates now automatically restart everyone's streams.
   Previously editing a channel's audio config persisted and broadcast a
   ChannelEvent::UPDATED, but no layer restarted streams — encoders/decoders are
   frozen at announce time. handle_channel_event now detects audio-config changes
   on the user's current channel and stop->starts each active local stream. The
   server reads the updated config on re-announce; peers wire up fresh decoders
   at the new ssrc.

All 29 CTest tests pass; Windows DLL + C# client build clean. Apple clients not
yet compile-verified (Windows environment).
This commit is contained in:
2026-06-24 14:29:39 +02:00
parent 2baefddbe4
commit 6fe7bf0158
40 changed files with 701 additions and 71 deletions

View File

@@ -10,6 +10,38 @@ up instantly. Newest status at the top.
## ▶ Where we left off / next action ## ▶ Where we left off / next action
- **Done (2026-06-24):** **Three bug fixes — voice join/leave, channel edit defaults, channel-update stream restart.**
1. **Join/Leave Voice now truly subscribes/unsubscribes from the voice plane.** Previously
"Join Voice" only started the local mic — receiving was always on (gated by channel
membership alone). Added a protocol-level voice subscription concept: new
`SubscribeVoiceRequest`/`UnsubscribeVoiceRequest`/`VoiceSubscriptionResult` proto messages
(`core/proto/voicecat.proto`), `User.voice_subscribed` field, `vc_join_voice`/`vc_leave_voice`
C ABI functions (`core/include/voicecat.h`), `VC_EVENT_VOICE_STATE` event, server-side
`voice_subscribed_` flag on `ConnSession` checked by the SFU relay's recipient filter
(`SessionRegistry::find_channel_sessions` excludes non-subscribers; `MediaRelay::on_udp_frame`
also skips non-subscribed senders). The core client gates `sync_remote_streams` on
`voice_subscribed_`, tears down all remote decoders + stops local streams on leave, and
re-syncs from the session model on join. All three clients (Windows/macOS/iOS) rewired
their Join/Leave Voice button to call `joinVoice`+start mic / `leaveVoice`+core stops mic.
The configured input mode (PTT/VAD/AlwaysOn) takes effect on join — no extra mic button.
Text chat works regardless of voice subscription. **Apple clients not yet compile-verified
(Windows environment).**
2. **Channel edit dialog now shows the channel's actual current settings.** The read struct
`vc_channel` (`voicecat.h`) was missing `sort_order` and `audio` fields — only the write
struct `vc_channel_info` had them. Extended `vc_channel` with both (additive, no ABI break),
updated the session model (`session::Channel`) and `apply_snapshot`/`apply_channel_event`
to populate them, and updated `vc_list_channels` marshaling. All three clients now build
the edit descriptor from the actual channel info instead of hardcoded defaults.
3. **Channel parameter updates now automatically restart everyone's streams.** Previously
editing a channel's audio config (codec/bitrate/sample-rate/FEC/DTX/etc.) persisted and
broadcast a `ChannelEvent::UPDATED`, but no layer restarted streams — encoders/decoders
are frozen at announce time. `handle_channel_event` (`core/src/core/client.cpp`) now
detects audio-config changes on the user's current channel and calls
`restart_active_streams_for_channel`, which stop→starts each active local stream. The
server reads the updated channel config on re-announce, and peers' `sync_remote_streams`
wire up fresh decoders at the new ssrc. The `LocalStream` struct now retains the stream
label across restarts. No server or protocol change needed.
- **[ ] Soon — jitter buffer should measure REAL arrival jitter (RFC 3550), not sender - **[ ] Soon — jitter buffer should measure REAL arrival jitter (RFC 3550), not sender
timestamps.** `JitterBuffer::push` (`core/src/audio/audio_engine.cpp:84-108`) estimates timestamps.** `JitterBuffer::push` (`core/src/audio/audio_engine.cpp:84-108`) estimates
jitter from `gap = ts - last_push_ts_`, where `ts` is the **sender's timestamp** — which is jitter from `gap = ts - last_push_ts_`, where `ts` is the **sender's timestamp** — which is

View File

@@ -133,6 +133,8 @@ public enum VoiceCatEventType: UInt32, Sendable, Equatable {
case genericResult = 14 case genericResult = 14
/// M5: reply to `requestAccountList()` call `listAccounts()` to read. /// M5: reply to `requestAccountList()` call `listAccounts()` to read.
case accountList = 15 case accountList = 15
/// Voice-plane subscription state. `u32a` = 1 (subscribed) or 0 (unsubscribed).
case voiceState = 16
public init(_ cValue: vc_event_type) { public init(_ cValue: vc_event_type) {
self = VoiceCatEventType(rawValue: cValue.rawValue) ?? .error self = VoiceCatEventType(rawValue: cValue.rawValue) ?? .error

View File

@@ -38,7 +38,8 @@ internal enum Marshaling {
let c = items.advanced(by: i).pointee let c = items.advanced(by: i).pointee
result.append(Channel(id: c.id, parentId: c.parent_id, name: string(c.name), result.append(Channel(id: c.id, parentId: c.parent_id, name: string(c.name),
topic: string(c.topic), passwordProtected: c.password_protected != 0, topic: string(c.topic), passwordProtected: c.password_protected != 0,
maxUsers: c.max_users)) maxUsers: c.max_users, sortOrder: c.sort_order,
audio: audioConfig(c.audio)))
} }
vc_free_channel_list(&list) vc_free_channel_list(&list)
return result return result
@@ -53,7 +54,8 @@ internal enum Marshaling {
result.append(User(id: u.id, nickname: string(u.nickname), isGuest: u.is_guest != 0, result.append(User(id: u.id, nickname: string(u.nickname), isGuest: u.is_guest != 0,
channelId: u.channel_id, selfMicMuted: u.self_mic_muted != 0, channelId: u.channel_id, selfMicMuted: u.self_mic_muted != 0,
selfDeafened: u.self_deafened != 0, serverMuted: u.server_muted != 0, selfDeafened: u.self_deafened != 0, serverMuted: u.server_muted != 0,
serverDeafened: u.server_deafened != 0)) serverDeafened: u.server_deafened != 0,
voiceSubscribed: u.voice_subscribed != 0))
} }
vc_free_user_list(&list) vc_free_user_list(&list)
return result return result

View File

@@ -16,11 +16,17 @@ public struct Channel: Sendable, Equatable, Identifiable {
public let passwordProtected: Bool public let passwordProtected: Bool
/// 0 = unlimited. /// 0 = unlimited.
public let maxUsers: UInt32 public let maxUsers: UInt32
public let sortOrder: UInt32
/// Authoritative channel Opus params (docs/voice.md §3). Populated from the Channel proto
/// so the edit dialog can read back the current config.
public let audio: AudioConfig
public init(id: UInt32, parentId: UInt32, name: String, topic: String, public init(id: UInt32, parentId: UInt32, name: String, topic: String,
passwordProtected: Bool, maxUsers: UInt32) { passwordProtected: Bool, maxUsers: UInt32, sortOrder: UInt32,
audio: AudioConfig) {
self.id = id; self.parentId = parentId; self.name = name; self.topic = topic self.id = id; self.parentId = parentId; self.name = name; self.topic = topic
self.passwordProtected = passwordProtected; self.maxUsers = maxUsers self.passwordProtected = passwordProtected; self.maxUsers = maxUsers
self.sortOrder = sortOrder; self.audio = audio
} }
} }
@@ -57,13 +63,15 @@ public struct User: Sendable, Equatable, Identifiable {
public let selfDeafened: Bool public let selfDeafened: Bool
public let serverMuted: Bool public let serverMuted: Bool
public let serverDeafened: Bool public let serverDeafened: Bool
public let voiceSubscribed: Bool
public init(id: UInt32, nickname: String, isGuest: Bool, channelId: UInt32, public init(id: UInt32, nickname: String, isGuest: Bool, channelId: UInt32,
selfMicMuted: Bool, selfDeafened: Bool, serverMuted: Bool, selfMicMuted: Bool, selfDeafened: Bool, serverMuted: Bool,
serverDeafened: Bool) { serverDeafened: Bool, voiceSubscribed: Bool) {
self.id = id; self.nickname = nickname; self.isGuest = isGuest; self.channelId = channelId self.id = id; self.nickname = nickname; self.isGuest = isGuest; self.channelId = channelId
self.selfMicMuted = selfMicMuted; self.selfDeafened = selfDeafened self.selfMicMuted = selfMicMuted; self.selfDeafened = selfDeafened
self.serverMuted = serverMuted; self.serverDeafened = serverDeafened self.serverMuted = serverMuted; self.serverDeafened = serverDeafened
self.voiceSubscribed = voiceSubscribed
} }
} }

View File

@@ -256,6 +256,16 @@ public final class VoiceCatClient {
VoiceCatResult(vc_leave_channel(handle)) VoiceCatResult(vc_leave_channel(handle))
} }
@discardableResult
public func joinVoice() -> VoiceCatResult {
VoiceCatResult(vc_join_voice(handle))
}
@discardableResult
public func leaveVoice() -> VoiceCatResult {
VoiceCatResult(vc_leave_voice(handle))
}
/// Pull the current channel tree. Re-call after `.channelList`/`.userJoined`/`.userLeft`/ /// Pull the current channel tree. Re-call after `.channelList`/`.userJoined`/`.userLeft`/
/// `.userUpdated` events. The native list is freed inside this call callers never /// `.userUpdated` events. The native list is freed inside this call callers never
/// manage native lifetime. /// manage native lifetime.

View File

@@ -98,7 +98,7 @@ final class AppState {
} }
func disconnect() { func disconnect() {
session?.stopMicStream() session?.leaveVoice()
session?.client.disconnect() session?.client.disconnect()
IOSAudioEngine.shared.stop() IOSAudioEngine.shared.stop()
AudioSessionManager.shared.deactivateSession() AudioSessionManager.shared.deactivateSession()

View File

@@ -20,6 +20,7 @@ struct ActivityEntry: Identifiable {
struct VoiceState { struct VoiceState {
var micActive = false var micActive = false
var voiceSubscribed = false
var selfMuted = false var selfMuted = false
var selfDeafened = false var selfDeafened = false
var serverMuted = false var serverMuted = false
@@ -151,8 +152,6 @@ final class SessionState {
addActivity("Sharing screen audio (\(channels == 2 ? "stereo" : "mono"))") addActivity("Sharing screen audio (\(channels == 2 ? "stereo" : "mono"))")
break break
} }
// A remote user started a stream ensure the audio session is active so we can
// hear them even if we haven't joined voice ourselves.
if ev.userId != selfUserId { if ev.userId != selfUserId {
do { do {
try AudioSessionManager.shared.ensureSessionActive() try AudioSessionManager.shared.ensureSessionActive()
@@ -163,7 +162,25 @@ final class SessionState {
AudioSessionManager.shared.logSessionState("stream started (user \(ev.userId))") AudioSessionManager.shared.logSessionState("stream started (user \(ev.userId))")
addActivity("Stream started (user \(ev.userId))") addActivity("Stream started (user \(ev.userId))")
case .streamStopped: case .streamStopped:
if ev.userId == selfUserId {
if ev.streamId == voiceState.localStreamId {
voiceState.localStreamId = 0
voiceState.micActive = false
voiceState.level = 0
}
} else {
addActivity("Stream stopped (user \(ev.userId))") addActivity("Stream stopped (user \(ev.userId))")
}
case .voiceState:
let subscribed = ev.u32a != 0
voiceState.voiceSubscribed = subscribed
if subscribed {
doStartMicStream()
} else {
voiceState.micActive = false
voiceState.level = 0
EventFeedback.shared.play(.voiceOff)
}
case .joinResult: case .joinResult:
if ev.result == .ok { if ev.result == .ok {
currentChannelId = ev.channelId currentChannelId = ev.channelId
@@ -234,12 +251,15 @@ final class SessionState {
currentChannelId = 0 currentChannelId = 0
} }
func startMicStream() { func joinVoice() {
AVAudioApplication.requestRecordPermission { [weak self] granted in AVAudioApplication.requestRecordPermission { [weak self] granted in
DispatchQueue.main.async { DispatchQueue.main.async {
guard let self else { return } guard let self else { return }
if granted { if granted {
self.doStartMicStream() let result = self.client.joinVoice()
if result != .ok {
self.addActivity("Failed to join voice: \(result.description)")
}
} else { } else {
self.addActivity("Microphone permission denied — grant in Settings > Privacy > Microphone") self.addActivity("Microphone permission denied — grant in Settings > Privacy > Microphone")
} }
@@ -255,10 +275,6 @@ final class SessionState {
return return
} }
// Unified iOS path: the core is always external (set at connect via setExternalPlayback +
// every MIC stream external_feed), and `IOSAudioEngine` drives capture + playback. So the
// mic stream is just started with external_feed=true and the engine is told the mic is now
// active no setExternalPlayback toggle, no audioRestart ordering, no VPIO/miniaudio fork.
let desc = StreamDescriptor(kind: .mic, deviceId: voiceState.currentDeviceId, label: "Mic", let desc = StreamDescriptor(kind: .mic, deviceId: voiceState.currentDeviceId, label: "Mic",
externalFeed: true) externalFeed: true)
let (result, streamId) = client.startStream(desc) let (result, streamId) = client.startStream(desc)
@@ -274,25 +290,14 @@ final class SessionState {
if channels != 1 { if channels != 1 {
client.setCaptureChannels(streamId: streamId, channels: channels) client.setCaptureChannels(streamId: streamId, channels: channels)
} }
// Engage the mic: installs the input tap and (per preset) VPIO, in one engine rebuild.
IOSAudioEngine.shared.startMic(streamId: streamId, channels: channels) IOSAudioEngine.shared.startMic(streamId: streamId, channels: channels)
} }
func stopMicStream() { func leaveVoice() {
// Disengage the mic (removes the tap + VPIO) but keep the engine running for any remaining if voiceState.screenStreamId != 0 { stopScreenShare() }
// remote audio. Then stop the core's MIC stream. The core stays external throughout no client.setPushToTalk(false)
// setExternalPlayback toggle, no audioRestart.
IOSAudioEngine.shared.stopMic() IOSAudioEngine.shared.stopMic()
if voiceState.localStreamId != 0 { client.leaveVoice()
client.stopStream(voiceState.localStreamId)
voiceState.localStreamId = 0
EventFeedback.shared.play(.voiceOff)
}
voiceState.micActive = false
voiceState.level = 0
// Do NOT deactivate the AVAudioSession here the user may still want to hear
// remote audio (other people talking). The session is deactivated only when
// disconnecting from the server (see AppState.disconnect / .disconnected event).
} }
// MARK: - Screen audio share // MARK: - Screen audio share

View File

@@ -15,8 +15,7 @@ struct ChannelEditView: View {
@State private var maxUsers = "0" @State private var maxUsers = "0"
@State private var sortOrder = "0" @State private var sortOrder = "0"
// Audio (Opus). Note: the channel list does not carry the current audio config, so when // Audio (Opus) populated from the channel's current config when editing.
// editing an existing channel these start from the codec defaults (same as macOS/Windows).
@State private var stereo = false @State private var stereo = false
@State private var bitrate = "64000" @State private var bitrate = "64000"
@State private var sampleRate = "48000" @State private var sampleRate = "48000"
@@ -122,6 +121,17 @@ struct ChannelEditView: View {
parentId = ch.parentId parentId = ch.parentId
passwordProtected = ch.passwordProtected passwordProtected = ch.passwordProtected
maxUsers = "\(ch.maxUsers)" maxUsers = "\(ch.maxUsers)"
sortOrder = "\(ch.sortOrder)"
stereo = ch.audio.stereo
bitrate = "\(ch.audio.bitrateBps)"
sampleRate = "\(ch.audio.sampleRate)"
frameMs = ch.audio.frameMs
application = ch.audio.application
packetLoss = "\(ch.audio.expectedPacketLoss)"
complexity = ch.audio.complexity
fec = ch.audio.fec
dtx = ch.audio.dtx
dred = ch.audio.dred
} }
private func save() { private func save() {

View File

@@ -294,7 +294,7 @@ struct SettingsView: View {
// MARK: - Server // MARK: - Server
Section("Server") { Section("Server") {
Button(role: .destructive) { Button(role: .destructive) {
session.stopMicStream() session.leaveVoice()
appState.disconnect() appState.disconnect()
} label: { } label: {
Label("Disconnect", systemImage: "phone.down") Label("Disconnect", systemImage: "phone.down")

View File

@@ -14,9 +14,9 @@ struct VoiceControlsView: View {
} else { } else {
Button { Button {
if session.voiceState.micActive { if session.voiceState.micActive {
session.stopMicStream() session.leaveVoice()
} else { } else {
session.startMicStream() session.joinVoice()
} }
} label: { } label: {
Text(session.voiceState.micActive ? "Leave Voice" : "Join Voice") Text(session.voiceState.micActive ? "Leave Voice" : "Join Voice")
@@ -28,7 +28,7 @@ struct VoiceControlsView: View {
.foregroundStyle(session.voiceState.micActive ? .green : .accentColor) .foregroundStyle(session.voiceState.micActive ? .green : .accentColor)
} }
.disabled(session.currentChannelId == 0) .disabled(session.currentChannelId == 0)
.accessibilityLabel(session.voiceState.micActive ? "Leave Voice — stop sending microphone audio" : "Join Voice — start sending microphone audio") .accessibilityLabel(session.voiceState.micActive ? "Leave Voice" : "Join Voice")
} }
// Level meter // Level meter
@@ -81,7 +81,7 @@ struct VoiceControlsView: View {
// Disconnect // Disconnect
Button(role: .destructive) { Button(role: .destructive) {
session.stopMicStream() session.leaveVoice()
session.client.disconnect() session.client.disconnect()
} label: { } label: {
Image(systemName: "phone.down.fill") Image(systemName: "phone.down.fill")
@@ -115,13 +115,13 @@ private struct PTTButton: View {
.updating($isPressing) { _, state, _ in state = true } .updating($isPressing) { _, state, _ in state = true }
.onChanged { _ in .onChanged { _ in
if !isPressing { return } if !isPressing { return }
if !session.voiceState.micActive { session.startMicStream() } if !session.voiceState.micActive { session.joinVoice() }
session.setPushToTalk(true) session.setPushToTalk(true)
UIImpactFeedbackGenerator(style: .medium).impactOccurred() UIImpactFeedbackGenerator(style: .medium).impactOccurred()
} }
.onEnded { _ in .onEnded { _ in
session.setPushToTalk(false) session.setPushToTalk(false)
session.stopMicStream() session.leaveVoice()
} }
) )
.accessibilityLabel("Push to talk, hold to transmit") .accessibilityLabel("Push to talk, hold to transmit")

View File

@@ -497,7 +497,8 @@ final class MainWindowController: NSWindowController, NSWindowDelegate {
isGuest: true, isGuest: true,
channelId: event.channelId, channelId: event.channelId,
selfMicMuted: false, selfDeafened: false, selfMicMuted: false, selfDeafened: false,
serverMuted: false, serverDeafened: false) serverMuted: false, serverDeafened: false,
voiceSubscribed: false)
users[event.userId] = u users[event.userId] = u
refreshChannelTree(); refreshUserList() refreshChannelTree(); refreshUserList()
if event.channelId == currentChannelId && event.userId != selfUserId { if event.channelId == currentChannelId && event.userId != selfUserId {
@@ -540,7 +541,8 @@ final class MainWindowController: NSWindowController, NSWindowDelegate {
users[selfUserId] = User(id: self_.id, nickname: self_.nickname, isGuest: self_.isGuest, users[selfUserId] = User(id: self_.id, nickname: self_.nickname, isGuest: self_.isGuest,
channelId: event.channelId, channelId: event.channelId,
selfMicMuted: self_.selfMicMuted, selfDeafened: self_.selfDeafened, selfMicMuted: self_.selfMicMuted, selfDeafened: self_.selfDeafened,
serverMuted: self_.serverMuted, serverDeafened: self_.serverDeafened) serverMuted: self_.serverMuted, serverDeafened: self_.serverDeafened,
voiceSubscribed: self_.voiceSubscribed)
} }
refreshChannelTree(); refreshUserList(); updateStatusLabel() refreshChannelTree(); refreshUserList(); updateStatusLabel()
let name = channels.first(where: { $0.id == event.channelId })?.name ?? "Channel #\(event.channelId)" let name = channels.first(where: { $0.id == event.channelId })?.name ?? "Channel #\(event.channelId)"
@@ -589,10 +591,18 @@ final class MainWindowController: NSWindowController, NSWindowDelegate {
addActivity("\(u.nickname) started \(kindStr) stream") addActivity("\(u.nickname) started \(kindStr) stream")
case .streamStopped: case .streamStopped:
if let u = users[event.userId], u.channelId == currentChannelId { if event.userId == selfUserId {
if event.streamId == micStreamId {
micStreamId = 0
settingsWindowController?.resetLevel()
}
} else if let u = users[event.userId], u.channelId == currentChannelId {
addActivity("\(u.nickname) stopped a stream") addActivity("\(u.nickname) stopped a stream")
} }
case .voiceState:
handleVoiceState(event)
case .disconnected: case .disconnected:
handleDisconnected(event) handleDisconnected(event)
@@ -783,14 +793,27 @@ final class MainWindowController: NSWindowController, NSWindowDelegate {
@objc private func micToggleClicked() { @objc private func micToggleClicked() {
if micStreamId == 0 { if micStreamId == 0 {
let result = client.joinVoice()
if result != .ok {
addActivity("Failed to join voice: \(result)")
}
} else {
stopAuxStream()
if screenStreamId != 0 { stopScreenAudio() }
client.setPushToTalk(false)
client.leaveVoice()
}
}
private func handleVoiceState(_ event: VoiceCatEvent) {
let subscribed = event.u32a != 0
if subscribed {
let (result, streamId) = client.startStream(StreamDescriptor(kind: .mic, deviceId: nil, label: "Microphone")) let (result, streamId) = client.startStream(StreamDescriptor(kind: .mic, deviceId: nil, label: "Microphone"))
if result == .ok { if result == .ok {
micStreamId = streamId micStreamId = streamId
if let devId = selectedInputDeviceId { if let devId = selectedInputDeviceId {
client.setInputDevice(streamId: streamId, deviceId: devId) client.setInputDevice(streamId: streamId, deviceId: devId)
} }
// Stored on the stream before the announce round-trip completes, so the core's
// first capture-device open (ensure_audio_running) picks up the channel count.
client.setCaptureChannels(streamId: streamId, channels: stereoMic ? 2 : 1) client.setCaptureChannels(streamId: streamId, channels: stereoMic ? 2 : 1)
client.setInputMode(selectedInputMode) client.setInputMode(selectedInputMode)
if selectedInputMode == .voiceActivation { if selectedInputMode == .voiceActivation {
@@ -803,14 +826,11 @@ final class MainWindowController: NSWindowController, NSWindowDelegate {
EventFeedback.shared.play(.voiceOn) EventFeedback.shared.play(.voiceOn)
NSAccessibility.post(element: logTextView, notification: .announcementRequested, NSAccessibility.post(element: logTextView, notification: .announcementRequested,
userInfo: [.announcement: "Joined voice", .priority: NSAccessibilityPriorityLevel.medium]) userInfo: [.announcement: "Joined voice", .priority: NSAccessibilityPriorityLevel.medium])
startAuxStream() // no-op unless the aux stream is enabled in settings startAuxStream()
} else { } else {
addActivity("Failed to start microphone: \(result)") addActivity("Failed to start microphone: \(result)")
} }
} else { } else {
stopAuxStream()
client.setPushToTalk(false)
client.stopStream(micStreamId)
micStreamId = 0 micStreamId = 0
settingsWindowController?.resetLevel() settingsWindowController?.resetLevel()
setVoiceJoinedState(false) setVoiceJoinedState(false)
@@ -933,6 +953,15 @@ final class MainWindowController: NSWindowController, NSWindowDelegate {
} }
} }
private func stopScreenAudio() {
guard screenStreamId != 0 else { return }
stopScreenCapture()
client.stopStream(screenStreamId)
screenStreamId = 0
setShareScreenButton(active: false)
addActivity("Stopped sharing screen audio")
}
/// Announce the SCREEN_AUDIO stream with the chosen selection in hand. ScreenCaptureKit /// Announce the SCREEN_AUDIO stream with the chosen selection in hand. ScreenCaptureKit
/// capture starts once the server's StreamAnnounceResult lands (the .streamStarted event), /// capture starts once the server's StreamAnnounceResult lands (the .streamStarted event),
/// when the effective audio config and thus the channel count is known. See /// when the effective audio config and thus the channel count is known. See
@@ -1449,7 +1478,7 @@ extension MainWindowController: NSMenuDelegate {
let info = ChannelEdit(id: ch.id, parentId: ch.parentId, name: ch.name, let info = ChannelEdit(id: ch.id, parentId: ch.parentId, name: ch.name,
topic: ch.topic, passwordProtected: ch.passwordProtected, topic: ch.topic, passwordProtected: ch.passwordProtected,
password: nil, maxUsers: ch.maxUsers, password: nil, maxUsers: ch.maxUsers,
sortOrder: 0, audio: AudioConfig()) sortOrder: ch.sortOrder, audio: ch.audio)
let sheet = ChannelEditSheet(channels: channels, editing: info) let sheet = ChannelEditSheet(channels: channels, editing: info)
sheet.onComplete = { [weak self] edited in sheet.onComplete = { [weak self] edited in
guard let edited else { return } guard let edited else { return }

View File

@@ -349,10 +349,21 @@ public partial class MainForm : Form
HandleStreamStarted(ev); HandleStreamStarted(ev);
break; break;
case VcEventType.StreamStopped: case VcEventType.StreamStopped:
if (_users.TryGetValue(ev.UserId, out var stUser) && if (ev.UserId == _selfUserId)
{
if (ev.StreamId == _micStreamId)
{
_micStreamId = 0;
pbLevel.Value = 0;
}
}
else if (_users.TryGetValue(ev.UserId, out var stUser) &&
stUser.ChannelId == _currentChannelId) stUser.ChannelId == _currentChannelId)
AddActivity($"{stUser.Nickname} stopped a stream"); AddActivity($"{stUser.Nickname} stopped a stream");
break; break;
case VcEventType.VoiceState:
HandleVoiceState(ev);
break;
case VcEventType.Disconnected: case VcEventType.Disconnected:
HandleDisconnected(ev); HandleDisconnected(ev);
break; break;
@@ -378,7 +389,7 @@ public partial class MainForm : Form
private void HandleUserJoined(VoiceCatEvent ev) private void HandleUserJoined(VoiceCatEvent ev)
{ {
var user = new UserInfo(ev.UserId, ev.Text ?? $"User#{ev.UserId}", false, ev.ChannelId, var user = new UserInfo(ev.UserId, ev.Text ?? $"User#{ev.UserId}", false, ev.ChannelId,
false, false, false, false); false, false, false, false, false);
_users[ev.UserId] = user; _users[ev.UserId] = user;
RefreshChannelTree(); RefreshChannelTree();
RefreshUserList(); RefreshUserList();
@@ -663,6 +674,32 @@ public partial class MainForm : Form
private void BtnMicToggle_Click(object? sender, EventArgs e) private void BtnMicToggle_Click(object? sender, EventArgs e)
{ {
if (_micStreamId == 0) if (_micStreamId == 0)
{
var result = _client.JoinVoice();
if (result != VcResult.Ok)
AddActivity($"Failed to join voice: {result}");
}
else
{
StopAuxStream();
if (_screenStreamId != 0) StopScreenAudio();
_client.SetPushToTalk(false);
_client.LeaveVoice();
}
}
private void SetVoiceJoinedState(bool joined)
{
tsbJoinVoice.Text = joined ? "Leave Voice" : "Join Voice";
_miJoinVoice.Text = joined ? "Leave &Voice" : "&Join Voice";
chkMute.Enabled = joined;
chkDeafen.Enabled = joined;
}
private void HandleVoiceState(VoiceCatEvent ev)
{
bool subscribed = ev.U32a != 0;
if (subscribed)
{ {
var (result, streamId) = _client.StartStream(VcStreamKind.Mic, "Microphone"); var (result, streamId) = _client.StartStream(VcStreamKind.Mic, "Microphone");
if (result == VcResult.Ok) if (result == VcResult.Ok)
@@ -680,7 +717,7 @@ public partial class MainForm : Form
SetVoiceJoinedState(true); SetVoiceJoinedState(true);
AddActivity("Joined voice — microphone active"); AddActivity("Joined voice — microphone active");
_feedback.PlaySound(SoundEvent.VoiceOn); _feedback.PlaySound(SoundEvent.VoiceOn);
StartAuxStream(); // no-op unless the aux stream is enabled in settings StartAuxStream();
} }
else else
{ {
@@ -689,25 +726,12 @@ public partial class MainForm : Form
} }
else else
{ {
StopAuxStream();
_client.SetPushToTalk(false);
_client.StopStream(_micStreamId);
_micStreamId = 0;
pbLevel.Value = 0;
SetVoiceJoinedState(false); SetVoiceJoinedState(false);
AddActivity("Left voice"); AddActivity("Left voice");
_feedback.PlaySound(SoundEvent.VoiceOff); _feedback.PlaySound(SoundEvent.VoiceOff);
} }
} }
private void SetVoiceJoinedState(bool joined)
{
tsbJoinVoice.Text = joined ? "Leave Voice" : "Join Voice";
_miJoinVoice.Text = joined ? "Leave &Voice" : "&Join Voice";
chkMute.Enabled = joined;
chkDeafen.Enabled = joined;
}
private void ApplySelfMute() => private void ApplySelfMute() =>
_client.SetSelfMute(chkMute.Checked, chkDeafen.Checked); _client.SetSelfMute(chkMute.Checked, chkDeafen.Checked);
@@ -1122,8 +1146,8 @@ public partial class MainForm : Form
var editInfo = new ChannelEditInfo( var editInfo = new ChannelEditInfo(
channel.Id, channel.ParentId, channel.Name, channel.Topic, channel.Id, channel.ParentId, channel.Name, channel.Topic,
channel.PasswordProtected, null, channel.MaxUsers, 0, channel.PasswordProtected, null, channel.MaxUsers, channel.SortOrder,
new AudioConfigInfo(0, false, 48000, 0, 20, 0, true, 0, false, 10, false)); channel.Audio);
using var dlg = new ChannelEditDialog(_channels, editInfo); using var dlg = new ChannelEditDialog(_channels, editInfo);
if (dlg.ShowDialog(this) != DialogResult.OK || dlg.Result is null) return; if (dlg.ShowDialog(this) != DialogResult.OK || dlg.Result is null) return;

View File

@@ -92,6 +92,8 @@ public enum VcEventType
GenericResult = 14, GenericResult = 14,
/// <summary>M5: reply to VoiceCatClient.RequestAccountList — call ListAccounts() to read.</summary> /// <summary>M5: reply to VoiceCatClient.RequestAccountList — call ListAccounts() to read.</summary>
AccountList = 15, AccountList = 15,
/// <summary>Voice-plane subscription state. u32a = 1 (subscribed) or 0 (unsubscribed).</summary>
VoiceState = 16,
} }
/// <summary> /// <summary>

View File

@@ -38,7 +38,9 @@ internal static class Marshaling
Marshal.PtrToStringUTF8(raw.Name) ?? string.Empty, Marshal.PtrToStringUTF8(raw.Name) ?? string.Empty,
Marshal.PtrToStringUTF8(raw.Topic) ?? string.Empty, Marshal.PtrToStringUTF8(raw.Topic) ?? string.Empty,
raw.PasswordProtected != 0, raw.PasswordProtected != 0,
raw.MaxUsers)); raw.MaxUsers,
raw.SortOrder,
ToManaged(in raw.Audio)));
} }
NativeMethods.vc_free_channel_list(ref native); NativeMethods.vc_free_channel_list(ref native);
return result; return result;
@@ -59,7 +61,8 @@ internal static class Marshaling
raw.SelfMicMuted != 0, raw.SelfMicMuted != 0,
raw.SelfDeafened != 0, raw.SelfDeafened != 0,
raw.ServerMuted != 0, raw.ServerMuted != 0,
raw.ServerDeafened != 0)); raw.ServerDeafened != 0,
raw.VoiceSubscribed != 0));
} }
NativeMethods.vc_free_user_list(ref native); NativeMethods.vc_free_user_list(ref native);
return result; return result;

View File

@@ -9,7 +9,9 @@ public sealed record ChannelInfo(
string Name, string Name,
string Topic, string Topic,
bool PasswordProtected, bool PasswordProtected,
uint MaxUsers); uint MaxUsers,
uint SortOrder,
AudioConfigInfo Audio);
public sealed record ChannelEditInfo( public sealed record ChannelEditInfo(
uint Id, uint Id,
@@ -30,7 +32,8 @@ public sealed record UserInfo(
bool SelfMicMuted, bool SelfMicMuted,
bool SelfDeafened, bool SelfDeafened,
bool ServerMuted, bool ServerMuted,
bool ServerDeafened); bool ServerDeafened,
bool VoiceSubscribed);
public sealed record PermissionsInfo( public sealed record PermissionsInfo(
bool CanCreateTempChannel, bool CanCreateTempChannel,

View File

@@ -54,6 +54,12 @@ internal static partial class NativeMethods
[LibraryImport(LibName)] [LibraryImport(LibName)]
internal static partial VcResult vc_leave_channel(nint c); internal static partial VcResult vc_leave_channel(nint c);
[LibraryImport(LibName)]
internal static partial VcResult vc_join_voice(nint c);
[LibraryImport(LibName)]
internal static partial VcResult vc_leave_voice(nint c);
// ── Local media streams ───────────────────────────────────────────────────────────────── // ── Local media streams ─────────────────────────────────────────────────────────────────
[LibraryImport(LibName)] [LibraryImport(LibName)]
internal static partial VcResult vc_stream_start(nint c, in VcStreamDescNative desc, internal static partial VcResult vc_stream_start(nint c, in VcStreamDescNative desc,

View File

@@ -116,6 +116,8 @@ internal struct VcChannelNative
public IntPtr Topic; public IntPtr Topic;
public int PasswordProtected; public int PasswordProtected;
public uint MaxUsers; public uint MaxUsers;
public uint SortOrder;
public VcAudioConfigNative Audio;
} }
[StructLayout(LayoutKind.Sequential)] [StructLayout(LayoutKind.Sequential)]
@@ -136,6 +138,7 @@ internal struct VcUserNative
public int SelfDeafened; public int SelfDeafened;
public int ServerMuted; public int ServerMuted;
public int ServerDeafened; public int ServerDeafened;
public int VoiceSubscribed;
} }
[StructLayout(LayoutKind.Sequential)] [StructLayout(LayoutKind.Sequential)]

View File

@@ -152,6 +152,12 @@ public sealed class VoiceCatClient : IDisposable
public VcResult LeaveChannel() => public VcResult LeaveChannel() =>
NativeMethods.vc_leave_channel(_handle.DangerousGetHandle()); NativeMethods.vc_leave_channel(_handle.DangerousGetHandle());
public VcResult JoinVoice() =>
NativeMethods.vc_join_voice(_handle.DangerousGetHandle());
public VcResult LeaveVoice() =>
NativeMethods.vc_leave_voice(_handle.DangerousGetHandle());
public List<ChannelInfo> ListChannels() public List<ChannelInfo> ListChannels()
{ {
NativeMethods.vc_list_channels(_handle.DangerousGetHandle(), out var native); NativeMethods.vc_list_channels(_handle.DangerousGetHandle(), out var native);

View File

@@ -145,6 +145,10 @@ typedef enum vc_event_type {
vc_delete_channel/vc_create_account/vc_reset_password/ vc_delete_channel/vc_create_account/vc_reset_password/
vc_delete_account. */ vc_delete_account. */
VC_EVENT_ACCOUNT_LIST = 15, /* Reply to vc_list_accounts. */ VC_EVENT_ACCOUNT_LIST = 15, /* Reply to vc_list_accounts. */
VC_EVENT_VOICE_STATE = 16, /* u32a = subscribed(0/1). Reply to vc_join_voice()/
vc_leave_voice(), and also emitted when the server
changes your voice-subscription state. The user list
(vc_user) carries per-user voice_subscribed. */
} vc_event_type; } vc_event_type;
/* TOFU server-identity classification (M4) — see VC_EVENT_SERVER_IDENTITY and /* TOFU server-identity classification (M4) — see VC_EVENT_SERVER_IDENTITY and
@@ -292,6 +296,10 @@ typedef struct vc_channel {
const char* topic; const char* topic;
int password_protected; /* bool */ int password_protected; /* bool */
uint32_t max_users; /* 0 = unlimited */ uint32_t max_users; /* 0 = unlimited */
uint32_t sort_order; /* channel sort order */
/* Authoritative channel Opus params (docs/voice.md §3). Populated from the Channel proto
* so the edit dialog can read back the current config without a separate round-trip. */
vc_audio_config audio;
} vc_channel; } vc_channel;
typedef struct vc_channel_list { typedef struct vc_channel_list {
@@ -308,6 +316,7 @@ typedef struct vc_user {
int self_deafened; /* bool */ int self_deafened; /* bool */
int server_muted; /* bool */ int server_muted; /* bool */
int server_deafened; /* bool */ int server_deafened; /* bool */
int voice_subscribed; /* bool — true when on the voice plane */
} vc_user; } vc_user;
typedef struct vc_user_list { typedef struct vc_user_list {
@@ -366,6 +375,16 @@ VC_API vc_result vc_join_channel(vc_client* c, uint32_t channel_id,
const char* password /* nullable */); const char* password /* nullable */);
VC_API vc_result vc_leave_channel(vc_client* c); VC_API vc_result vc_leave_channel(vc_client* c);
/* ── Voice-plane subscription ─────────────────────────────────────────────── */
/* Joining voice subscribes to the voice plane: the server starts relaying voice frames
* to you, and the core wires up remote-stream decoders so you hear other users. Leaving
* voice unsubscribes: the server stops relaying voice to you, the core tears down all
* remote decoders (playback stops), and any active local mic/screen/aux streams are
* stopped. Text chat is unaffected either way. Result arrives as VC_EVENT_VOICE_STATE
* (u32a = 1 for subscribed, 0 for unsubscribed). */
VC_API vc_result vc_join_voice(vc_client* c);
VC_API vc_result vc_leave_voice(vc_client* c);
/* ── Local media streams (mic / screen audio / aux) ───────────────────────── */ /* ── Local media streams (mic / screen audio / aux) ───────────────────────── */
VC_API vc_result vc_stream_start(vc_client* c, const vc_stream_desc* desc, VC_API vc_result vc_stream_start(vc_client* c, const vc_stream_desc* desc,
uint32_t* out_stream_id); uint32_t* out_stream_id);

View File

@@ -44,6 +44,9 @@ message Envelope {
StreamStop stream_stop = 42; StreamStop stream_stop = 42;
StreamStateUpdate stream_state = 43; StreamStateUpdate stream_state = 43;
UdpBinding udp_binding = 44; UdpBinding udp_binding = 44;
SubscribeVoiceRequest subscribe_voice = 45;
UnsubscribeVoiceRequest unsubscribe_voice = 46;
VoiceSubscriptionResult voice_subscription_result = 47;
// Text (5059) // Text (5059)
TextMessage text_message = 50; TextMessage text_message = 50;
@@ -120,6 +123,7 @@ message User {
bool server_muted = 7; bool server_muted = 7;
repeated StreamInfo streams = 8; repeated StreamInfo streams = 8;
bool server_deafened = 9; // M5: server-imposed deafen bool server_deafened = 9; // M5: server-imposed deafen
bool voice_subscribed = 10; // true when the user is on the voice plane (hears + can send)
} }
message Permissions { message Permissions {
@@ -216,6 +220,14 @@ message StreamStop { uint32 stream_id = 1; }
message StreamStateUpdate { uint32 user_id = 1; uint32 stream_id = 2; bool muted = 3; bool talking = 4; } message StreamStateUpdate { uint32 user_id = 1; uint32 stream_id = 2; bool muted = 3; bool talking = 4; }
message UdpBinding { bytes udp_token = 1; bool ack = 2; } message UdpBinding { bytes udp_token = 1; bool ack = 2; }
// Voice-plane subscription. Joining voice = you can hear and send voice in your channel;
// leaving voice = the server stops relaying voice frames to you and you stop receiving.
// Text chat is unaffected. The server toggles a per-session flag checked by the SFU relay's
// recipient filter (docs/protocol.md §5). Result arrives as voice_subscription_result.
message SubscribeVoiceRequest {}
message UnsubscribeVoiceRequest {}
message VoiceSubscriptionResult { bool ok = 1; string error = 2; bool subscribed = 3; }
// ── Text (ephemeral — server does not persist history, docs/protocol.md §5) ────── // ── Text (ephemeral — server does not persist history, docs/protocol.md §5) ──────
message TextMessage { message TextMessage {
TextScope scope = 1; TextScope scope = 1;

View File

@@ -500,6 +500,9 @@ void vc_client::handle_envelope(const voicecat::v1::Envelope& env) {
case voicecat::v1::Envelope::kStreamAnnounceResult: case voicecat::v1::Envelope::kStreamAnnounceResult:
handle_stream_announce_result(env.request_id(), env.stream_announce_result()); handle_stream_announce_result(env.request_id(), env.stream_announce_result());
break; break;
case voicecat::v1::Envelope::kVoiceSubscriptionResult:
handle_voice_subscription_result(env.voice_subscription_result());
break;
case voicecat::v1::Envelope::kPong: { case voicecat::v1::Envelope::kPong: {
// Correlate the echoed nonce to measure RTT (docs/protocol.md §7). // Correlate the echoed nonce to measure RTT (docs/protocol.md §7).
uint64_t nonce = env.pong().nonce(); uint64_t nonce = env.pong().nonce();
@@ -618,10 +621,44 @@ void vc_client::handle_server_state(const voicecat::v1::ServerStateSnapshot& sna
} }
void vc_client::handle_channel_event(const voicecat::v1::ChannelEvent& ce) { void vc_client::handle_channel_event(const voicecat::v1::ChannelEvent& ce) {
bool audio_changed = false;
uint32_t updated_channel_id = 0;
{ {
std::lock_guard<std::mutex> lk(session_model_mu_); std::lock_guard<std::mutex> lk(session_model_mu_);
// Detect audio-config changes on our current channel BEFORE apply_channel_event
// overwrites the stored config. Encoders/decoders are frozen at StreamAnnounceResult
// time (docs/voice.md §3), so a channel audio change requires restarting active
// local streams to pick up the new params.
if (ce.kind() == voicecat::v1::ChannelEvent::UPDATED) {
updated_channel_id = ce.channel().id();
const auto* self_user = session_model_.find_user(self_user_id_);
if (self_user && self_user->channel_id == updated_channel_id) {
const auto* old_ch = session_model_.find_channel(updated_channel_id);
const auto& na = ce.channel().audio();
if (old_ch) {
audio_changed =
old_ch->audio_sample_rate != (na.sample_rate() ? na.sample_rate() : 48000) ||
old_ch->audio_frame_ms != (na.frame_ms() ? na.frame_ms() : 20) ||
old_ch->audio_mode != static_cast<uint32_t>(na.mode()) ||
old_ch->audio_bitrate_bps != na.bitrate_bps() ||
old_ch->audio_application != static_cast<uint32_t>(na.application()) ||
old_ch->audio_fec != na.fec() ||
old_ch->audio_expected_packet_loss != na.expected_packet_loss() ||
old_ch->audio_dtx != na.dtx() ||
old_ch->audio_complexity != (na.complexity() ? na.complexity() : 10) ||
old_ch->audio_dred != na.dred();
}
}
}
session_model_.apply_channel_event(ce); session_model_.apply_channel_event(ce);
} }
if (audio_changed) {
restart_active_streams_for_channel(updated_channel_id);
}
// Same "go look" signal vc_list_channels' callers already poll on after the initial // Same "go look" signal vc_list_channels' callers already poll on after the initial
// snapshot — see voicecat.h's VC_EVENT_CHANNEL_LIST doc comment. // snapshot — see voicecat.h's VC_EVENT_CHANNEL_LIST doc comment.
vc_event ev{}; vc_event ev{};
@@ -1215,6 +1252,10 @@ void vc_client::ensure_audio_running() {
void vc_client::sync_remote_streams(const voicecat::v1::User& user) { void vc_client::sync_remote_streams(const voicecat::v1::User& user) {
if (user.id() == self_user_id_) return; if (user.id() == self_user_id_) return;
// Don't wire up remote decoders when not on the voice plane — the server isn't relaying
// voice to us, and decoders would sit idle consuming memory. resync_remote_streams()
// wires them up on voice join.
if (!voice_subscribed_.load(std::memory_order_acquire)) return;
std::vector<uint32_t> current_ssrcs; std::vector<uint32_t> current_ssrcs;
for (const auto& si : user.streams()) current_ssrcs.push_back(si.ssrc()); for (const auto& si : user.streams()) current_ssrcs.push_back(si.ssrc());
@@ -1285,6 +1326,7 @@ vc_result vc_client::stream_start(const vc_stream_desc& desc, uint32_t* out_stre
ls.stream_id = sid; ls.stream_id = sid;
ls.pending = true; ls.pending = true;
ls.external_feed = (desc.external_feed != 0); ls.external_feed = (desc.external_feed != 0);
ls.label = desc.label ? desc.label : "";
if (out_stream_id) *out_stream_id = sid; if (out_stream_id) *out_stream_id = sid;
req_id = next_req_id_++; req_id = next_req_id_++;
@@ -1441,6 +1483,141 @@ vc_client::LocalStream* vc_client::find_local_stream_by_id(uint32_t stream_id) {
return nullptr; return nullptr;
} }
void vc_client::restart_active_streams_for_channel(uint32_t channel_id) {
// Collect active stream info under the lock, then stop→start outside the lock (stream_stop
// and stream_start both acquire local_streams_mu_). capture_device_id and capture_channels
// survive the restart: stream_stop doesn't clear them, and stream_start reuses the existing
// LocalStream entry (auto& ls = local_streams_[kind]) without resetting them.
struct ActiveStreamInfo {
int kind;
std::string label;
int external_feed;
uint32_t stream_id;
};
std::vector<ActiveStreamInfo> active;
{
std::lock_guard lk(local_streams_mu_);
for (auto& [kind, ls] : local_streams_) {
if (!ls.active.load(std::memory_order_acquire)) continue;
active.push_back({kind, ls.label, ls.external_feed ? 1 : 0, ls.stream_id});
}
}
for (const auto& info : active) {
stream_stop(info.stream_id);
vc_stream_desc desc{};
desc.kind = static_cast<vc_stream_kind>(info.kind);
desc.device_id = nullptr; // default — capture_device_id is retained on the LocalStream
desc.label = info.label.c_str();
desc.external_feed = info.external_feed;
uint32_t new_sid = 0;
stream_start(desc, &new_sid);
}
}
// ── Voice-plane subscription ──────────────────────────────────────────────────
vc_result vc_client::join_voice() {
if (state_net_.load(std::memory_order_acquire) != VC_STATE_CONNECTED) return VC_ERR_NOT_CONNECTED;
voicecat::v1::Envelope req;
req.set_request_id(next_req_id_++);
req.mutable_subscribe_voice();
queue_envelope(req);
return VC_OK;
}
vc_result vc_client::leave_voice() {
if (state_net_.load(std::memory_order_acquire) != VC_STATE_CONNECTED) return VC_ERR_NOT_CONNECTED;
voicecat::v1::Envelope req;
req.set_request_id(next_req_id_++);
req.mutable_unsubscribe_voice();
queue_envelope(req);
return VC_OK;
}
void vc_client::handle_voice_subscription_result(
const voicecat::v1::VoiceSubscriptionResult& msg) {
bool subscribed = msg.subscribed();
voice_subscribed_.store(subscribed, std::memory_order_release);
if (subscribed) {
// Wire up remote-stream decoders for all users already in the session model.
// While we were unsubscribed, sync_remote_streams was gated off, so no decoders
// exist — resync from the current session state.
resync_remote_streams();
} else {
// Stop all active local streams (emits VC_EVENT_STREAM_STOPPED for each) and tear
// down every remote decoder so playback ceases. The server has already stopped
// relaying voice to us.
stop_all_local_streams();
clear_all_remote_streams();
}
vc_event ev{};
ev.type = VC_EVENT_VOICE_STATE;
ev.u32a = subscribed ? 1 : 0;
emit(ev);
}
void vc_client::resync_remote_streams() {
bool any_added = false;
{
std::lock_guard<std::mutex> sm_lk(session_model_mu_);
std::lock_guard rs_lk(remote_streams_mu_);
for (const auto& u : session_model_.users()) {
if (u.id == self_user_id_) continue;
for (const auto& s : u.streams) {
if (remote_streams_.count(s.ssrc)) continue;
remote_streams_[s.ssrc] = {u.id, s.stream_id};
voicecat::v1::AudioConfig audio;
audio.set_sample_rate(s.sample_rate);
audio.set_frame_ms(s.frame_ms);
audio.set_mode(static_cast<voicecat::v1::ChannelMode>(s.mode));
audio.set_bitrate_bps(s.bitrate_bps);
audio.set_application(static_cast<voicecat::v1::OpusApplication>(s.application));
audio.set_fec(s.fec);
audio.set_expected_packet_loss(s.expected_packet_loss);
audio.set_dtx(s.dtx);
audio.set_complexity(s.complexity);
audio.set_dred(s.dred);
voicecat::codec::OpusParams p = opus_params_from_audio_config(audio);
bool is_voice = (s.kind == static_cast<int>(voicecat::v1::STREAM_MIC));
audio_engine_.init_recv_stream(s.ssrc, p, u.id, s.stream_id, is_voice);
bool muted = self_deafened_.load(std::memory_order_acquire) ||
server_deafened_.load(std::memory_order_acquire);
audio_engine_.set_stream_mute(s.ssrc, muted);
any_added = true;
}
}
}
if (any_added) ensure_audio_running();
}
void vc_client::clear_all_remote_streams() {
std::lock_guard lk(remote_streams_mu_);
for (auto& [ssrc, info] : remote_streams_) {
audio_engine_.remove_stream(ssrc);
}
remote_streams_.clear();
}
void vc_client::stop_all_local_streams() {
std::vector<uint32_t> active_ids;
{
std::lock_guard lk(local_streams_mu_);
for (auto& [kind, ls] : local_streams_) {
if (ls.active.load(std::memory_order_acquire)) {
active_ids.push_back(ls.stream_id);
}
}
}
for (uint32_t sid : active_ids) {
stream_stop(sid);
}
}
vc_result vc_client::set_input_device(uint32_t stream_id, const char* device_id) { vc_result vc_client::set_input_device(uint32_t stream_id, const char* device_id) {
int kind = -1; int kind = -1;
{ {
@@ -1734,6 +1911,19 @@ vc_result vc_client::list_channels(vc_channel_list* out) {
items[i].topic = topic; items[i].topic = topic;
items[i].password_protected = ch.password_protected ? 1 : 0; items[i].password_protected = ch.password_protected ? 1 : 0;
items[i].max_users = ch.max_users; items[i].max_users = ch.max_users;
items[i].sort_order = ch.sort_order;
auto& ac = items[i].audio;
ac.codec = 0; // OPUS
ac.mode = ch.audio_mode;
ac.sample_rate = ch.audio_sample_rate;
ac.bitrate_bps = ch.audio_bitrate_bps;
ac.frame_ms = ch.audio_frame_ms;
ac.application = ch.audio_application;
ac.fec = ch.audio_fec ? 1 : 0;
ac.expected_packet_loss = ch.audio_expected_packet_loss;
ac.dtx = ch.audio_dtx ? 1 : 0;
ac.complexity = ch.audio_complexity;
ac.dred = ch.audio_dred ? 1 : 0;
} }
out->items = items; out->items = items;
out->count = channels.size(); out->count = channels.size();
@@ -1756,6 +1946,7 @@ vc_result vc_client::list_users(vc_user_list* out) {
items[i].self_deafened = u.self_deafened ? 1 : 0; items[i].self_deafened = u.self_deafened ? 1 : 0;
items[i].server_muted = u.server_muted ? 1 : 0; items[i].server_muted = u.server_muted ? 1 : 0;
items[i].server_deafened = u.server_deafened ? 1 : 0; items[i].server_deafened = u.server_deafened ? 1 : 0;
items[i].voice_subscribed = u.voice_subscribed ? 1 : 0;
} }
out->items = items; out->items = items;
out->count = users.size(); out->count = users.size();
@@ -2036,6 +2227,8 @@ vc_result vc_client::authenticate_guest(const char*) { return VC_ERR_NOT_
vc_result vc_client::authenticate_user(const char*, const char*) { return VC_ERR_NOT_IMPLEMENTED; } vc_result vc_client::authenticate_user(const char*, const char*) { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::join_channel(uint32_t, const char*) { return VC_ERR_NOT_IMPLEMENTED; } vc_result vc_client::join_channel(uint32_t, const char*) { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::leave_channel() { return VC_ERR_NOT_IMPLEMENTED; } vc_result vc_client::leave_channel() { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::join_voice() { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::leave_voice() { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::stream_start(const vc_stream_desc&, uint32_t*) { return VC_ERR_NOT_IMPLEMENTED; } vc_result vc_client::stream_start(const vc_stream_desc&, uint32_t*) { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::stream_stop(uint32_t) { return VC_ERR_NOT_IMPLEMENTED; } vc_result vc_client::stream_stop(uint32_t) { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::set_input_device(uint32_t, const char*) { return VC_ERR_NOT_IMPLEMENTED; } vc_result vc_client::set_input_device(uint32_t, const char*) { return VC_ERR_NOT_IMPLEMENTED; }

View File

@@ -46,6 +46,9 @@ struct vc_client {
vc_result join_channel(uint32_t channel_id, const char* password); vc_result join_channel(uint32_t channel_id, const char* password);
vc_result leave_channel(); vc_result leave_channel();
vc_result join_voice();
vc_result leave_voice();
vc_result stream_start(const vc_stream_desc& desc, uint32_t* out_stream_id); vc_result stream_start(const vc_stream_desc& desc, uint32_t* out_stream_id);
vc_result stream_stop(uint32_t stream_id); vc_result stream_stop(uint32_t stream_id);
vc_result set_input_device(uint32_t stream_id, const char* device_id); vc_result set_input_device(uint32_t stream_id, const char* device_id);
@@ -161,6 +164,11 @@ struct vc_client {
uint64_t server_session_id_{0}; uint64_t server_session_id_{0};
std::atomic<uint64_t> next_req_id_{1}; std::atomic<uint64_t> next_req_id_{1};
// Voice-plane subscription state. When false, the server does not relay voice frames to
// us, and we do not wire up remote-stream decoders (sync_remote_streams is gated on this).
// Toggled by vc_join_voice / vc_leave_voice; confirmed via VoiceSubscriptionResult.
std::atomic<bool> voice_subscribed_{false};
// Keepalive: client sends a Ping every ~15s (docs/protocol.md §7) so the server's // Keepalive: client sends a Ping every ~15s (docs/protocol.md §7) so the server's
// last_seen stays fresh and the reaper doesn't drop us. Pong echoes the nonce, which // last_seen stays fresh and the reaper doesn't drop us. Pong echoes the nonce, which
// we correlate to measure RTT. The read loop's 50ms TLS timeout means it spins fast // we correlate to measure RTT. The read loop's 50ms TLS timeout means it spins fast
@@ -265,6 +273,11 @@ struct vc_client {
// time and checked in handle_stream_announce_result / stream_stop. // time and checked in handle_stream_announce_result / stream_stop.
bool external_feed = false; bool external_feed = false;
// Stream label (from vc_stream_desc.label at stream_start time). Retained so the
// channel-update stream restart (restart_active_streams_for_channel) can re-announce
// with the same label without the caller's involvement.
std::string label;
// Reframe buffer: the AudioEngine clock is fixed at 48 kHz / 20 ms, so capture/feed // Reframe buffer: the AudioEngine clock is fixed at 48 kHz / 20 ms, so capture/feed
// always delivers 960-sample frames — but the channel's frame_ms (docs/voice.md §3) // always delivers 960-sample frames — but the channel's frame_ms (docs/voice.md §3)
// can be 2.5…60 ms, so the encoder needs frame_samples per call (480 @10ms, 1920 @40ms, // can be 2.5…60 ms, so the encoder needs frame_samples per call (480 @10ms, 1920 @40ms,
@@ -344,12 +357,27 @@ struct vc_client {
void handle_server_state(const voicecat::v1::ServerStateSnapshot& snap); void handle_server_state(const voicecat::v1::ServerStateSnapshot& snap);
void handle_user_event(const voicecat::v1::UserEvent& ue); void handle_user_event(const voicecat::v1::UserEvent& ue);
void handle_channel_event(const voicecat::v1::ChannelEvent& ce); void handle_channel_event(const voicecat::v1::ChannelEvent& ce);
// Called from handle_channel_event when the current channel's audio config changed:
// stop→start every active local stream so the new Opus params take effect (encoders are
// frozen at StreamAnnounceResult time — docs/voice.md §3). capture_device_id and
// capture_channels survive the restart (stream_stop doesn't clear them; stream_start
// reuses the existing LocalStream entry). The server reads the updated channel config
// on re-announce and returns new effective_audio; peers' sync_remote_streams wire up
// fresh decoders at the new ssrc.
void restart_active_streams_for_channel(uint32_t channel_id);
void handle_join_channel_result(const voicecat::v1::JoinChannelResult& msg); void handle_join_channel_result(const voicecat::v1::JoinChannelResult& msg);
void handle_text_message(const voicecat::v1::TextMessage& msg); void handle_text_message(const voicecat::v1::TextMessage& msg);
void handle_disconnect(const voicecat::v1::Disconnect& msg); void handle_disconnect(const voicecat::v1::Disconnect& msg);
void handle_udp_binding_ack(const voicecat::v1::UdpBinding& msg); void handle_udp_binding_ack(const voicecat::v1::UdpBinding& msg);
void handle_stream_announce_result(uint64_t req_id, void handle_stream_announce_result(uint64_t req_id,
const voicecat::v1::StreamAnnounceResult& msg); const voicecat::v1::StreamAnnounceResult& msg);
void handle_voice_subscription_result(const voicecat::v1::VoiceSubscriptionResult& msg);
// Wire up remote-stream decoders for all users in the session model (used on voice join).
void resync_remote_streams();
// Tear down all remote-stream decoders and clear remote_streams_ (used on voice leave).
void clear_all_remote_streams();
// Stop every active local stream (used on voice leave — emits STREAM_STOPPED for each).
void stop_all_local_streams();
// ── M2: UDP / media helpers ────────────────────────────────────────────────── // ── M2: UDP / media helpers ──────────────────────────────────────────────────
// Kicks off TCP UdpBinding request; called once after a successful AuthResult. // Kicks off TCP UdpBinding request; called once after a successful AuthResult.

View File

@@ -26,6 +26,20 @@ std::pair<const User*, const Stream*> SessionModel::find_user_by_ssrc(uint32_t s
#ifdef VOICECAT_HAS_NET #ifdef VOICECAT_HAS_NET
namespace { namespace {
void copy_channel_audio(Channel& ch, const voicecat::v1::AudioConfig& a) {
ch.audio_sample_rate = a.sample_rate() ? a.sample_rate() : 48000;
ch.audio_frame_ms = a.frame_ms() ? a.frame_ms() : 20;
ch.audio_mode = static_cast<uint32_t>(a.mode());
ch.audio_bitrate_bps = a.bitrate_bps();
ch.audio_application = static_cast<uint32_t>(a.application());
ch.audio_fec = a.fec();
ch.audio_expected_packet_loss = a.expected_packet_loss();
ch.audio_dtx = a.dtx();
ch.audio_complexity = a.complexity();
ch.audio_dred = a.dred();
}
std::vector<Stream> copy_streams( std::vector<Stream> copy_streams(
const google::protobuf::RepeatedPtrField<voicecat::v1::StreamInfo>& src) { const google::protobuf::RepeatedPtrField<voicecat::v1::StreamInfo>& src) {
std::vector<Stream> out; std::vector<Stream> out;
@@ -62,6 +76,8 @@ void SessionModel::apply_snapshot(const voicecat::v1::ServerStateSnapshot& snap)
ch.topic = pb.topic(); ch.topic = pb.topic();
ch.password_protected = pb.password_protected(); ch.password_protected = pb.password_protected();
ch.max_users = pb.max_users(); ch.max_users = pb.max_users();
ch.sort_order = static_cast<uint32_t>(pb.order());
copy_channel_audio(ch, pb.audio());
channels_.push_back(std::move(ch)); channels_.push_back(std::move(ch));
} }
@@ -76,6 +92,7 @@ void SessionModel::apply_snapshot(const voicecat::v1::ServerStateSnapshot& snap)
u.self_deafened = pb.self_deafened(); u.self_deafened = pb.self_deafened();
u.server_muted = pb.server_muted(); u.server_muted = pb.server_muted();
u.server_deafened = pb.server_deafened(); u.server_deafened = pb.server_deafened();
u.voice_subscribed = pb.voice_subscribed();
u.streams = copy_streams(pb.streams()); u.streams = copy_streams(pb.streams());
users_.push_back(std::move(u)); users_.push_back(std::move(u));
} }
@@ -95,6 +112,7 @@ void SessionModel::apply_user_event(const voicecat::v1::UserEvent& ev) {
u.self_deafened = pb.self_deafened(); u.self_deafened = pb.self_deafened();
u.server_muted = pb.server_muted(); u.server_muted = pb.server_muted();
u.server_deafened = pb.server_deafened(); u.server_deafened = pb.server_deafened();
u.voice_subscribed = pb.voice_subscribed();
u.streams = copy_streams(pb.streams()); u.streams = copy_streams(pb.streams());
auto it = std::find_if(users_.begin(), users_.end(), auto it = std::find_if(users_.begin(), users_.end(),
@@ -122,6 +140,8 @@ void SessionModel::apply_channel_event(const voicecat::v1::ChannelEvent& ev) {
ch.topic = pb.topic(); ch.topic = pb.topic();
ch.password_protected = pb.password_protected(); ch.password_protected = pb.password_protected();
ch.max_users = pb.max_users(); ch.max_users = pb.max_users();
ch.sort_order = static_cast<uint32_t>(pb.order());
copy_channel_audio(ch, pb.audio());
auto it = std::find_if(channels_.begin(), channels_.end(), auto it = std::find_if(channels_.begin(), channels_.end(),
[&](const Channel& x) { return x.id == ch.id; }); [&](const Channel& x) { return x.id == ch.id; });

View File

@@ -25,6 +25,21 @@ struct Channel {
std::string topic; std::string topic;
bool password_protected{false}; bool password_protected{false};
uint32_t max_users{0}; uint32_t max_users{0};
uint32_t sort_order{0};
// Authoritative channel AudioConfig (docs/voice.md §3) — mirrors
// voicecat::v1::AudioConfig field-for-field, same flat shape as Stream below. Populated
// from ServerStateSnapshot / ChannelEvent so the channel-edit dialog can read back the
// current config (the vc_channel read struct now carries these too).
uint32_t audio_sample_rate{48000};
uint32_t audio_frame_ms{20};
uint32_t audio_mode{0}; // 0 = mono, 1 = stereo
uint32_t audio_bitrate_bps{0};
uint32_t audio_application{0}; // 0=VOIP, 1=AUDIO, 2=LOWDELAY
bool audio_fec{false};
uint32_t audio_expected_packet_loss{0};
bool audio_dtx{false};
uint32_t audio_complexity{0};
bool audio_dred{false};
}; };
struct Stream { struct Stream {
@@ -57,6 +72,7 @@ struct User {
bool self_deafened{false}; bool self_deafened{false};
bool server_muted{false}; bool server_muted{false};
bool server_deafened{false}; bool server_deafened{false};
bool voice_subscribed{false};
std::vector<Stream> streams; std::vector<Stream> streams;
}; };

View File

@@ -82,6 +82,16 @@ vc_result vc_leave_channel(vc_client* c) {
return c->leave_channel(); return c->leave_channel();
} }
vc_result vc_join_voice(vc_client* c) {
if (c == nullptr) return VC_ERR_INVALID_ARG;
return c->join_voice();
}
vc_result vc_leave_voice(vc_client* c) {
if (c == nullptr) return VC_ERR_INVALID_ARG;
return c->leave_voice();
}
vc_result vc_stream_start(vc_client* c, const vc_stream_desc* desc, uint32_t* out_stream_id) { vc_result vc_stream_start(vc_client* c, const vc_stream_desc* desc, uint32_t* out_stream_id) {
if (c == nullptr || desc == nullptr) return VC_ERR_INVALID_ARG; if (c == nullptr || desc == nullptr) return VC_ERR_INVALID_ARG;
return c->stream_start(*desc, out_stream_id); return c->stream_start(*desc, out_stream_id);

View File

@@ -82,6 +82,9 @@ message Envelope {
StreamStop stream_stop = 42; StreamStop stream_stop = 42;
StreamStateUpdate stream_state = 43; // talking/muted indicator StreamStateUpdate stream_state = 43; // talking/muted indicator
UdpBinding udp_binding = 44; // token to bind the UDP 5-tuple UdpBinding udp_binding = 44; // token to bind the UDP 5-tuple
SubscribeVoiceRequest subscribe_voice = 45; // join the voice plane
UnsubscribeVoiceRequest unsubscribe_voice = 46; // leave the voice plane
VoiceSubscriptionResult voice_subscription_result = 47; // ack with subscribed flag
// ── Text ────────────────────────────────────────────── // ── Text ──────────────────────────────────────────────
TextMessage text_message = 50; TextMessage text_message = 50;

View File

@@ -111,6 +111,14 @@ void ConnSession::on_frame(std::vector<uint8_t> frame) {
if (st == State::Authenticated) if (st == State::Authenticated)
handle_stream_stop(env.stream_stop()); handle_stream_stop(env.stream_stop());
break; break;
case voicecat::v1::Envelope::kSubscribeVoice:
if (st == State::Authenticated)
handle_subscribe_voice(env.request_id());
break;
case voicecat::v1::Envelope::kUnsubscribeVoice:
if (st == State::Authenticated)
handle_unsubscribe_voice(env.request_id());
break;
// ── M5 moderation / admin ───────────────────────────────────────────── // ── M5 moderation / admin ─────────────────────────────────────────────
case voicecat::v1::Envelope::kKick: case voicecat::v1::Envelope::kKick:
@@ -521,6 +529,14 @@ void ConnSession::handle_udp_binding(uint64_t req_id, const voicecat::v1::UdpBin
void ConnSession::handle_stream_announce(uint64_t req_id, void ConnSession::handle_stream_announce(uint64_t req_id,
const voicecat::v1::StreamAnnounce& msg) { const voicecat::v1::StreamAnnounce& msg) {
if (!voice_subscribed_.load(std::memory_order_acquire)) {
auto env = make_env(req_id);
auto* res = env.mutable_stream_announce_result();
res->set_ok(false);
res->set_error("not subscribed to voice — call vc_join_voice first");
send_envelope(env);
return;
}
uint32_t ssrc = registry_->assign_ssrc(session_id_); uint32_t ssrc = registry_->assign_ssrc(session_id_);
uint32_t stream_id = next_stream_id_++; uint32_t stream_id = next_stream_id_++;
@@ -595,6 +611,66 @@ void ConnSession::handle_stream_stop(const voicecat::v1::StreamStop& msg) {
} }
} }
void ConnSession::handle_subscribe_voice(uint64_t req_id) {
voice_subscribed_.store(true, std::memory_order_release);
registry_->set_user_voice_subscribed(user_id_.load(), true);
// Reply with the result.
auto env = make_env(req_id);
auto* res = env.mutable_voice_subscription_result();
res->set_ok(true);
res->set_subscribed(true);
send_envelope(env);
// Broadcast the updated user proto so peers see voice_subscribed = true.
if (auto updated_user = registry_->user_snapshot_user(user_id_.load())) {
auto bcast = make_env();
auto* ue = bcast.mutable_user_event();
ue->set_kind(voicecat::v1::UserEvent::UPDATED);
*ue->mutable_user() = *updated_user;
registry_->broadcast(bcast, /*exclude*/ 0);
}
}
void ConnSession::handle_unsubscribe_voice(uint64_t req_id) {
voice_subscribed_.store(false, std::memory_order_release);
// Clear all the user's streams so peers stop receiving from them. Each clear broadcasts
// a UserEvent::UPDATED; we collect the ids first to avoid mutating announced_stream_ids_
// while iterating.
auto stream_ids = announced_stream_ids_;
for (uint32_t sid : stream_ids) {
auto it = std::find(announced_stream_ids_.begin(), announced_stream_ids_.end(), sid);
if (it != announced_stream_ids_.end()) announced_stream_ids_.erase(it);
auto updated = registry_->clear_user_stream(user_id_.load(), sid);
if (updated) {
auto bcast = make_env();
auto* ue = bcast.mutable_user_event();
ue->set_kind(voicecat::v1::UserEvent::UPDATED);
*ue->mutable_user() = *updated;
registry_->broadcast(bcast, /*exclude*/ 0);
}
}
registry_->set_user_voice_subscribed(user_id_.load(), false);
// Reply with the result.
auto env = make_env(req_id);
auto* res = env.mutable_voice_subscription_result();
res->set_ok(true);
res->set_subscribed(false);
send_envelope(env);
// Broadcast the updated user proto so peers see voice_subscribed = false.
if (auto updated_user = registry_->user_snapshot_user(user_id_.load())) {
auto bcast = make_env();
auto* ue = bcast.mutable_user_event();
ue->set_kind(voicecat::v1::UserEvent::UPDATED);
*ue->mutable_user() = *updated_user;
registry_->broadcast(bcast, /*exclude*/ 0);
}
}
// ── M5 handlers ────────────────────────────────────────────────────────────── // ── M5 handlers ──────────────────────────────────────────────────────────────
void ConnSession::handle_kick_request(uint64_t req_id, const voicecat::v1::KickRequest& msg) { void ConnSession::handle_kick_request(uint64_t req_id, const voicecat::v1::KickRequest& msg) {

View File

@@ -82,6 +82,7 @@ class ConnSession : public std::enable_shared_from_this<ConnSession> {
State state() const { return state_.load(); } State state() const { return state_.load(); }
uint64_t session_id() const { return session_id_; } uint64_t session_id() const { return session_id_; }
uint32_t user_id() const { return user_id_; } uint32_t user_id() const { return user_id_; }
bool voice_subscribed() const { return voice_subscribed_.load(std::memory_order_acquire); }
// Keepalive: bump last_seen to "now" on any inbound activity (TCP frame or UDP voice // Keepalive: bump last_seen to "now" on any inbound activity (TCP frame or UDP voice
// frame). The server's reaper sweeps sessions whose last_seen is older than the // frame). The server's reaper sweeps sessions whose last_seen is older than the
@@ -105,6 +106,8 @@ class ConnSession : public std::enable_shared_from_this<ConnSession> {
void handle_stream_announce(uint64_t req_id, const voicecat::v1::StreamAnnounce& msg); void handle_stream_announce(uint64_t req_id, const voicecat::v1::StreamAnnounce& msg);
void handle_stream_stop(const voicecat::v1::StreamStop& msg); void handle_stream_stop(const voicecat::v1::StreamStop& msg);
void handle_leave_channel(); void handle_leave_channel();
void handle_subscribe_voice(uint64_t req_id);
void handle_unsubscribe_voice(uint64_t req_id);
// M5 handlers // M5 handlers
void handle_kick_request(uint64_t req_id, const voicecat::v1::KickRequest& msg); void handle_kick_request(uint64_t req_id, const voicecat::v1::KickRequest& msg);
@@ -172,6 +175,11 @@ class ConnSession : public std::enable_shared_from_this<ConnSession> {
// M5: permissions granted at auth time (server-side authority). // M5: permissions granted at auth time (server-side authority).
voicecat::v1::Permissions permissions_; voicecat::v1::Permissions permissions_;
// Voice-plane subscription. When false, the SFU relay excludes this session from the
// recipient set (find_channel_sessions), and handle_stream_announce rejects new streams.
// Toggled by SubscribeVoiceRequest / UnsubscribeVoiceRequest.
std::atomic<bool> voice_subscribed_{false};
}; };
} // namespace voicecat::server } // namespace voicecat::server

View File

@@ -81,6 +81,12 @@ void MediaRelay::on_udp_frame(const uint8_t* data, size_t len,
auto sender_session = registry_->find_by_udp_endpoint(sender); auto sender_session = registry_->find_by_udp_endpoint(sender);
if (!sender_session) { note_drop("unmapped-endpoint"); return; } if (!sender_session) { note_drop("unmapped-endpoint"); return; }
// Defense in depth: don't relay voice from sessions not on the voice plane.
// handle_stream_announce already rejects announces from non-subscribers, so a
// non-subscribed session shouldn't have any active streams — but this catches any
// edge case (e.g. a stream announced just before unsubscribing).
if (!sender_session->voice_subscribed()) { note_drop("not-subscribed"); return; }
auto* recv_crypto = sender_session->recv_crypto(); auto* recv_crypto = sender_session->recv_crypto();
if (!recv_crypto) { note_drop("no-recv-crypto"); return; } if (!recv_crypto) { note_drop("no-recv-crypto"); return; }

View File

@@ -145,6 +145,13 @@ bool SessionRegistry::set_user_channel(uint32_t user_id, uint32_t channel_id) {
return true; return true;
} }
void SessionRegistry::set_user_voice_subscribed(uint32_t user_id, bool subscribed) {
std::unique_lock lk(mu_);
auto it = users_.find(user_id);
if (it == users_.end()) return;
it->second.proto.set_voice_subscribed(subscribed);
}
std::vector<voicecat::v1::Channel> SessionRegistry::channel_snapshot() const { std::vector<voicecat::v1::Channel> SessionRegistry::channel_snapshot() const {
std::shared_lock lk(mu_); std::shared_lock lk(mu_);
std::vector<voicecat::v1::Channel> result; std::vector<voicecat::v1::Channel> result;
@@ -492,7 +499,11 @@ std::vector<std::shared_ptr<ConnSession>> SessionRegistry::find_channel_sessions
if (entry.session_id == exclude_session_id) continue; if (entry.session_id == exclude_session_id) continue;
auto sit = sessions_.find(entry.session_id); auto sit = sessions_.find(entry.session_id);
if (sit == sessions_.end()) continue; if (sit == sessions_.end()) continue;
if (auto sess = sit->second.lock()) result.push_back(sess); if (auto sess = sit->second.lock()) {
// Only relay voice to sessions that are on the voice plane.
if (!sess->voice_subscribed()) continue;
result.push_back(sess);
}
} }
return result; return result;
} }

View File

@@ -75,6 +75,9 @@ class SessionRegistry {
// Move a user to a channel. Returns false if channel doesn't exist. // Move a user to a channel. Returns false if channel doesn't exist.
bool set_user_channel(uint32_t user_id, uint32_t channel_id); bool set_user_channel(uint32_t user_id, uint32_t channel_id);
// Set the user's voice-plane subscription flag on their proto (broadcast-ready).
void set_user_voice_subscribed(uint32_t user_id, bool subscribed);
// Snapshot for ServerStateSnapshot message. // Snapshot for ServerStateSnapshot message.
std::vector<voicecat::v1::Channel> channel_snapshot() const; std::vector<voicecat::v1::Channel> channel_snapshot() const;
std::vector<voicecat::v1::User> user_snapshot() const; std::vector<voicecat::v1::User> user_snapshot() const;

View File

@@ -46,6 +46,7 @@ struct EventStore {
bool channel_list_received{false}; bool channel_list_received{false};
bool generic_result_received{false}; bool generic_result_received{false};
bool generic_ok{false}; bool generic_ok{false};
bool voice_subscribed{false};
std::vector<std::pair<uint32_t, uint32_t>> streams_started; // (user_id, stream_id) std::vector<std::pair<uint32_t, uint32_t>> streams_started; // (user_id, stream_id)
vc_client* client{nullptr}; vc_client* client{nullptr};
const char* label{nullptr}; const char* label{nullptr};
@@ -61,6 +62,7 @@ static void on_event(void* user, const vc_event* ev) {
s->self_user_id = ev->user_id; s->self_user_id = ev->user_id;
break; break;
case VC_EVENT_CHANNEL_LIST: s->channel_list_received = true; break; case VC_EVENT_CHANNEL_LIST: s->channel_list_received = true; break;
case VC_EVENT_VOICE_STATE: s->voice_subscribed = (ev->u32a == 1); break;
case VC_EVENT_GENERIC_RESULT: case VC_EVENT_GENERIC_RESULT:
s->generic_result_received = true; s->generic_result_received = true;
s->generic_ok = (ev->result == VC_OK); s->generic_ok = (ev->result == VC_OK);
@@ -90,6 +92,8 @@ static bool connect_guest(vc_client*& client, const char* name, const char* labe
if (vc_authenticate_guest(client, name) != VC_OK) return false; if (vc_authenticate_guest(client, name) != VC_OK) return false;
if (!wait_for(ev, [](EventStore& s) { return s.auth_ok; }, 8000)) return false; if (!wait_for(ev, [](EventStore& s) { return s.auth_ok; }, 8000)) return false;
if (!wait_for(ev, [](EventStore& s) { return s.channel_list_received; }, 3000)) return false; if (!wait_for(ev, [](EventStore& s) { return s.channel_list_received; }, 3000)) return false;
if (vc_join_voice(client) != VC_OK) return false;
if (!wait_for(ev, [](EventStore& s) { return s.voice_subscribed; }, 5000)) return false;
return true; return true;
} }

View File

@@ -59,6 +59,7 @@ struct EventStore {
bool channel_list_received{false}; bool channel_list_received{false};
std::vector<StreamEvent> stream_events; std::vector<StreamEvent> stream_events;
std::vector<JoinResult> join_results; std::vector<JoinResult> join_results;
bool voice_subscribed{false};
const char* label{nullptr}; const char* label{nullptr};
vc_client* client{nullptr}; vc_client* client{nullptr};
@@ -78,6 +79,9 @@ static void on_event(void* user, const vc_event* ev) {
case VC_EVENT_CHANNEL_LIST: case VC_EVENT_CHANNEL_LIST:
s->channel_list_received = true; s->channel_list_received = true;
break; break;
case VC_EVENT_VOICE_STATE:
s->voice_subscribed = (ev->u32a == 1);
break;
case VC_EVENT_STREAM_STARTED: case VC_EVENT_STREAM_STARTED:
s->stream_events.push_back({true, ev->user_id, ev->stream_id}); s->stream_events.push_back({true, ev->user_id, ev->stream_id});
break; break;
@@ -195,6 +199,9 @@ static void test_live_channel_user_list() {
CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000)); CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000));
CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000)); CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(clientA) == VC_OK);
CHECK(wait_for(evA, [](EventStore& s) { return s.voice_subscribed; }, 5000));
EventStore evB; EventStore evB;
evB.label = "B"; evB.label = "B";
vc_callbacks cbB{on_event, nullptr, &evB}; vc_callbacks cbB{on_event, nullptr, &evB};
@@ -207,6 +214,9 @@ static void test_live_channel_user_list() {
CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000)); CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000));
CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000)); CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(clientB) == VC_OK);
CHECK(wait_for(evB, [](EventStore& s) { return s.voice_subscribed; }, 5000));
uint32_t a_uid = 0, b_uid = 0; uint32_t a_uid = 0, b_uid = 0;
{ std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; } { std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; }
{ std::lock_guard lk(evB.mu); b_uid = evB.self_user_id; } { std::lock_guard lk(evB.mu); b_uid = evB.self_user_id; }

View File

@@ -35,6 +35,7 @@ struct EventStore {
bool channel_list_received{false}; bool channel_list_received{false};
bool generic_result_received{false}; bool generic_result_received{false};
bool generic_ok{false}; bool generic_ok{false};
bool voice_subscribed{false};
std::vector<std::pair<uint32_t,uint32_t>> streams_started; // (user_id, stream_id) std::vector<std::pair<uint32_t,uint32_t>> streams_started; // (user_id, stream_id)
vc_client* client{nullptr}; vc_client* client{nullptr};
const char* label{nullptr}; const char* label{nullptr};
@@ -54,6 +55,9 @@ static void on_event(void* user, const vc_event* ev) {
case VC_EVENT_CHANNEL_LIST: case VC_EVENT_CHANNEL_LIST:
s->channel_list_received = true; s->channel_list_received = true;
break; break;
case VC_EVENT_VOICE_STATE:
s->voice_subscribed = (ev->u32a == 1);
break;
case VC_EVENT_GENERIC_RESULT: case VC_EVENT_GENERIC_RESULT:
s->generic_result_received = true; s->generic_result_received = true;
s->generic_ok = (ev->result == VC_OK); s->generic_ok = (ev->result == VC_OK);
@@ -199,6 +203,9 @@ int main() {
CHECK(wait_for(evGuest, [](EventStore& s){ return s.auth_ok; }, 8000)); CHECK(wait_for(evGuest, [](EventStore& s){ return s.auth_ok; }, 8000));
CHECK(wait_for(evGuest, [](EventStore& s){ return s.channel_list_received; }, 3000)); CHECK(wait_for(evGuest, [](EventStore& s){ return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(guest) == VC_OK);
CHECK(wait_for(evGuest, [](EventStore& s){ return s.voice_subscribed; }, 5000));
// Move the guest into the DRED channel. // Move the guest into the DRED channel.
uint32_t guest_uid = 0; uint32_t guest_uid = 0;
{ std::lock_guard lk(evGuest.mu); guest_uid = evGuest.self_user_id; } { std::lock_guard lk(evGuest.mu); guest_uid = evGuest.self_user_id; }

View File

@@ -47,6 +47,7 @@ struct EventStore {
uint32_t self_user_id{0}; uint32_t self_user_id{0};
bool channel_list_received{false}; bool channel_list_received{false};
bool join_ok{false}; bool join_ok{false};
bool voice_subscribed{false};
std::vector<std::pair<uint32_t, uint32_t>> streams_started; // (user_id, stream_id) std::vector<std::pair<uint32_t, uint32_t>> streams_started; // (user_id, stream_id)
vc_client* client{nullptr}; vc_client* client{nullptr};
const char* label{nullptr}; const char* label{nullptr};
@@ -66,6 +67,9 @@ static void on_event(void* user, const vc_event* ev) {
case VC_EVENT_CHANNEL_LIST: case VC_EVENT_CHANNEL_LIST:
s->channel_list_received = true; s->channel_list_received = true;
break; break;
case VC_EVENT_VOICE_STATE:
s->voice_subscribed = (ev->u32a == 1);
break;
case VC_EVENT_JOIN_RESULT: case VC_EVENT_JOIN_RESULT:
s->join_ok = (ev->result == VC_OK); s->join_ok = (ev->result == VC_OK);
break; break;
@@ -96,6 +100,8 @@ static bool connect_guest(vc_client*& client, const char* name, const char* labe
if (vc_authenticate_guest(client, name) != VC_OK) return false; if (vc_authenticate_guest(client, name) != VC_OK) return false;
if (!wait_for(ev, [](EventStore& s) { return s.auth_ok; }, 8000)) return false; if (!wait_for(ev, [](EventStore& s) { return s.auth_ok; }, 8000)) return false;
if (!wait_for(ev, [](EventStore& s) { return s.channel_list_received; }, 3000)) return false; if (!wait_for(ev, [](EventStore& s) { return s.channel_list_received; }, 3000)) return false;
if (vc_join_voice(client) != VC_OK) return false;
if (!wait_for(ev, [](EventStore& s) { return s.voice_subscribed; }, 5000)) return false;
return true; return true;
} }

View File

@@ -51,6 +51,7 @@ struct EventStore {
bool channel_list_received{false}; bool channel_list_received{false};
bool generic_result_received{false}; bool generic_result_received{false};
bool generic_ok{false}; bool generic_ok{false};
bool voice_subscribed{false};
std::vector<std::pair<uint32_t, uint32_t>> streams_started; // (user_id, stream_id) std::vector<std::pair<uint32_t, uint32_t>> streams_started; // (user_id, stream_id)
vc_client* client{nullptr}; vc_client* client{nullptr};
const char* label{nullptr}; const char* label{nullptr};
@@ -66,6 +67,7 @@ static void on_event(void* user, const vc_event* ev) {
s->self_user_id = ev->user_id; s->self_user_id = ev->user_id;
break; break;
case VC_EVENT_CHANNEL_LIST: s->channel_list_received = true; break; case VC_EVENT_CHANNEL_LIST: s->channel_list_received = true; break;
case VC_EVENT_VOICE_STATE: s->voice_subscribed = (ev->u32a == 1); break;
case VC_EVENT_GENERIC_RESULT: case VC_EVENT_GENERIC_RESULT:
s->generic_result_received = true; s->generic_result_received = true;
s->generic_ok = (ev->result == VC_OK); s->generic_ok = (ev->result == VC_OK);
@@ -95,6 +97,8 @@ static bool connect_guest(vc_client*& client, const char* name, const char* labe
if (vc_authenticate_guest(client, name) != VC_OK) return false; if (vc_authenticate_guest(client, name) != VC_OK) return false;
if (!wait_for(ev, [](EventStore& s) { return s.auth_ok; }, 8000)) return false; if (!wait_for(ev, [](EventStore& s) { return s.auth_ok; }, 8000)) return false;
if (!wait_for(ev, [](EventStore& s) { return s.channel_list_received; }, 3000)) return false; if (!wait_for(ev, [](EventStore& s) { return s.channel_list_received; }, 3000)) return false;
if (vc_join_voice(client) != VC_OK) return false;
if (!wait_for(ev, [](EventStore& s) { return s.voice_subscribed; }, 5000)) return false;
return true; return true;
} }

View File

@@ -291,6 +291,23 @@ struct TestClient {
return true; return true;
} }
bool do_subscribe_voice() {
{
v1::Envelope env;
env.set_request_id(3);
env.mutable_subscribe_voice();
if (!tcp_send_envelope(*tls, env)) return false;
}
for (int i = 0; i < 20; ++i) {
v1::Envelope env;
if (!tcp_recv_envelope(*tls, codec, env)) return false;
if (env.has_voice_subscription_result()) {
return env.voice_subscription_result().subscribed();
}
}
return false;
}
bool do_stream_announce() { bool do_stream_announce() {
// Send StreamAnnounce. // Send StreamAnnounce.
{ {
@@ -396,6 +413,7 @@ int main() {
CHECK(A.do_hello_and_auth("SenderBob")); CHECK(A.do_hello_and_auth("SenderBob"));
CHECK(A.server_udp_port == udp_port.load()); CHECK(A.server_udp_port == udp_port.load());
CHECK(A.do_udp_binding()); CHECK(A.do_udp_binding());
CHECK(A.do_subscribe_voice());
CHECK(A.do_stream_announce()); CHECK(A.do_stream_announce());
CHECK(A.stream_ok); CHECK(A.stream_ok);
std::printf("m2_voice: A ssrc=%u local_udp=%u\n", A.assigned_ssrc, A.udp_local_port); std::printf("m2_voice: A ssrc=%u local_udp=%u\n", A.assigned_ssrc, A.udp_local_port);
@@ -408,6 +426,7 @@ int main() {
CHECK(B.do_hello_and_auth("ReceiverAlice")); CHECK(B.do_hello_and_auth("ReceiverAlice"));
CHECK(B.server_udp_port == udp_port.load()); CHECK(B.server_udp_port == udp_port.load());
CHECK(B.do_udp_binding()); CHECK(B.do_udp_binding());
CHECK(B.do_subscribe_voice());
// B doesn't need to announce a stream to receive relayed frames // B doesn't need to announce a stream to receive relayed frames
if (g_failures > 0) { if (g_failures > 0) {

View File

@@ -59,6 +59,7 @@ struct EventStore {
bool channel_list_received{false}; bool channel_list_received{false};
std::vector<StreamEvent> stream_events; std::vector<StreamEvent> stream_events;
std::vector<TalkEvent> talk_events; std::vector<TalkEvent> talk_events;
bool voice_subscribed{false};
bool disconnected{false}; bool disconnected{false};
const char* label{nullptr}; const char* label{nullptr};
@@ -85,6 +86,9 @@ static void on_event(void* user, const vc_event* ev) {
case VC_EVENT_CHANNEL_LIST: case VC_EVENT_CHANNEL_LIST:
s->channel_list_received = true; s->channel_list_received = true;
break; break;
case VC_EVENT_VOICE_STATE:
s->voice_subscribed = (ev->u32a == 1);
break;
case VC_EVENT_STREAM_STARTED: case VC_EVENT_STREAM_STARTED:
s->stream_events.push_back({true, ev->user_id, ev->stream_id}); s->stream_events.push_back({true, ev->user_id, ev->stream_id});
break; break;
@@ -191,6 +195,9 @@ int main() {
CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000)); CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000));
CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000)); CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(clientA) == VC_OK);
CHECK(wait_for(evA, [](EventStore& s) { return s.voice_subscribed; }, 5000));
// ── Client B: guest "M3-B" ──────────────────────────────────────────────── // ── Client B: guest "M3-B" ────────────────────────────────────────────────
EventStore evB; EventStore evB;
evB.label = "clientB"; evB.label = "clientB";
@@ -205,6 +212,9 @@ int main() {
CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000)); CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000));
CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000)); CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(clientB) == VC_OK);
CHECK(wait_for(evB, [](EventStore& s) { return s.voice_subscribed; }, 5000));
uint32_t a_uid = 0; uint32_t a_uid = 0;
{ std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; } { std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; }

View File

@@ -59,6 +59,7 @@ struct EventStore {
bool channel_list_received{false}; bool channel_list_received{false};
bool saw_stream_started{false}; bool saw_stream_started{false};
std::vector<TalkEvent> talk_events; std::vector<TalkEvent> talk_events;
bool voice_subscribed{false};
bool disconnected{false}; bool disconnected{false};
const char* label{nullptr}; const char* label{nullptr};
@@ -83,6 +84,9 @@ static void on_event(void* user, const vc_event* ev) {
case VC_EVENT_CHANNEL_LIST: case VC_EVENT_CHANNEL_LIST:
s->channel_list_received = true; s->channel_list_received = true;
break; break;
case VC_EVENT_VOICE_STATE:
s->voice_subscribed = (ev->u32a == 1);
break;
case VC_EVENT_STREAM_STARTED: case VC_EVENT_STREAM_STARTED:
s->saw_stream_started = true; s->saw_stream_started = true;
break; break;
@@ -490,6 +494,9 @@ static void test_vad_and_ptt_gate() {
CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000)); CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000));
CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000)); CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(clientA) == VC_OK);
CHECK(wait_for(evA, [](EventStore& s) { return s.voice_subscribed; }, 5000));
EventStore evB; EventStore evB;
evB.label = "B"; evB.label = "B";
vc_callbacks cbB{on_event, nullptr, &evB}; vc_callbacks cbB{on_event, nullptr, &evB};
@@ -502,6 +509,9 @@ static void test_vad_and_ptt_gate() {
CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000)); CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000));
CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000)); CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(clientB) == VC_OK);
CHECK(wait_for(evB, [](EventStore& s) { return s.voice_subscribed; }, 5000));
uint32_t a_uid = 0; uint32_t a_uid = 0;
{ std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; } { std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; }

View File

@@ -40,6 +40,7 @@ struct EventStore {
uint32_t self_user_id{0}; uint32_t self_user_id{0};
bool channel_list_received{false}; bool channel_list_received{false};
std::vector<StreamEvent> stream_events; std::vector<StreamEvent> stream_events;
bool voice_subscribed{false};
const char* label{nullptr}; const char* label{nullptr};
bool disconnected{false}; bool disconnected{false};
@@ -66,6 +67,9 @@ static void on_event(void* user, const vc_event* ev) {
case VC_EVENT_CHANNEL_LIST: case VC_EVENT_CHANNEL_LIST:
s->channel_list_received = true; s->channel_list_received = true;
break; break;
case VC_EVENT_VOICE_STATE:
s->voice_subscribed = (ev->u32a == 1);
break;
case VC_EVENT_STREAM_STARTED: case VC_EVENT_STREAM_STARTED:
s->stream_events.push_back({true, ev->user_id, ev->stream_id}); s->stream_events.push_back({true, ev->user_id, ev->stream_id});
break; break;
@@ -159,6 +163,9 @@ int main() {
CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000)); CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000));
CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000)); CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(clientA) == VC_OK);
CHECK(wait_for(evA, [](EventStore& s) { return s.voice_subscribed; }, 5000));
// ── Client B: guest "VoiceB" ─────────────────────────────────────────────── // ── Client B: guest "VoiceB" ───────────────────────────────────────────────
EventStore evB; EventStore evB;
evB.label = "clientB"; evB.label = "clientB";
@@ -173,6 +180,9 @@ int main() {
CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000)); CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000));
CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000)); CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000));
CHECK(vc_join_voice(clientB) == VC_OK);
CHECK(wait_for(evB, [](EventStore& s) { return s.voice_subscribed; }, 5000));
uint32_t a_uid = 0; uint32_t a_uid = 0;
{ std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; } { std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; }