Skip to content

Commit 07140d7

Browse files
committed
refactor(async): integrate AsyncSampleIterator into AsyncSCStream
- Rename SampleBufferStream to AsyncSampleIterator - Integrate iterator directly into AsyncSCStream struct - Simplify API: stream.next().await instead of frames.next().await - Add buffered_count(), clear_buffer(), try_next() methods - Update example to use simplified API
1 parent b7d2946 commit 07140d7

2 files changed

Lines changed: 62 additions & 77 deletions

File tree

examples/06_async.rs

Lines changed: 2 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -49,19 +49,13 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
4949
.set_height(1080)?;
5050

5151
// Create async stream with 30-frame buffer
52-
let (stream, frames) = AsyncSCStream::new(
53-
&filter,
54-
&config,
55-
30,
56-
SCStreamOutputType::Screen
57-
);
58-
52+
let stream = AsyncSCStream::new(&filter, &config, 30, SCStreamOutputType::Screen);
5953
stream.start_capture().await?;
6054

6155
// Capture 10 frames
6256
let mut count = 0;
6357
while count < 10 {
64-
if let Some(_frame) = frames.next().await {
58+
if let Some(_frame) = stream.next().await {
6559
count += 1;
6660
println!(" Frame {}", count);
6761
}

src/async_api.rs

Lines changed: 60 additions & 69 deletions
Original file line numberDiff line numberDiff line change
@@ -325,7 +325,7 @@ impl AsyncSCContentSharingPicker {
325325
}
326326
}
327327

328-
/// Async wrapper for `SCStream` with async frame iteration
328+
/// Async wrapper for `SCStream` with integrated frame iteration
329329
///
330330
/// Provides async methods for stream lifecycle and frame iteration.
331331
/// **Executor-agnostic** - works with any async runtime.
@@ -344,41 +344,42 @@ impl AsyncSCContentSharingPicker {
344344
/// let filter = SCContentFilter::build().display(display).exclude_windows(&[]).build();
345345
/// let config = SCStreamConfiguration::build();
346346
///
347-
/// let (mut stream, frames) = AsyncSCStream::new(&filter, &config, 30, SCStreamOutputType::Screen);
347+
/// let mut stream = AsyncSCStream::new(&filter, &config, 30, SCStreamOutputType::Screen);
348348
/// stream.start_capture().await?;
349349
///
350350
/// // Process frames asynchronously
351-
/// while let Some(frame) = frames.next().await {
351+
/// while let Some(frame) = stream.next().await {
352352
/// println!("Got frame!");
353353
/// }
354354
/// # Ok(())
355355
/// # }
356356
/// ```
357357
pub struct AsyncSCStream {
358358
stream: crate::stream::SCStream,
359+
iterator: AsyncSampleIterator,
359360
}
360361

361-
/// Async stream of sample buffers
362+
/// Async iterator over sample buffers
362363
///
363364
/// Provides async iteration over captured frames.
364-
/// Created by [`AsyncSCStream::new`].
365-
pub struct SampleBufferStream {
366-
inner: Arc<Mutex<SampleBufferStreamState>>,
365+
/// Access via [`AsyncSCStream::next()`].
366+
pub struct AsyncSampleIterator {
367+
inner: Arc<Mutex<AsyncSampleIteratorState>>,
367368
}
368369

369-
struct SampleBufferStreamState {
370+
struct AsyncSampleIteratorState {
370371
buffer: std::collections::VecDeque<crate::cm::CMSampleBuffer>,
371372
waker: Option<Waker>,
372373
closed: bool,
373374
capacity: usize,
374375
}
375376

376-
/// Internal sender for sample buffer stream
377-
struct SampleBufferSender {
378-
inner: Arc<Mutex<SampleBufferStreamState>>,
377+
/// Internal sender for async sample iterator
378+
struct AsyncSampleSender {
379+
inner: Arc<Mutex<AsyncSampleIteratorState>>,
379380
}
380381

381-
impl crate::stream::output_trait::SCStreamOutputTrait for SampleBufferSender {
382+
impl crate::stream::output_trait::SCStreamOutputTrait for AsyncSampleSender {
382383
fn did_output_sample_buffer(
383384
&self,
384385
sample_buffer: crate::cm::CMSampleBuffer,
@@ -401,7 +402,7 @@ impl crate::stream::output_trait::SCStreamOutputTrait for SampleBufferSender {
401402
}
402403
}
403404

404-
impl Drop for SampleBufferSender {
405+
impl Drop for AsyncSampleSender {
405406
fn drop(&mut self) {
406407
if let Ok(mut state) = self.inner.lock() {
407408
state.closed = true;
@@ -412,56 +413,16 @@ impl Drop for SampleBufferSender {
412413
}
413414
}
414415

415-
impl SampleBufferStream {
416-
/// Get the next sample buffer asynchronously
417-
///
418-
/// Returns `None` when the stream is closed.
419-
pub fn next(&self) -> SampleBufferNext<'_> {
420-
SampleBufferNext { stream: self }
421-
}
422-
423-
/// Try to get a sample without waiting
424-
#[must_use]
425-
pub fn try_next(&self) -> Option<crate::cm::CMSampleBuffer> {
426-
self.inner.lock().ok()?.buffer.pop_front()
427-
}
428-
429-
/// Check if the stream has been closed
430-
#[must_use]
431-
pub fn is_closed(&self) -> bool {
432-
self.inner.lock().map(|s| s.closed).unwrap_or(true)
433-
}
434-
435-
/// Get the number of buffered samples
436-
#[must_use]
437-
pub fn len(&self) -> usize {
438-
self.inner.lock().map(|s| s.buffer.len()).unwrap_or(0)
439-
}
440-
441-
/// Check if the buffer is empty
442-
#[must_use]
443-
pub fn is_empty(&self) -> bool {
444-
self.len() == 0
445-
}
446-
447-
/// Clear all buffered samples
448-
pub fn clear(&self) {
449-
if let Ok(mut state) = self.inner.lock() {
450-
state.buffer.clear();
451-
}
452-
}
453-
}
454-
455416
/// Future for getting the next sample buffer
456-
pub struct SampleBufferNext<'a> {
457-
stream: &'a SampleBufferStream,
417+
pub struct NextSample<'a> {
418+
iterator: &'a AsyncSampleIterator,
458419
}
459420

460-
impl Future for SampleBufferNext<'_> {
421+
impl Future for NextSample<'_> {
461422
type Output = Option<crate::cm::CMSampleBuffer>;
462423

463424
fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
464-
let Ok(mut state) = self.stream.inner.lock() else {
425+
let Ok(mut state) = self.iterator.inner.lock() else {
465426
return Poll::Ready(None);
466427
};
467428

@@ -478,15 +439,13 @@ impl Future for SampleBufferNext<'_> {
478439
}
479440
}
480441

481-
unsafe impl Send for SampleBufferStream {}
482-
unsafe impl Sync for SampleBufferStream {}
483-
unsafe impl Send for SampleBufferSender {}
484-
unsafe impl Sync for SampleBufferSender {}
442+
unsafe impl Send for AsyncSampleIterator {}
443+
unsafe impl Sync for AsyncSampleIterator {}
444+
unsafe impl Send for AsyncSampleSender {}
445+
unsafe impl Sync for AsyncSampleSender {}
485446

486447
impl AsyncSCStream {
487-
/// Create a new async stream with frame iteration
488-
///
489-
/// Returns both the stream and an async iterator for frames.
448+
/// Create a new async stream
490449
///
491450
/// # Arguments
492451
///
@@ -500,24 +459,56 @@ impl AsyncSCStream {
500459
config: &SCStreamConfiguration,
501460
buffer_capacity: usize,
502461
output_type: crate::stream::output_type::SCStreamOutputType,
503-
) -> (Self, SampleBufferStream) {
504-
let state = Arc::new(Mutex::new(SampleBufferStreamState {
462+
) -> Self {
463+
let state = Arc::new(Mutex::new(AsyncSampleIteratorState {
505464
buffer: std::collections::VecDeque::with_capacity(buffer_capacity),
506465
waker: None,
507466
closed: false,
508467
capacity: buffer_capacity,
509468
}));
510469

511-
let sender = SampleBufferSender {
470+
let sender = AsyncSampleSender {
512471
inner: Arc::clone(&state),
513472
};
514473

515474
let mut stream = crate::stream::SCStream::new(filter, config);
516475
stream.add_output_handler(sender, output_type);
517476

518-
let receiver = SampleBufferStream { inner: state };
477+
let iterator = AsyncSampleIterator { inner: state };
478+
479+
Self { stream, iterator }
480+
}
481+
482+
/// Get the next sample buffer asynchronously
483+
///
484+
/// Returns `None` when the stream is closed.
485+
pub fn next(&self) -> NextSample<'_> {
486+
NextSample { iterator: &self.iterator }
487+
}
519488

520-
(Self { stream }, receiver)
489+
/// Try to get a sample without waiting
490+
#[must_use]
491+
pub fn try_next(&self) -> Option<crate::cm::CMSampleBuffer> {
492+
self.iterator.inner.lock().ok()?.buffer.pop_front()
493+
}
494+
495+
/// Check if the stream has been closed
496+
#[must_use]
497+
pub fn is_closed(&self) -> bool {
498+
self.iterator.inner.lock().map(|s| s.closed).unwrap_or(true)
499+
}
500+
501+
/// Get the number of buffered samples
502+
#[must_use]
503+
pub fn buffered_count(&self) -> usize {
504+
self.iterator.inner.lock().map(|s| s.buffer.len()).unwrap_or(0)
505+
}
506+
507+
/// Clear all buffered samples
508+
pub fn clear_buffer(&self) {
509+
if let Ok(mut state) = self.iterator.inner.lock() {
510+
state.buffer.clear();
511+
}
521512
}
522513

523514
/// Start capture asynchronously

0 commit comments

Comments
 (0)