スキル一覧に戻る
co-labs-co

combine-video-frame-distribution

by co-labs-co

Central skills repository for ContextHarness - curated agent skills and workflows

0🍴 0📅 2026年1月12日
GitHubで見るManusで実行

SKILL.md


name: combine-video-frame-distribution version: 0.1.0 author: cmtzco description: Thread-safe distribution of video frames (CVPixelBuffer) to multiple subscribers using Apple's Combine framework. Use this skill for real-time video streaming to multiple consumers, multi-subscriber broadcasting, cross-thread frame delivery, or replacing timer-based polling with push-based delivery.

Skill: Combine Video Frame Distribution

Efficient, thread-safe distribution of video frames (CVPixelBuffer) from a camera capture pipeline to multiple concurrent subscribers using Apple's Combine framework.

When to Use

  • Real-time video streaming: Camera feed to multiple consumers (preview, encoder, analytics)
  • Multi-subscriber broadcasting: Same frame data needed by different components
  • Cross-thread delivery: Camera background queue to main thread UI updates
  • Event-driven architecture: Replace timer-based polling with push-based frame delivery
  • SwiftUI integration: Bridging camera callbacks to @Published properties

Key Concepts

PassthroughSubject vs CurrentValueSubject

AspectPassthroughSubjectCurrentValueSubject
StorageNone - immediate forwardingHolds current value
New SubscribersReceive nothing until next sendImmediately receive current value
MemoryMinimal overheadStores one value
Video Use Case✅ Correct choice❌ Wasteful for transient frames

Why PassthroughSubject for video: Frames are transient (30-60 fps), new subscribers don't need stale frames.

Thread Safety Model

PassthroughSubject is thread-safe for send(), but callbacks execute on the sender's thread:

// Subscriber receives on cameraQueue (the sending thread)
subject.sink { buffer in 
    // ⚠️ This runs on cameraQueue, NOT main thread
}

// Use receive(on:) for thread control
subject
    .receive(on: DispatchQueue.main)
    .sink { buffer in
        // ✅ This runs on main thread
    }

Backpressure Strategy

Video frames are "hot" data—they arrive regardless of subscriber readiness:

StrategyVideo Suitability
No buffering (drop if busy)✅ Best for real-time
Buffer with limit⚠️ Adds latency
Unlimited buffer❌ Will crash

Implementation Guide

Step 1: Define the Frame Publisher

import Combine
import AVFoundation

@MainActor
class CameraManager: NSObject, ObservableObject {
    // Thread-safe subject for frame broadcasting
    nonisolated(unsafe) private let frameSubject = PassthroughSubject<CVPixelBuffer, Never>()
    
    nonisolated var framePublisher: AnyPublisher<CVPixelBuffer, Never> {
        frameSubject.eraseToAnyPublisher()
    }
    
    private let cameraQueue = DispatchQueue(label: "camera.capture", qos: .userInteractive)
}

Step 2: Publish Frames from Camera Callback

extension CameraManager: AVCaptureVideoDataOutputSampleBufferDelegate {
    nonisolated func captureOutput(
        _ output: AVCaptureOutput, 
        didOutput sampleBuffer: CMSampleBuffer, 
        from connection: AVCaptureConnection
    ) {
        guard let pixelBuffer = CMSampleBufferGetImageBuffer(sampleBuffer) else { return }
        
        // Broadcast to all subscribers (thread-safe)
        frameSubject.send(pixelBuffer)
    }
}

Step 3: Subscribe for UI Updates (Main Thread)

class StreamingService: ObservableObject {
    private var cancellables = Set<AnyCancellable>()
    @Published var currentFrame: CVPixelBuffer?
    
    func setupPreviewSubscriber(camera: CameraManager) {
        camera.framePublisher
            .receive(on: DispatchQueue.main)  // Thread hop for UI safety
            .sink { [weak self] buffer in
                self?.currentFrame = buffer
            }
            .store(in: &cancellables)
    }
}

Step 4: Subscribe for Encoding (Same Thread)

class VideoEncoder {
    private var frameSubscription: AnyCancellable?
    
    func startEncoding(from camera: CameraManager) {
        // No receive(on:) - stay on camera queue for performance
        frameSubscription = camera.framePublisher
            .sink { [weak self] buffer in
                self?.encode(buffer)  // Runs on camera queue
            }
    }
}

Step 5: Throttle UI Updates

@MainActor
class P2PStreamingManager: ObservableObject {
    @Published var decodedPixelBuffer: CVPixelBuffer?
    
    // Backing store updated at full frame rate
    private(set) var latestDecodedFrame: CVPixelBuffer?
    
    // Throttle @Published to ~30Hz
    private var lastFramePublishTime: CFAbsoluteTime = 0
    private let framePublishInterval: CFAbsoluteTime = 1.0 / 30.0
    
    func handleDecodedFrame(_ buffer: CVPixelBuffer) {
        latestDecodedFrame = buffer
        
        let now = CFAbsoluteTimeGetCurrent()
        if now - lastFramePublishTime >= framePublishInterval {
            decodedPixelBuffer = buffer
            lastFramePublishTime = now
        }
    }
}

Step 6: Multiple Subscriber Pattern

let framePublisher = camera.framePublisher

// UI Preview (main thread, throttled)
framePublisher
    .receive(on: DispatchQueue.main)
    .throttle(for: .milliseconds(33), scheduler: DispatchQueue.main, latest: true)
    .sink { updatePreview($0) }
    .store(in: &cancellables)

// H.264 Encoder (camera queue, full rate)
framePublisher
    .sink { encoder.encode($0) }
    .store(in: &cancellables)

// Statistics (sampled)
framePublisher
    .throttle(for: .seconds(1), scheduler: DispatchQueue.global(), latest: true)
    .sink { updateStatistics($0) }
    .store(in: &cancellables)

Common Pitfalls

1. Using CurrentValueSubject for Video

// ❌ Wrong: Stores stale frames, wastes memory
let frameSubject = CurrentValueSubject<CVPixelBuffer?, Never>(nil)

// ✅ Correct: No storage for transient data
let frameSubject = PassthroughSubject<CVPixelBuffer, Never>()

2. Blocking the Camera Queue

// ❌ Wrong: Heavy work blocks frame capture
framePublisher.sink { buffer in
    let processed = expensiveOperation(buffer)
}

// ✅ Correct: Move work off camera queue
framePublisher
    .receive(on: DispatchQueue.global(qos: .userInitiated))
    .sink { buffer in
        let processed = expensiveOperation(buffer)
    }

3. Publishing @Published Too Frequently

// ❌ Wrong: iOS rate-limits at ~120Hz, causes warnings
@Published var frame: CVPixelBuffer?
framePublisher.sink { self.frame = $0 }  // 60fps = warning spam

// ✅ Correct: Throttle @Published updates
private(set) var latestFrame: CVPixelBuffer?  // Full rate
@Published var displayFrame: CVPixelBuffer?   // Throttled to 30Hz

4. Strong Reference Cycles

// ❌ Wrong: Strong self capture
framePublisher.sink { buffer in
    self.process(buffer)
}

// ✅ Correct: Weak self
framePublisher.sink { [weak self] buffer in
    self?.process(buffer)
}

5. Not Storing Subscriptions

// ❌ Wrong: Subscription immediately cancelled
func subscribe(to camera: CameraManager) {
    camera.framePublisher.sink { ... }  // Discarded!
}

// ✅ Correct: Store in Set
private var cancellables = Set<AnyCancellable>()

func subscribe(to camera: CameraManager) {
    camera.framePublisher
        .sink { ... }
        .store(in: &cancellables)
}

Performance Tips

  1. Avoid unnecessary thread hops: If encoder runs on camera queue, don't add receive(on:)
  2. Throttle UI updates: SwiftUI doesn't need 60fps; 30Hz is sufficient
  3. Share expensive computations: Use .share() if multiple subscribers need same transform
  4. Pre-allocate resources: Create encoders/decoders before subscribing

References


Derived from CarSeet project - CameraManager.swift, P2PStreamingManager.swift

スコア

総合スコア

60/100

リポジトリの品質指標に基づく評価

SKILL.md

SKILL.mdファイルが含まれている

+20
LICENSE

ライセンスが設定されている

+10
説明文

100文字以上の説明がある

0/10
人気

GitHub Stars 100以上

0/15
最近の活動

3ヶ月以内に更新がある

0/10
フォーク

10回以上フォークされている

0/5
Issue管理

オープンIssueが50未満

+5
言語

プログラミング言語が設定されている

+5
タグ

1つ以上のタグが設定されている

0/5

レビュー

💬

レビュー機能は近日公開予定です