Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -6,3 +6,4 @@ zig-pkg
/*.h264
/*.bmp
/*.mp4
/*.mov
47 changes: 32 additions & 15 deletions src/audio/audio_replay_buffer.zig
Original file line number Diff line number Diff line change
Expand Up @@ -64,7 +64,7 @@ pub fn add_data(self: *Self, data: Arc(AudioCaptureData)) !void {
var ready_packets = self.timeline.take_ready_packets();
defer deinitPacketList(&ready_packets);
self.append_packets(&ready_packets);
self.trim_packets();
self.trim_packets(.{});
}

/// Flush any remaining packets in the timeline.
Expand All @@ -74,12 +74,12 @@ pub fn finalize(self: *Self) !void {
var ready_packets = self.timeline.take_ready_packets();
defer deinitPacketList(&ready_packets);
self.append_packets(&ready_packets);
self.trim_packets();
self.trim_packets(.{});
}

pub fn set_replay_seconds(self: *Self, replay_seconds: u32) void {
self.replay_seconds = replay_seconds;
self.trim_packets();
self.trim_packets(.{});
}

pub fn packet_iterator(self: *Self) LinkedListIterator(EncodedAudioPacketNode) {
Expand All @@ -99,19 +99,36 @@ fn append_packets(self: *Self, packets: *std.DoublyLinkedList) void {
}
}

/// Remove packets that are older than the configured replay duration.
fn trim_packets(self: *Self) void {
const retention_samples = self.replay_retention_samples();
const oldest_sample = self.timeline.encoded_until_sample - retention_samples;
fn remove_first_packet(self: *Self) void {
if (self.packets.popFirst()) |first| {
const packet_node: *EncodedAudioPacketNode = @fieldParentPtr("node", first);
self.size -= @intCast(packet_node.data.*.size);
self.len -= 1;
packet_node.deinit();
}
}

/// Remove packets that are older than the configured replay duration, or older
/// than `oldest_time_ns` when provided.
pub fn trim_packets(self: *Self, args: struct {
oldest_time_ns: ?i128 = null,
}) void {
const oldest_sample = if (args.oldest_time_ns) |oldest_time_ns| blk: {
break :blk self.timeline.timestamp_to_sample_floor(oldest_time_ns);
} else blk: {
const retention_samples = self.replay_retention_samples();
break :blk self.timeline.encoded_until_sample - retention_samples;
};

if (oldest_sample == null) {
return;
}

while (self.packets.first) |first| {
const packet_node: *EncodedAudioPacketNode = @fieldParentPtr("node", first);
const packet_end = packet_node.data.*.pts + packet_node.data.*.duration;
if (packet_end <= oldest_sample) {
_ = self.packets.popFirst();
self.len -= 1;
self.size -= @intCast(packet_node.data.*.size);
packet_node.deinit();
if (packet_end <= oldest_sample.?) {
self.remove_first_packet();
} else {
break;
}
Expand All @@ -122,8 +139,8 @@ fn replay_retention_samples(self: *Self) i64 {
return self.replay_seconds * SAMPLE_RATE;
}

const TestUtil = struct {
fn create_audio_capture_data(
pub const TestUtil = struct {
pub fn create_audio_capture_data(
allocator: Allocator,
timestamp_ns: i128,
samples_per_channel: usize,
Expand Down Expand Up @@ -159,7 +176,7 @@ test "AudioReplayBuffer - add_data - encodes audio before export and exposes pac
try replay_buffer.finalize();

try std.testing.expect(replay_buffer.has_packets());
const packet_window = replay_buffer.timeline.get_sample_window(
const packet_window = replay_buffer.timeline.get_unclamped_sample_window(
std.time.ns_per_s,
std.time.ns_per_s + std.time.ns_per_s / 10,
) orelse return error.ExpectedSampleWindow;
Expand Down
13 changes: 3 additions & 10 deletions src/audio/audio_timeline.zig
Original file line number Diff line number Diff line change
Expand Up @@ -138,7 +138,7 @@ pub const AudioTimeline = struct {
device.value_ptr.* = .{};
}

var start_sample = self.timestamp_to_sample_floor(audio_capture_data.start_ns());
var start_sample = self.timestamp_to_sample_floor(audio_capture_data.start_ns()).?;

// Chunks don't always arrive at exactly the timestamps they are
// expected to. Small jitter here can cause static once multiple devices
Expand Down Expand Up @@ -193,14 +193,6 @@ pub const AudioTimeline = struct {
};
}

pub fn get_sample_window(self: *Self, start_time_ns: i128, end_time_ns: i128) ?SampleWindow {
if (self.timeline_origin_ns == null) return null;
return .{
.start_sample = self.timestamp_to_sample_floor(start_time_ns),
.end_sample = self.timestamp_to_sample_ceil(end_time_ns),
};
}

/// Export alignment sometimes needs a window that begins before the first
/// captured audio sample. Keep that offset negative so muxing preserves the
/// initial gap instead of shifting audio earlier to time zero. This can occur
Expand Down Expand Up @@ -292,7 +284,8 @@ pub const AudioTimeline = struct {
}

/// Floor is used for starts so a chunk never begins after its true time.
fn timestamp_to_sample_floor(self: *Self, timestamp_ns: i128) i64 {
pub fn timestamp_to_sample_floor(self: *Self, timestamp_ns: i128) ?i64 {
if (self.timeline_origin_ns == null) return null;
return @max(self.timestamp_to_sample_floor_unclamped(timestamp_ns), 0);
}

Expand Down
2 changes: 1 addition & 1 deletion src/store/audio_session.zig
Original file line number Diff line number Diff line change
Expand Up @@ -142,7 +142,7 @@ pub const AudioSession = struct {
try replay_buffer.add_data(data.clone());
self.store.dispatch(.{
.capture = .{
.update_replay_buffer_size = .{ .audio_bytes = replay_buffer.size },
.update_replay_buffer_metrics = .{ .audio_bytes = replay_buffer.size },
},
});
}
Expand Down
97 changes: 92 additions & 5 deletions src/store/capture_store.zig
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ const Muxer = @import("../video/muxer.zig").Muxer;
const Mutex = @import("../mutex.zig").Mutex;
const VideoCaptureSelection = @import("../capture/video/video_capture.zig").VideoCaptureSelection;
const VideoReplayBuffer = @import("../video/video_replay_buffer.zig").VideoReplayBuffer;
const AudioReplayBuffer = @import("../audio/audio_replay_buffer.zig");
const exporter = @import("../exporter.zig");
const AudioCaptureData = @import("../capture/audio/audio_capture_data.zig");
const Arc = @import("../arc.zig").Arc;
Expand Down Expand Up @@ -75,7 +76,7 @@ pub const CaptureStore = struct {
update_audio_device_level: UpdateAudioDeviceLevelPayload,
start_audio_capture_thread,

update_replay_buffer_size: union(enum) {
update_replay_buffer_metrics: union(enum) {
audio_bytes: u64,
video: struct { bytes: u64, start_time: ?std.Io.Timestamp },
},
Expand All @@ -87,6 +88,12 @@ pub const CaptureStore = struct {
stop_replay_buffer_success,
stop_replay_buffer_fail,

/// When the video replay buffer state changes, sync the audio replay buffer with it.
/// This is really only necessary in the case that the replay length is infinite, but
/// the user relies on max bytes to trim the replay buffer. Without this, the audio
/// replay buffer would grow infinitely.
sync_replay_buffers,

save_replay,
save_replay_success,
save_replay_fail,
Expand Down Expand Up @@ -137,6 +144,7 @@ pub const CaptureStore = struct {
.stop_video_capture = .{effect_stop_video_capture},
.select_video_source = .{effect_select_video_source},
.select_video_source_prepared = .{effect_select_video_source_prepared},
.sync_replay_buffers = .{effect_sync_replay_buffers},
.save_replay = .{effect_save_replay},
};

Expand Down Expand Up @@ -363,7 +371,7 @@ pub const CaptureStore = struct {
break;
}
},
.update_replay_buffer_size => |payload| {
.update_replay_buffer_metrics => |payload| {
switch (payload) {
.audio_bytes => |audio_bytes| {
state.capture.replay_buffer_metrics.audio_bytes = audio_bytes;
Expand Down Expand Up @@ -392,7 +400,11 @@ pub const CaptureStore = struct {

pub fn effect_sync_replay_buffer_with_user_settings(store: *Store, replay_seconds: u32) void {
store.capture_store.audio_session.set_replay_buffer_seconds(replay_seconds);
store.capture_store.video_session.set_replay_buffer_seconds(replay_seconds);
store.capture_store.video_session.set_replay_buffer_values(.{ .replay_seconds = replay_seconds });
}

pub fn effect_sync_replay_buffer_max_bytes(store: *Store, replay_max_bytes: u64) void {
store.capture_store.video_session.set_replay_buffer_values(.{ .replay_max_bytes = replay_max_bytes });
}
// ----------------------------------------------------------------------------

Expand Down Expand Up @@ -557,6 +569,7 @@ pub const CaptureStore = struct {
.fps = state.user_settings.user_settings.capture_fps,
.bit_rate = state.user_settings.user_settings.capture_bit_rate,
.replay_seconds = state.user_settings.user_settings.replay_seconds,
.replay_max_bytes = state.user_settings.user_settings.replay_max_bytes,
};
};

Expand All @@ -568,6 +581,7 @@ pub const CaptureStore = struct {
state_local.fps,
state_local.bit_rate,
state_local.replay_seconds,
state_local.replay_max_bytes,
);
store.dispatch(.{ .capture = .start_replay_buffer_success });
}
Expand Down Expand Up @@ -672,12 +686,33 @@ pub const CaptureStore = struct {
store.dispatch(.{ .capture = .stop_recording_to_disk_success });
}

/// See sync_replay_buffers message type for details.
fn effect_sync_replay_buffers(store: *Store, _: anytype) !void {
const self = &store.capture_store;
const video_start_ns = blk: {
var replay_buffer_locked = self.video_session.video_replay_buffer.lock();
defer replay_buffer_locked.unlock();
const replay_buffer = replay_buffer_locked.unwrap() orelse return;
const start_time = replay_buffer.get_start_time() orelse return;
break :blk start_time.nanoseconds;
};
const audio_bytes = blk: {
var replay_buffer_locked = self.audio_session.audio_replay_buffer.lock();
defer replay_buffer_locked.unlock();
const replay_buffer = replay_buffer_locked.unwrap() orelse return;
replay_buffer.trim_packets(.{ .oldest_time_ns = video_start_ns });
break :blk replay_buffer.size;
};
store.dispatch(.{ .capture = .{ .update_replay_buffer_metrics = .{ .audio_bytes = audio_bytes } } });
}

fn effect_save_replay(store: *Store, _: anytype) !void {
const self = &store.capture_store;
errdefer store.dispatch(.{ .capture = .save_replay_fail });

var fps: u32 = 0;
var replay_seconds: u32 = 0;
var replay_max_bytes: u64 = 0;
var video_output_directory: ?String = null;
defer {
if (video_output_directory) |*_video_output_directory| _video_output_directory.deinit();
Expand All @@ -695,6 +730,7 @@ pub const CaptureStore = struct {
const settings = state.user_settings.user_settings;
fps = settings.capture_fps;
replay_seconds = settings.replay_seconds;
replay_max_bytes = settings.replay_max_bytes;
// video_output_directory should never be null at this point. If so, there is
// something seriously wrong.
assert(settings.video_output_directory != null);
Expand All @@ -712,6 +748,7 @@ pub const CaptureStore = struct {

const video_replay_buffer: ?*VideoReplayBuffer = (try self.video_session.take_and_swap_replay_buffer(
replay_seconds,
replay_max_bytes,
));
defer if (video_replay_buffer) |_video_replay_buffer| _video_replay_buffer.deinit();

Expand Down Expand Up @@ -980,12 +1017,12 @@ test "CaptureStore - update_replay_buffer_size" {

const size = 1024 * 1024 * 10; // 10MB

store.dispatch(.{ .capture = .{ .update_replay_buffer_size = .{ .audio_bytes = size } } });
store.dispatch(.{ .capture = .{ .update_replay_buffer_metrics = .{ .audio_bytes = size } } });
store.run(.{ .once = true, .wait_for_effects = true });
const start_time = std.Io.Timestamp.now(std.testing.io, .awake);
store.dispatch(.{
.capture = .{
.update_replay_buffer_size = .{
.update_replay_buffer_metrics = .{
.video = .{
.bytes = size,
.start_time = start_time,
Expand All @@ -1001,6 +1038,56 @@ test "CaptureStore - update_replay_buffer_size" {
try std.testing.expectEqual(20, state.capture.replay_buffer_metrics.size_in_mb(.total));
}

test "CaptureStore - sync_replay_buffers - should remove audio frames when the video replay buffer is trimmed" {
const TestStore = @import("./store.zig").TestStore;
const AudioReplayBufferTestUtil = @import("../audio/audio_replay_buffer.zig").TestUtil;
const allocator = std.testing.allocator;
const test_store = try TestStore.init(allocator);
defer test_store.deinit();
const store = test_store.store;
const state = &store.state.private.value;

const video_start_ns = (2 * std.time.ns_per_s);

var audio_replay_buffer = try AudioReplayBuffer.init(allocator, 10);
store.capture_store.audio_session.audio_replay_buffer.set(audio_replay_buffer);

for (0..4) |second| {
const chunk = try AudioReplayBufferTestUtil.create_audio_capture_data(
allocator,
(@as(i128, @intCast(second)) * std.time.ns_per_s),
2048,
0.1,
);

try audio_replay_buffer.add_data(chunk);
}
try audio_replay_buffer.finalize();
const audio_bytes_before = audio_replay_buffer.size;
try std.testing.expect(audio_bytes_before > 0);

var video_replay_buffer = try VideoReplayBuffer.init(allocator, 10, 0, &.{});
store.capture_store.video_session.video_replay_buffer.set(video_replay_buffer);

try video_replay_buffer.add_frame(&.{1}, video_start_ns, true);
try video_replay_buffer.add_frame(&.{2}, video_start_ns + std.time.ns_per_s, false);

store.dispatch(.{ .capture = .sync_replay_buffers });
store.run(.{ .once = true, .wait_for_effects = true });
store.run(.{ .once = true, .wait_for_effects = true });

try std.testing.expect(audio_replay_buffer.size < audio_bytes_before);
try std.testing.expectEqual(audio_replay_buffer.size, state.capture.replay_buffer_metrics.audio_bytes);

// Ensure that there are no audio packets that occurred before the oldest video frame.
const oldest_audio_sample = audio_replay_buffer.timeline.timestamp_to_sample_floor(video_start_ns) orelse return error.ExpectedOldestAudioSample;
var iter = audio_replay_buffer.packet_iterator();
while (iter.next()) |packet| {
const packet_end = packet.data.*.pts + packet.data.*.duration;
try std.testing.expect(packet_end > oldest_audio_sample);
}
}

// ----------------------------------------------------------------------------
// TODO: Still need to write tests for the rest of the message types.
// ----------------------------------------------------------------------------
Expand Down
3 changes: 2 additions & 1 deletion src/store/store.zig
Original file line number Diff line number Diff line change
Expand Up @@ -288,7 +288,8 @@ pub const Store = struct {
effect_fn(store, effect_payload);
},
}
log.debug("[execute_registered_effects] effect: {s}", .{@typeName(@TypeOf(effect_fn))});
// This is super chatty - we probably don't want this on anymore.
// log.debug("[execute_registered_effects] effect: {s}", .{@typeName(@TypeOf(effect_fn))});
} else {
@compileError(@typeName(@TypeOf(effect_fn)) ++ " has no return type");
}
Expand Down
Loading
Loading