FP Perception 0.1.2
Loading...
Searching...
No Matches
audio_buffer.hpp
Go to the documentation of this file.
1#pragma once
2
3#include <algorithm>
4#include <chrono>
5#include <condition_variable>
6#include <cstddef>
7#include <cstdint>
8#include <mutex>
9
10#include <rclcpp/rclcpp.hpp>
11
14
16{
17
19{
20public:
21 static constexpr const char* kExpiredAudioSliceError = "error: requested audio slice is expired. for future "
22 "interactions, increase the audio retention window parameter.";
23
24 void append(const audio_data& data, int max_duration_seconds)
25 {
26 append(data, max_duration_seconds, rclcpp::Clock().now());
27 }
28
29 void append(const audio_data& data, int max_duration_seconds, const rclcpp::Time& end_time)
30 {
31 if (data.samples.empty())
32 return;
33
34 std::unique_lock<std::mutex> lock(mutex_);
35
36 const int channels = std::max(1, data.channels);
37 const int sample_rate = std::max(1, data.sample_rate);
38 max_duration_seconds_ = std::max(1, max_duration_seconds);
39 const size_t incoming_frames = data.samples.size() / static_cast<size_t>(channels);
40 if (incoming_frames == 0)
41 return;
42
44 {
45 if (!buffer_.samples.empty())
46 {
47 RCLCPP_WARN(rclcpp::get_logger("AudioBuffer"),
48 "Resetting public audio buffer due to format change: sample_rate %d -> %d, channels %d -> %d, "
49 "buffered_samples=%zu",
51 }
52 buffer_.samples.clear();
54 total_frames_ = 0;
55 initialized_time_ = false;
56 buffer_end_time_ = rclcpp::Time(0, 0, RCL_ROS_TIME);
57 }
58
59 const rclcpp::Duration incoming_duration = framesToDuration(incoming_frames, sample_rate);
60
62 {
63 buffer_end_time_ = end_time;
64 buffer_start_time_ = buffer_end_time_ - incoming_duration;
65 initialized_time_ = true;
66 }
67 else if (end_time > buffer_end_time_)
68 {
69 const auto observed_gap = end_time - buffer_end_time_;
70 if (observed_gap > incoming_duration)
71 {
72 const auto missing_duration = observed_gap - incoming_duration;
73 const size_t missing_frames = durationToFrames(missing_duration, sample_rate);
74 const size_t missing_samples = missing_frames * static_cast<size_t>(channels);
75 if (missing_samples > 0)
76 {
77 buffer_.samples.insert(buffer_.samples.end(), missing_samples, 0);
78 total_samples_ += static_cast<uint64_t>(missing_samples);
79 total_frames_ += static_cast<uint64_t>(missing_frames);
80 }
81 }
82 }
83
89
90 buffer_.samples.insert(buffer_.samples.end(), data.samples.begin(), data.samples.end());
91 total_samples_ += static_cast<uint64_t>(data.samples.size());
92 total_frames_ += static_cast<uint64_t>(incoming_frames);
93
94 const size_t max_buffer_size = static_cast<size_t>(std::max(1, data.sample_rate)) *
95 static_cast<size_t>(std::max(1, data.channels)) *
96 static_cast<size_t>(std::max(1, max_duration_seconds));
97
98 if (buffer_.samples.size() > max_buffer_size)
99 {
100 const size_t excess_samples = buffer_.samples.size() - max_buffer_size;
101 const size_t excess_frames = excess_samples / static_cast<size_t>(channels);
102 const size_t samples_to_remove = excess_frames * static_cast<size_t>(channels);
103 if (samples_to_remove > 0)
104 {
105 buffer_.samples.erase(buffer_.samples.begin(),
106 buffer_.samples.begin() + static_cast<std::ptrdiff_t>(samples_to_remove));
107 buffer_start_time_ = buffer_start_time_ + framesToDuration(excess_frames, sample_rate);
108 }
109 }
110
111 buffer_end_time_ = end_time;
112 const size_t buffered_frames = buffer_.samples.size() / static_cast<size_t>(channels);
113 buffer_start_time_ = buffer_end_time_ - framesToDuration(buffered_frames, sample_rate);
114
115 const size_t denom =
116 static_cast<size_t>(std::max(1, data.chunk_size)) * static_cast<size_t>(std::max(1, data.channels));
117 buffer_.chunk_count = denom ? (buffer_.samples.size() / denom) : 0;
118
119 const double buffered_seconds = static_cast<double>(buffered_frames) / static_cast<double>(sample_rate);
120 if (buffered_seconds < 0.5)
121 {
122 RCLCPP_DEBUG(rclcpp::get_logger("AudioBuffer"),
123 "Public audio buffer remains small after append: added_frames=%zu buffered_frames=%zu "
124 "buffered_seconds=%.3f retention=%d start=%.9f end=%.9f",
125 incoming_frames, buffered_frames, buffered_seconds, max_duration_seconds,
126 buffer_start_time_.seconds(), buffer_end_time_.seconds());
127 }
128
129 lock.unlock();
130 condition_variable_.notify_all();
131 }
132
133 void waitForAudio(std::chrono::seconds timeout = std::chrono::seconds(5))
134 {
135 std::unique_lock<std::mutex> lock(mutex_);
136
138 return;
139
140 condition_variable_.wait_for(lock, timeout, [this] { return isInitializedLocked(); });
141
142 if (!isInitializedLocked())
143 throw fp_perception_exception("public audio buffer not initialized");
144 }
145
146 audio_data readLatest(int duration_seconds)
147 {
148 duration_seconds = std::max(1, duration_seconds);
149 waitForAudio();
150
151 std::lock_guard<std::mutex> lock(mutex_);
152
153 const int sample_rate = buffer_.sample_rate;
154 const int channels = std::max(1, buffer_.channels);
155 const size_t requested_samples =
156 static_cast<size_t>(sample_rate) * static_cast<size_t>(channels) * static_cast<size_t>(duration_seconds);
157 const size_t available_samples = buffer_.samples.size();
158
159 if (available_samples == 0)
160 throw fp_perception_exception("public audio buffer is empty");
161
162 audio_data out = buffer_;
163 out.sample_rate = sample_rate;
164 out.channels = channels;
165 out.chunk_count = 1;
166
167 const size_t samples_to_copy = std::min(requested_samples, available_samples);
168 const size_t start_index = available_samples - samples_to_copy;
169
170 out.samples.assign(buffer_.samples.begin() + static_cast<std::ptrdiff_t>(start_index), buffer_.samples.end());
171 out.chunk_size = static_cast<int>(samples_to_copy / static_cast<size_t>(channels));
172
173 return out;
174 }
175
176 audio_data readWindow(const rclcpp::Time& start_time, int duration_seconds)
177 {
178 duration_seconds = std::max(1, duration_seconds);
179 waitForAudio();
180
181 std::lock_guard<std::mutex> lock(mutex_);
182
183 if (!initialized_time_ || buffer_.samples.empty())
184 throw fp_perception_exception("public audio buffer timestamp state is not initialized");
185
186 const int sample_rate = std::max(1, buffer_.sample_rate);
187 const int channels = std::max(1, buffer_.channels);
188 const rclcpp::Time buffer_end_time = currentEndTimeLocked();
189 const rclcpp::Time requested_end_time = start_time + rclcpp::Duration::from_seconds(duration_seconds);
190
191 audio_data out = buffer_;
192 const size_t requested_frames = static_cast<size_t>(sample_rate) * static_cast<size_t>(duration_seconds);
193 out.samples.clear();
194 out.chunk_size = 0;
195 out.chunk_count = 1;
196 out.sample_rate = sample_rate;
197 out.channels = channels;
198
199 const rclcpp::Time overlap_start = start_time < buffer_start_time_ ? buffer_start_time_ : start_time;
200 const rclcpp::Time overlap_end = requested_end_time > buffer_end_time ? buffer_end_time : requested_end_time;
201 if (overlap_end <= overlap_start)
202 return out;
203
204 const size_t source_start_frame = durationToFrames(overlap_start - buffer_start_time_, sample_rate);
205 const size_t output_start_frame = durationToFrames(overlap_start - start_time, sample_rate);
206 size_t frames_to_copy = durationToFrames(overlap_end - overlap_start, sample_rate);
207
208 const size_t available_frames = buffer_.samples.size() / static_cast<size_t>(channels);
209 if (source_start_frame >= available_frames || output_start_frame >= requested_frames)
210 return out;
211
212 frames_to_copy = std::min(frames_to_copy, available_frames - source_start_frame);
213 frames_to_copy = std::min(frames_to_copy, requested_frames - output_start_frame);
214 if (frames_to_copy == 0)
215 return out;
216
217 const size_t source_start_sample = source_start_frame * static_cast<size_t>(channels);
218 const size_t samples_to_copy = frames_to_copy * static_cast<size_t>(channels);
219 out.samples.assign(buffer_.samples.begin() + static_cast<std::ptrdiff_t>(source_start_sample),
220 buffer_.samples.begin() + static_cast<std::ptrdiff_t>(source_start_sample + samples_to_copy));
221 out.chunk_size = static_cast<int>(frames_to_copy);
222
223 return out;
224 }
225
226 void waitForWindow(const rclcpp::Time& start_time, int duration_seconds,
227 std::chrono::seconds timeout = std::chrono::seconds(30))
228 {
229 duration_seconds = std::max(1, duration_seconds);
230 waitForAudio(timeout);
231
232 const rclcpp::Time requested_end_time = start_time + rclcpp::Duration::from_seconds(duration_seconds);
233 std::unique_lock<std::mutex> lock(mutex_);
234
235 condition_variable_.wait_for(lock, timeout, [this, &start_time, &requested_end_time] {
237 return false;
238
239 if (isExpiredLocked(start_time))
240 return true;
241
242 return currentEndTimeLocked() >= requested_end_time;
243 });
244
246 throw fp_perception_exception("public audio buffer timestamp state is not initialized");
247
248 if (isExpiredLocked(start_time))
249 {
250 const auto current_end_time = currentEndTimeLocked();
252 std::string(kExpiredAudioSliceError) + " requested=[" + std::to_string(start_time.seconds()) + ", " +
253 std::to_string(requested_end_time.seconds()) + "] buffered=[" + std::to_string(buffer_start_time_.seconds()) +
254 ", " + std::to_string(current_end_time.seconds()) + "]");
255 }
256
257 if (currentEndTimeLocked() < requested_end_time)
258 {
259 const auto current_end_time = currentEndTimeLocked();
261 "Timeout waiting for requested audio slice. requested=[" + std::to_string(start_time.seconds()) + ", " +
262 std::to_string(requested_end_time.seconds()) + "] buffered=[" + std::to_string(buffer_start_time_.seconds()) +
263 ", " + std::to_string(current_end_time.seconds()) + "]");
264 }
265 }
266
267 rclcpp::Time startTime() const
268 {
269 std::lock_guard<std::mutex> lock(mutex_);
270 return buffer_start_time_;
271 }
272
273 rclcpp::Time endTime() const
274 {
275 std::lock_guard<std::mutex> lock(mutex_);
276 return currentEndTimeLocked();
277 }
278
279private:
280 static rclcpp::Duration framesToDuration(size_t frames, int sample_rate)
281 {
282 const int64_t nanoseconds = static_cast<int64_t>((static_cast<long double>(frames) * 1000000000.0L) /
283 static_cast<long double>(std::max(1, sample_rate)));
284 return rclcpp::Duration::from_nanoseconds(nanoseconds);
285 }
286
287 static size_t durationToFrames(const rclcpp::Duration& duration, int sample_rate)
288 {
289 if (duration.nanoseconds() <= 0)
290 return 0;
291 return static_cast<size_t>(
292 (static_cast<long double>(duration.nanoseconds()) * static_cast<long double>(std::max(1, sample_rate))) /
293 1000000000.0L);
294 }
295
297 {
298 return buffer_.sample_rate > 0 && buffer_.channels > 0 && !buffer_.samples.empty();
299 }
300
301 rclcpp::Time currentEndTimeLocked() const
302 {
303 return buffer_end_time_;
304 }
305
306 bool isExpiredLocked(const rclcpp::Time& start_time) const
307 {
308 if (max_duration_seconds_ <= 0)
309 return false;
310
311 const auto retention_start = currentEndTimeLocked() - rclcpp::Duration::from_seconds(max_duration_seconds_);
312 return retention_start > start_time;
313 }
314
316 mutable std::mutex mutex_;
317 std::condition_variable condition_variable_;
318 uint64_t total_samples_{ 0 };
319 uint64_t total_frames_{ 0 };
321 bool initialized_time_{ false };
322 rclcpp::Time buffer_start_time_{ 0, 0, RCL_ROS_TIME };
323 rclcpp::Time buffer_end_time_{ 0, 0, RCL_ROS_TIME };
324};
325
326} // namespace fp_perception
Definition audio_buffer.hpp:19
uint64_t total_samples_
Definition audio_buffer.hpp:318
rclcpp::Time endTime() const
Definition audio_buffer.hpp:273
rclcpp::Time startTime() const
Definition audio_buffer.hpp:267
bool isInitializedLocked() const
Definition audio_buffer.hpp:296
void append(const audio_data &data, int max_duration_seconds, const rclcpp::Time &end_time)
Definition audio_buffer.hpp:29
static size_t durationToFrames(const rclcpp::Duration &duration, int sample_rate)
Definition audio_buffer.hpp:287
uint64_t total_frames_
Definition audio_buffer.hpp:319
static rclcpp::Duration framesToDuration(size_t frames, int sample_rate)
Definition audio_buffer.hpp:280
rclcpp::Time buffer_start_time_
Definition audio_buffer.hpp:322
rclcpp::Time buffer_end_time_
Definition audio_buffer.hpp:323
std::condition_variable condition_variable_
Definition audio_buffer.hpp:317
bool isExpiredLocked(const rclcpp::Time &start_time) const
Definition audio_buffer.hpp:306
static constexpr const char * kExpiredAudioSliceError
Definition audio_buffer.hpp:21
audio_data readWindow(const rclcpp::Time &start_time, int duration_seconds)
Definition audio_buffer.hpp:176
void waitForWindow(const rclcpp::Time &start_time, int duration_seconds, std::chrono::seconds timeout=std::chrono::seconds(30))
Definition audio_buffer.hpp:226
bool initialized_time_
Definition audio_buffer.hpp:321
void waitForAudio(std::chrono::seconds timeout=std::chrono::seconds(5))
Definition audio_buffer.hpp:133
int max_duration_seconds_
Definition audio_buffer.hpp:320
std::mutex mutex_
Definition audio_buffer.hpp:316
audio_data readLatest(int duration_seconds)
Definition audio_buffer.hpp:146
audio_data buffer_
Definition audio_buffer.hpp:315
rclcpp::Time currentEndTimeLocked() const
Definition audio_buffer.hpp:301
void append(const audio_data &data, int max_duration_seconds)
Definition audio_buffer.hpp:24
Definition audio_buffer.hpp:16
Struct to hold audio data.
Definition structs.hpp:16
int chunk_count
Number of chunks in the audio data.
Definition structs.hpp:21
std::vector< int16_t > samples
Audio samples.
Definition structs.hpp:17
int sample_rate
Sample rate in Hz.
Definition structs.hpp:18
int channels
Number of audio channels.
Definition structs.hpp:19
bool override
if the system sample_rate, channels, and chunk_size should be overridden by message data
Definition structs.hpp:22
int chunk_size
Size of each audio chunk in samples.
Definition structs.hpp:20
Base class for driver exceptions.
Definition exceptions.hpp:14