Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
//
// Copyright 漏 2025 Stream.io Inc. All rights reserved.
//

import Foundation
import StreamWebRTC

public protocol AudioProcessingModule: RTCAudioProcessingModule, Sendable {

/// The currently active audio filter.
var activeAudioFilter: AudioFilter? { get }

/// Sets the audio filter to be used for audio processing.
/// - Parameter filter: The audio filter to set.
func setAudioFilter(_ filter: AudioFilter?)
}

final class AudioProcessingStore: RTCDefaultAudioProcessingModule, AudioProcessingModule, @unchecked Sendable {

private let store: Store<Namespace>

init() {
let initialState = Namespace.State.initial
store = Namespace.store(initialState: initialState)

super.init(
config: nil,
capturePostProcessingDelegate: initialState.capturePostProcessingDelegate,
renderPreProcessingDelegate: nil
)

store.dispatch(.load)
}

var activeAudioFilter: AudioFilter? { store.state.audioFilter }

func setAudioFilter(_ filter: AudioFilter?) {
store.dispatch(.setAudioFilter(filter))
}
}

extension AudioProcessingStore: InjectionKey {

nonisolated(unsafe) static var currentValue: AudioProcessingModule = AudioProcessingStore()
}

extension InjectedValues {
var audioFilterProcessingModule: AudioProcessingModule {
get { Self[AudioProcessingStore.self] }
set { Self[AudioProcessingStore.self] = newValue }
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
//
// Copyright 漏 2025 Stream.io Inc. All rights reserved.
//

import Combine
import Foundation
import StreamWebRTC

final class AudioCustomProcessingModule: NSObject, RTCAudioCustomProcessingDelegate, @unchecked Sendable {

enum Event {
case audioProcessingInitialize(sampleRateHz: Int, channels: Int)
case audioProcessingProcess(RTCAudioBuffer)
case audioProcessingRelease
}

private let subject: PassthroughSubject<Event, Never> = .init()
var publisher: AnyPublisher<Event, Never> { subject.eraseToAnyPublisher() }

func audioProcessingInitialize(
sampleRate sampleRateHz: Int,
channels: Int
) {
subject.send(
.audioProcessingInitialize(
sampleRateHz: sampleRateHz,
channels: channels
)
)
}

func audioProcessingProcess(audioBuffer: RTCAudioBuffer) {
subject.send(.audioProcessingProcess(audioBuffer))
}

func audioProcessingRelease() {
subject.send(.audioProcessingRelease)
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
//
// Copyright 漏 2025 Stream.io Inc. All rights reserved.
//

import Foundation

extension AudioProcessingStore.Namespace {

enum StoreAction: StoreActionBoxProtocol, Sendable {
case load
case setInitializedConfiguration(sampleRate: Int, channels: Int)
case setAudioFilter(AudioFilter?)
case setNumberOfCaptureChannels(Int)
case release
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
//
// Copyright 漏 2025 Stream.io Inc. All rights reserved.
//

import Foundation

extension AudioProcessingStore {

enum Namespace: StoreNamespace {
typealias State = StoreState

typealias Action = StoreAction

static let identifier: String = "io.getstream.audio.processing.store"

static func reducers() -> [Reducer<AudioProcessingStore.Namespace>] {
[
DefaultReducer()
]
}

static func middleware() -> [Middleware<AudioProcessingStore.Namespace>] {
[
CapturedChannelsMiddleware(),
AudioFilterMiddleware()
]
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
//
// Copyright 漏 2025 Stream.io Inc. All rights reserved.
//

import Foundation
import StreamWebRTC

extension AudioProcessingStore.Namespace {

struct StoreState: Equatable, Sendable {
var initializedSampleRate: Int
var initializedChannels: Int
var numberOfCaptureChannels: Int
var capturePostProcessingDelegate: AudioCustomProcessingModule
var audioFilter: AudioFilter?

static let initial = StoreState(
initializedSampleRate: 0,
initializedChannels: 0,
numberOfCaptureChannels: 0,
capturePostProcessingDelegate: .init(),
audioFilter: nil
)

static func == (
lhs: AudioProcessingStore.Namespace.StoreState,
rhs: AudioProcessingStore.Namespace.StoreState
) -> Bool {
lhs.numberOfCaptureChannels == rhs.numberOfCaptureChannels
&& lhs.capturePostProcessingDelegate === rhs.capturePostProcessingDelegate
&& lhs.audioFilter?.id == rhs.audioFilter?.id
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
//
// Copyright 漏 2025 Stream.io Inc. All rights reserved.
//

import Combine
import Foundation
import StreamWebRTC

extension AudioProcessingStore.Namespace {

final class AudioFilterMiddleware: Middleware<AudioProcessingStore.Namespace>, @unchecked Sendable {

private var cancellable: AnyCancellable?

override func apply(
state: AudioProcessingStore.Namespace.StoreState,
action: AudioProcessingStore.Namespace.StoreAction,
file: StaticString,
function: StaticString,
line: UInt
) {
switch action {
case let .setInitializedConfiguration(sampleRate, channels):
if let audioFilter = state.audioFilter {
audioFilter.initialize(
sampleRate: sampleRate,
channels: channels
)
}
case let .setAudioFilter(audioFilter):
if state.initializedSampleRate > 0, state.initializedChannels > 0 {
audioFilter?.initialize(
sampleRate: state.initializedSampleRate,
channels: state.initializedChannels
)
}
didUpdate(
audioFilter,
capturePostProcessingDelegate: state.capturePostProcessingDelegate
)
default:
break
}
}

// MARK: - Private Helpers

private func didUpdate(
_ audioFilter: AudioFilter?,
capturePostProcessingDelegate: AudioCustomProcessingModule
) {
cancellable?.cancel()
cancellable = nil

guard let audioFilter else {
return
}

cancellable = capturePostProcessingDelegate
.publisher
.compactMap {
guard case let .audioProcessingProcess(buffer) = $0 else {
return nil
}
return buffer
}
.sink { [weak self, audioFilter] in self?.process($0, on: audioFilter) }
}

private func process(
_ audioBuffer: RTCAudioBuffer,
on audioFilter: AudioFilter
) {
var audioBuffer = audioBuffer
audioFilter.applyEffect(to: &audioBuffer)
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
//
// Copyright 漏 2025 Stream.io Inc. All rights reserved.
//

import Combine
import Foundation

extension AudioProcessingStore.Namespace {

final class CapturedChannelsMiddleware: Middleware<AudioProcessingStore.Namespace>, @unchecked Sendable {

private var cancellable: AnyCancellable?

override func apply(
state: AudioProcessingStore.Namespace.StoreState,
action: AudioProcessingStore.Namespace.StoreAction,
file: StaticString,
function: StaticString,
line: UInt
) {
switch action {
case .load:
cancellable = state
.capturePostProcessingDelegate
.publisher
.sink { [weak self] in self?.didReceiveProcessingEvent($0) }

default:
break
}
}

// MARK: - Private Helpers

private func didReceiveProcessingEvent(
_ event: AudioCustomProcessingModule.Event
) {
switch event {
case let .audioProcessingInitialize(sampleRateHz, channels):
dispatcher?.dispatch(
.setInitializedConfiguration(
sampleRate: sampleRateHz,
channels: channels
)
)
case let .audioProcessingProcess(buffer):
if buffer.channels != stateProvider?()?.numberOfCaptureChannels {
dispatcher?.dispatch(.setNumberOfCaptureChannels(buffer.channels))
}
case .audioProcessingRelease:
dispatcher?.dispatch(.release)
}
}
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
//
// Copyright 漏 2025 Stream.io Inc. All rights reserved.
//

import Foundation

extension AudioProcessingStore.Namespace {

final class DefaultReducer: Reducer<AudioProcessingStore.Namespace>, @unchecked Sendable {

override func reduce(
state: AudioProcessingStore.Namespace.StoreState,
action: AudioProcessingStore.Namespace.StoreAction,
file: StaticString,
function: StaticString,
line: UInt
) throws -> AudioProcessingStore.Namespace.StoreState {
var updatedState = state

switch action {
case .load:
break

case let .setInitializedConfiguration(sampleRate, channels):
updatedState.initializedSampleRate = sampleRate
updatedState.initializedChannels = channels

case let .setAudioFilter(value):
updatedState.audioFilter = value

case let .setNumberOfCaptureChannels(value):
updatedState.numberOfCaptureChannels = value

case .release:
updatedState.initializedSampleRate = 0
updatedState.initializedChannels = 0
}

return updatedState
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,6 @@ extension AnyCancellable: @retroactive @unchecked Sendable {}
extension AVCaptureDevice: @retroactive @unchecked Sendable {}
extension AVCapturePhotoOutput: @retroactive @unchecked Sendable {}
extension AVCaptureVideoDataOutput: @retroactive @unchecked Sendable {}
extension CIImage: @retroactive @unchecked Sendable {}
extension CMSampleBuffer: @retroactive @unchecked Sendable {}
extension CXAnswerCallAction: @retroactive @unchecked Sendable {}
extension CXSetHeldCallAction: @retroactive @unchecked Sendable {}
Expand All @@ -34,7 +33,6 @@ extension AnyCancellable: @unchecked Sendable {}
extension AVCaptureDevice: @unchecked Sendable {}
extension AVCapturePhotoOutput: @unchecked Sendable {}
extension AVCaptureVideoDataOutput: @unchecked Sendable {}
extension CIImage: @unchecked Sendable {}
extension CMSampleBuffer: @unchecked Sendable {}
extension CXAnswerCallAction: @unchecked Sendable {}
extension CXSetHeldCallAction: @unchecked Sendable {}
Expand Down
Loading