[v3,27/41] libcamera: internal: Add a BufferQueue class to handle buffer queues
diff mbox series

Message ID 20260914140309.3354666-28-stefan.klug@ideasonboard.com
State New
Headers show
Series
  • rkisp1: pipeline rework for PFC
Related show

Commit Message

Stefan Klug Sept. 14, 2026, 2:02 p.m. UTC
Add a class that encapsulates a queue of v4l2 buffers and the typical
use-cases. This simplifies manual queue management and helps in cases
where a pre or postprocessing stage is needed.

Signed-off-by: Stefan Klug <stefan.klug@ideasonboard.com>

---

Changes in v3:
- Fixed incomplete reset in BufferQueue::releaseBuffers()
- Added missing disconnect in ~BufferQueueDelegate
- Added ASSERT in releaseBuffers()
- Updated state diagram to show that cancelled buffers jump the
  postprocessing stage
- Improved documentation

Changes in v2:
- Added this patch
---
 include/libcamera/internal/buffer_queue.h | 136 ++++++
 include/libcamera/internal/meson.build    |   1 +
 src/libcamera/buffer_queue.cpp            | 541 ++++++++++++++++++++++
 src/libcamera/meson.build                 |   1 +
 4 files changed, 679 insertions(+)
 create mode 100644 include/libcamera/internal/buffer_queue.h
 create mode 100644 src/libcamera/buffer_queue.cpp

Comments

Barnabás Pőcze Sept. 21, 2026, 2:47 p.m. UTC | #1
Hi

2026. 09. 14. 16:02 keltezéssel, Stefan Klug írta:
> Add a class that encapsulates a queue of v4l2 buffers and the typical
> use-cases. This simplifies manual queue management and helps in cases
> where a pre or postprocessing stage is needed.
> 
> Signed-off-by: Stefan Klug <stefan.klug@ideasonboard.com>
> 
> ---
> 
> Changes in v3:
> - Fixed incomplete reset in BufferQueue::releaseBuffers()
> - Added missing disconnect in ~BufferQueueDelegate
> - Added ASSERT in releaseBuffers()
> - Updated state diagram to show that cancelled buffers jump the
>    postprocessing stage
> - Improved documentation
> 
> Changes in v2:
> - Added this patch
> ---
>   include/libcamera/internal/buffer_queue.h | 136 ++++++
>   include/libcamera/internal/meson.build    |   1 +
>   src/libcamera/buffer_queue.cpp            | 541 ++++++++++++++++++++++
>   src/libcamera/meson.build                 |   1 +
>   4 files changed, 679 insertions(+)
>   create mode 100644 include/libcamera/internal/buffer_queue.h
>   create mode 100644 src/libcamera/buffer_queue.cpp
> 
> diff --git a/include/libcamera/internal/buffer_queue.h b/include/libcamera/internal/buffer_queue.h
> new file mode 100644
> index 000000000000..83a6a1f45be7
> --- /dev/null
> +++ b/include/libcamera/internal/buffer_queue.h
> @@ -0,0 +1,136 @@
> +/* SPDX-License-Identifier: LGPL-2.1-or-later */
> +/*
> + * Copyright (C) 2025, Ideas on Board
> + *
> + * Sequence sync helper
> + */
> +
> +#pragma once
> +
> +#include <map>
> +#include <memory>
> +
> +#include <libcamera/base/log.h>
> +#include <libcamera/base/signal.h>
> +
> +#include <libcamera/framebuffer.h>
> +
> +#include "sequence_sync_helper.h"
> +
> +namespace libcamera {
> +
> +LOG_DECLARE_CATEGORY(RkISP1Schedule)
> +
> +struct BufferQueueDelegateBase {
> +	virtual ~BufferQueueDelegateBase() = default;
> +	virtual int allocateBuffers(unsigned int count,
> +				    std::vector<std::unique_ptr<FrameBuffer>> *buffers) = 0;
> +	virtual int importBuffers(unsigned int count) = 0;
> +	virtual int releaseBuffers() = 0;
> +
> +	virtual int queueBuffer(FrameBuffer *buffer) = 0;
> +
> +	Signal<FrameBuffer *> bufferReady;
> +};
> +
> +template<typename T>
> +struct BufferQueueDelegate : public BufferQueueDelegateBase {

Would it make sense to call this `SimpleBufferQueueDelegate` or similar?
And keep the `BufferQueueDelegate` name for the interface?


> +	BufferQueueDelegate(T *target) : target_(target)

I don't think it's a big change, and I think it would make this type more
future-proof. You could take `T &&target` and move into `target_`. The type
could then support smart pointers without future changes. Or anything really
that implements `operator->()` (even `std::optional`...).


> +	{
> +		target_->bufferReady.connect(this, [this](FrameBuffer *buffer) {
> +			this->bufferReady.emit(buffer);
> +		});
> +	}
> +
> +	~BufferQueueDelegate()
> +	{
> +		target_->bufferReady.disconnect(this);
> +	}
> +
> +	int allocateBuffers(unsigned int count,
> +			    std::vector<std::unique_ptr<FrameBuffer>> *buffers) override
> +	{
> +		return target_->allocateBuffers(count, buffers);
> +	}
> +
> +	int importBuffers(unsigned int count) override
> +	{
> +		return target_->importBuffers(count);
> +	}
> +
> +	int queueBuffer(FrameBuffer *buffer) override
> +	{
> +		return target_->queueBuffer(buffer);
> +	}
> +
> +	int releaseBuffers() override
> +	{
> +		return target_->releaseBuffers();
> +	}
> +
> +private:
> +	T *target_;
> +};
> +
> +class BufferQueue
> +{
> +public:
> +	enum State {
> +		Idle = 0,
> +		Preparing,
> +		Capturing,
> +		Postprocessing
> +	};
> +
> +	enum Flags {
> +		PrepareStage = 1,
> +		PostprocessStage = 2

I think `1 << n` would be preferable for flags. But given that there are
only "stage flags", I think passing an instance of

   struct Stages {
     bool prepare = false;
     bool postprocess = false;
   };

to the constructor wouldn't be too bad either.


> +	};

I think `enum class` would be preferable for both.


> +
> +	BufferQueue(std::unique_ptr<BufferQueueDelegateBase> &&delegate, int flags = 0, std::string name = {});
> +
> +	int allocateBuffers(unsigned int count);
> +	int importBuffers(unsigned int count);
> +	int releaseBuffers();
> +
> +	int sequenceCorrection();
> +	uint32_t nextSequence();
> +
> +	int prepareBuffer(uint32_t *sequence = nullptr);
> +	int prepareBuffer(FrameBuffer *buffer, uint32_t *sequence = nullptr);
> +	int preparedBuffer();
> +
> +	int queueBuffer(uint32_t *sequence = nullptr);
> +	int queueBuffer(FrameBuffer *buffer, uint32_t *sequence = nullptr);
> +
> +	void postprocessedBuffer();
> +
> +	bool empty(State state);
> +
> +	FrameBuffer *front(State state);
> +
> +	unsigned int expectedSequence(FrameBuffer *buffer) const;
> +	const std::vector<std::unique_ptr<FrameBuffer>> &buffers() const;
> +
> +	Signal<FrameBuffer *> bufferReady;
> +
> +private:
> +	void onBufferReady(FrameBuffer *buffer);
> +
> +	int internalPrepareBuffer(FrameBuffer *buffer, uint32_t *sequence = nullptr);
> +	int internalPreparedBuffer();
> +	void internalPostprocessedBuffer();
> +
> +	std::map<State, std::list<FrameBuffer *>> bufferState_;

I think this should definitely be an `std::array<>`. Wrt. my earlier
`enum class` comment, there is the `RPi::Device` type that implements
essentially what I'm suggesting here.


> +	std::map<FrameBuffer *, unsigned int> expectedSequence_;
> +	std::vector<std::unique_ptr<FrameBuffer>> buffers_;
> +	SequenceSyncHelper syncHelper_;
> +	uint32_t nextSequence_;
> +	std::string name_;

The class should probably implement the `Loggable` interface
with `logPrefix()` being the name, instead of manually adding
the name to the log messages.


> +	bool ownsBuffers_;
> +	bool hasBuffers_;
> +	int flags_;

This should use `libcamera::Flags`, no?


> +	std::unique_ptr<BufferQueueDelegateBase> delegate_;
> +};
> +
> +} /* namespace libcamera */
> diff --git a/include/libcamera/internal/meson.build b/include/libcamera/internal/meson.build
> index 73635e31bad9..c44c57949382 100644
> --- a/include/libcamera/internal/meson.build
> +++ b/include/libcamera/internal/meson.build
> @@ -4,6 +4,7 @@ subdir('tracepoints')
>   
>   libcamera_internal_headers = files([
>       'bayer_format.h',
> +    'buffer_queue.h',
>       'byte_stream_buffer.h',
>       'camera.h',
>       'camera_controls.h',
> diff --git a/src/libcamera/buffer_queue.cpp b/src/libcamera/buffer_queue.cpp
> new file mode 100644
> index 000000000000..b92b633ca04d
> --- /dev/null
> +++ b/src/libcamera/buffer_queue.cpp
> @@ -0,0 +1,541 @@
> +/* SPDX-License-Identifier: LGPL-2.1-or-later */
> +/*
> + * Copyright (C) 2025, Ideas on Board
> + *
> + * BufferQueue implementation
> + */
> +
> +#include "libcamera/internal/buffer_queue.h"
> +
> +#include <libcamera/base/log.h>
> +
> +namespace libcamera {
> +
> +LOG_DEFINE_CATEGORY(BufferQueue)
> +
> +/**
> + * \struct BufferQueueDelegateBase
> + * \brief Abstract class that defines the delegate interface for the BufferQueue
> + *
> + * The BufferQueue needs to have access to the buffer handling functions of the
> + * object that wraps the underlying device. This is implemented by means of a
> + * delegate that is passed into the buffer queue at construction time. In most
> + * cases the \a BufferQueueDelegate template is sufficient.
> + *
> + * \fn BufferQueueDelegateBase::allocateBuffers
> + * \copydoc V4L2VideoDevice::allocateBuffers
> + *
> + * \fn BufferQueueDelegateBase::importBuffers
> + * \copydoc V4L2VideoDevice::importBuffers
> + *
> + * \fn BufferQueueDelegateBase::releaseBuffers
> + * \copydoc V4L2VideoDevice::releaseBuffers
> + *
> + * \fn BufferQueueDelegateBase::queueBuffer
> + * \brief Queue a buffer to the device
> + * \param[in] buffer The buffer to be queued
> + * Call queueBuffer on the underlying device.
> + *
> + * \see V4L2VideoDevice::queueBuffer
> + *
> + * \var BufferQueueDelegateBase::bufferReady
> + * \copydoc V4L2VideoDevice::bufferReady
> + *
> + * \struct BufferQueueDelegate
> + * \brief Template implementation of the BufferQueueDelegateBase interface
> + *
> + * This class implements the BufferQueueDelegateBase interface and forwards it
> + * to an object that provides the same functions (like V4L2VideoDevice).
> + *
> + * A pointer to the target object is passed in at contruction time. The user is

   construction
      ^


> + * responsible to ensure that the target object exists for the whole lifetime of
> + * the delegate.
> + *
> + * \fn BufferQueueDelegate::BufferQueueDelegate
> + * \param target The target to delegate to
> + *
> + * \class BufferQueue
> + * \brief Helper class to handle buffer queues
> + *
> + * Handling buffer queues is a common task when dealing with V4L2 video device.
> + * Depending on the specific use case, additional functionalities are required:
> + * - Supporting internally allocated and imported buffers
> + * - Estimate the sequence number that a buffer will have after dequeuing
> + * - Correct for errors in that sequence scheme due to frames beeing dropped in
> + *   the kernel.
> + * - Optionally adding a preparation stage for a buffer before it get's queued

   gets


> + *   to the device
> + * - Optionally adding a postprocessing stage after dequeueing
> + *
> + * This class encapsulates these functionalities in one component. To support
> + * arbitrary V4L2VideoDevice like classes, the actual access to the buffer
> + * related functions is handled through the BufferQueueDelegateBase interface
> + * with BufferQueueDelegate providing a default implementation that works with
> + * the V4L2VideDevice class.
> + *
> + * On construction time, it must be specified if a preparation stage and/or a
> + * postprocessing stage shall be added.
> + *
> + * Internally there are 4 queues, a buffer can be in:
> + * Idle, Preparing, Capturing, Postprocessing.
> + *
> + * After creation of the BufferQueue it is important to route all buffer related
> + * actions through the queue. That is:
> + * 1. Allocation/import of buffers
> + * 2. Queueing buffers
> + * 3. Releasing buffers
> + * 4. Handling the bufferReady signal.
> + *
> + * The callable functions depend on the flags passed to the constructor and are
> + * shown on the following state diagram.
> + *
> + *  *------*
> + *  | Idle |
> + *  *------*
> + *      |                        +----------------+
> + *      +-- prepareBuffer() ---->|  Preparing     |
> + *      |                        +----------------+
> + *      |                                 |
> + *   (no prepare stage)           preparedBuffer()
> + *	|		                  |
> + *	|		         +----------------+
> + *      \-- queueBuffer() ------>| Capturing      |
> + *                               +----------------+

I feel like "Capturing" might be a bit misleading given that
param buffers are not "captured into". Maybe "Queued" or something
similar.



> + *     /--- buffer was cancelled ------/  |
> + *     |                                  |
> + *     |                         +----------------+
> + *     |                         | Postprocessing | --> bufferReady signal
> + *     |                         +----------------+
> + *     |/--(no postprocess stage)------/  |
> + *     |                                  | postprocessedBuffer()
> + *     | /--------------------------------/
> + *  *------*
> + *  | Idle |
> + *  *------*
> + *
> + * Notes:
> + * - If buffers are not allocated by the queue but imported they never end up in
> + *   the idle queue but are passed in by prepareBuffer()/queueBuffer() and leave
> + *   the Queue after postprocessing.
> + * - If a preparing stage is used, queueBuffer can not be called.
> + * - If the postprocessing stage is disabled, it will still be used while
> + *   emitting the bufferReady signal but the transition to idle happens
> + *   automatically afterwards.
> + */
> +
> +/**
> + * \enum BufferQueue::State
> + * \brief The states a buffer can be in
> + *
> + * \var BufferQueue::Idle
> + * \brief Buffer is not queued
> + *
> + * \var BufferQueue::Preparing
> + * \brief Buffer is beeing prepared for queueing

   being


> + *
> + * \var BufferQueue::Capturing
> + * \brief Buffer is queued in
> + *
> + * \var BufferQueue::Postprocessing
> + * \brief Buffer beeing postprocessed

   being


> + */
> +
> +/**
> + * \enum BufferQueue::Flags
> + * \brief Flags for a BufferQueue
> + *
> + * \var BufferQueue::PrepareStage
> + * \brief The queue has a prepare stage
> + *
> + * \var BufferQueue::PostprocessStage
> + * \brief The queue has a postprocess stage
> + */
> +
> +/**
> + * \brief Construct a BufferQueue
> + * \param[in] delegate The delegate
> + * \param[in] flags Optional flags
> + * \param[in] name Optional name
> + *
> + * Construct a buffer queue using the given delegate to forward the buffer
> + * handling to. The default queue has only Idle and queued states. This can be
> + * changed using the \a flags parameter, to either add a prepare stage, a
> + * postprocessing stage or both.
> + */
> +BufferQueue::BufferQueue(std::unique_ptr<BufferQueueDelegateBase> &&delegate,
> +			 int flags, std::string name)
> +	: nextSequence_(0), name_(std::move(name)), ownsBuffers_(false),
> +	  hasBuffers_(false), flags_(flags), delegate_(std::move(delegate))
> +{
> +	delegate_->bufferReady.connect(this, &BufferQueue::onBufferReady);
> +}
> +
> +/**
> + * \brief Allocate buffers
> + * \param[in] count The number of buffers to allocate
> + *
> + * This function allocates \a count buffers by calling allocateBuffers() on the
> + * delegate and forwarding the return value. A non negative return code is
> + * treated as success. The buffers are owned by the BufferQueue.
> + *
> + * \return The value returned by the BufferQueueDelegateBase::allocateBuffers()
> + */
> +int BufferQueue::allocateBuffers(unsigned int count)
> +{
> +	ASSERT(!hasBuffers_);
> +	buffers_.clear();
> +	int ret = delegate_->allocateBuffers(count, &buffers_);
> +	if (ret < 0)
> +		return ret;
> +
> +	for (const auto &buffer : buffers_)
> +		bufferState_[Idle].push_back(buffer.get());
> +
> +	hasBuffers_ = true;
> +	ownsBuffers_ = true;
> +	return ret;
> +}
> +
> +/**
> + * \brief Import buffers
> + * \param[in] count The number of buffers to import
> + *
> + * This function imports \a count buffers by calling importBuffers() on the
> + * delegate and forwarding the return value. A non negative return code is
> + * treated as success.
> + *
> + * \return The value returned by the BufferQueueDelegateBase::importBuffers()
> + */
> +int BufferQueue::importBuffers(unsigned int count)
> +{
> +	int ret = delegate_->importBuffers(count);
> +	if (ret < 0)
> +		return ret;
> +
> +	hasBuffers_ = true;
> +	ownsBuffers_ = false;
> +	return 0;
> +}
> +
> +/**
> + * \brief Get the necessary correction
> + *
> + * \return The offset to add to get in sync again
> + */
> +int BufferQueue::sequenceCorrection()
> +{
> +	return syncHelper_.correction();
> +}
> +
> +/**
> + * \brief Get the seqence of the next buffer
> + *
> + * \return The sequence including necessary corrections
> + */
> +uint32_t BufferQueue::nextSequence()
> +{
> +	return nextSequence_ + syncHelper_.correction();
> +}

I'd probably inline the above two.


> +
> +/**
> + * \brief Move the next buffer to prepare state
> + * \param[out] sequence The expected sequence of the buffer
> + *
> + * This function moves the front buffer from the idle queue to prepare queue. If
> + * \a sequence is provided it is set to the expected sequence number of that
> + * buffer.
> + *
> + * \note This function must only be called if the queue has a prepare state and
> + * owns the buffers.
> + *
> + * \return 0 on success, a negative error code otherwise
> + */
> +int BufferQueue::prepareBuffer(uint32_t *sequence)
> +{
> +	ASSERT(hasBuffers_);
> +	ASSERT(ownsBuffers_);
> +	ASSERT(flags_ & PrepareStage);
> +	ASSERT(!bufferState_[Idle].empty());
> +
> +	FrameBuffer *buffer = bufferState_[Idle].front();
> +	return prepareBuffer(buffer, sequence);
> +}
> +
> +/**
> + * \brief Move a buffer to prepare state
> + * \param[in] buffer The buffer
> + * \param[out] sequence The expected sequence of the buffer
> + *
> + * This function moves \a buffer to prepare queue. If
> + * \a sequence is provided it is set to the expected sequence number of that
> + * buffer.
> + *
> + * \note This function must only be called if the queue has a prepare state. If
> + * the queue owns the buffer, \a buffer must point to the front buffer of the
> + * idle queue.
> + *
> + * \return 0 on success, a negative error code otherwise
> + */
> +int BufferQueue::prepareBuffer(FrameBuffer *buffer, uint32_t *sequence)
> +{
> +	ASSERT(flags_ & PrepareStage);
> +
> +	return internalPrepareBuffer(buffer, sequence);
> +}
> +
> +/**
> + * \brief Exit prepare state
> + *
> + * This function pops the frontmost buffer from the prepare queue and queues it
> + * on the underlying device by calling queueBuffer() on the delegate.
> + *
> + * \note This function must only be called if the queue has a prepare state.
> + *
> + * \return 0 on success, a negative error code otherwise
> + */
> +int BufferQueue::preparedBuffer()
> +{
> +	ASSERT(flags_ & PrepareStage);
> +
> +	return internalPreparedBuffer();
> +}
> +
> +/**
> + * \brief Queue the next buffer
> + * \param[out] sequence The expected sequence of the buffer
> + *
> + * This function queues the front buffer from the idle queue to the underlying
> + * device ba calling queueBuffer() on the delegate. If \a sequence is provided

   by


> + * it is set to the expected sequence number of that buffer.
> + *
> + * \note This function must only be called if the queue does not have a prepare
> + * state and owns the buffers.
> + *
> + * \return 0 on success, a negative error code otherwise
> + */
> +int BufferQueue::queueBuffer(uint32_t *sequence)
> +{
> +	ASSERT(hasBuffers_);
> +	ASSERT(ownsBuffers_);
> +	ASSERT(!bufferState_[Idle].empty());
> +
> +	FrameBuffer *buffer = bufferState_[Idle].front();
> +	return queueBuffer(buffer, sequence);
> +}
> +
> +/**
> + * \brief Queue a buffer
> + * \param[in] buffer The buffer
> + * \param[out] sequence The expected sequence of the buffer
> + *
> + * This function queues \a buffer to the underlying device ba calling
> + * queueBuffer() on the delegate. If \a sequence is provided it is set to the
> + * expected sequence number of that buffer.
> + *
> + * \note This function must only be called if the queue does not have a prepare
> + * state. If the queue owns the buffers, \a buffer must point to the front
> + * buffer of the idle queue.
> + *
> + * \return 0 on success, a negative error code otherwise
> + */
> +int BufferQueue::queueBuffer(FrameBuffer *buffer, uint32_t *sequence)
> +{
> +	ASSERT(!(flags_ & PrepareStage));
> +	return internalPrepareBuffer(buffer, sequence);
> +}
> +
> +/**
> + * \brief Exit postprocessed state
> + *
> + * This function pops the frontmost buffer from the postprocess queue and puts it
> + * back to the idle queue in case the buffers are owned by the queue
> + *
> + * \note This function must only be called if the queue has a prepare state.
> + */
> +void BufferQueue::postprocessedBuffer()
> +{
> +	ASSERT(hasBuffers_);
> +	ASSERT(flags_ & PostprocessStage);
> +	return internalPostprocessedBuffer();
> +}
> +
> +/**
> + * \brief Release buffers
> + *
> + * This function releases the allocated or imported buffers by calling
> + * releaseBuffers() on the delegate.
> + *
> + * \return 0 on success, a negative error code otherwise
> + */
> +int BufferQueue::releaseBuffers()
> +{
> +	ASSERT(bufferState_[BufferQueue::Idle].size() == buffers_.size());
> +	ASSERT(empty(Preparing) && empty(Capturing) && empty(Postprocessing));
> +
> +	bufferState_[BufferQueue::Idle] = {};

   .clear()

?


> +	buffers_.clear();
> +	hasBuffers_ = false;
> +	syncHelper_.reset();
> +	nextSequence_ = 0;
> +
> +	return delegate_->releaseBuffers();
> +}
> +
> +/**
> + * \brief Check if queue is empty
> + * \param[in] state The state
> + *
> + * \return True if the queue for the given state is empty, false otherwise
> + */
> +bool BufferQueue::empty(BufferQueue::State state)
> +{
> +	return bufferState_[state].empty();
> +}

I think this can be inlined.


> +
> +/**
> + * \brief Get the front buffer of a queue
> + * \param[in] state The state
> + *
> + * \return The front buffer of the queue, or null otherwise
> + */
> +FrameBuffer *BufferQueue::front(BufferQueue::State state)
> +{
> +	if (empty(state))
> +		return nullptr;
> +	return bufferState_[state].front();
> +}

Same here.



> +
> +/**
> + * \brief Get the expected sequence for a buffer
> + * \param[in] buffer The buffer
> + *
> + * \return The expected sequence
> + */
> +unsigned int BufferQueue::expectedSequence(FrameBuffer *buffer) const
> +{
> +	auto it = expectedSequence_.find(buffer);
> +	ASSERT(it != expectedSequence_.end());
> +	return it->second;
> +}
> +
> +/**
> + * \brief Get the allocated buffers
> + *
> + * \return The buffers owned by the queue
> + */
> +const std::vector<std::unique_ptr<FrameBuffer>> &BufferQueue::buffers() const
> +{
> +	return buffers_;
> +}

Same here.


> +
> +/**
> + * \var BufferQueue::bufferReady
> + * \brief A Signal emitted when a framebuffer completes
> + *
> + * When this signal is emitted the buffer will be in Postprocessing state. If
> + * the queue was constructed without a postprocessing stage, the buffer will
> + * automatically move to the idle state after the signal was emitted. Otherwise
> + * it will stay in postprocessing state until postprocessedBuffer() is called.
> + */
> +
> +void BufferQueue::onBufferReady(FrameBuffer *buffer)
> +{
> +	ASSERT(!empty(Capturing));
> +
> +	auto &meta = buffer->metadata();
> +	auto &queue = bufferState_[Capturing];
> +
> +	/*
> +         * V4L2 does not guarantee that buffers are dequeued in order. We expect
> +         * drivers to usually do so, and therefore warn, if a buffer is returned
> +         * out of order. After streamoff V4L2VideoDevice returns the buffers in
> +         * arbitrary order so there is no warning needed in that case.
> +         */
> +	auto it = std::find(queue.begin(), queue.end(), buffer);
> +	ASSERT(it != queue.end());
> +
> +	if (it != queue.begin() &&
> +	    meta.status != FrameMetadata::FrameCancelled)
> +		LOG(BufferQueue, Warning) << name_ << ": Dequeued buffer out of order " << buffer;
> +
> +	queue.erase(it);
> +	if (meta.status == FrameMetadata::FrameCancelled) {
> +		syncHelper_.cancelFrame();
> +		if (ownsBuffers_)
> +			bufferState_[Idle].push_back(buffer);
> +	} else {
> +		syncHelper_.receivedFrame(expectedSequence_[buffer], meta.sequence);
> +		bufferState_[Postprocessing].push_back(buffer);
> +	}
> +
> +	bufferReady.emit(buffer);
> +
> +	if (!(flags_ & PostprocessStage) &&
> +	    meta.status != FrameMetadata::FrameCancelled)
> +		internalPostprocessedBuffer();
> +}
> +
> +int BufferQueue::internalPrepareBuffer(FrameBuffer *buffer, uint32_t *sequence)
> +{
> +	ASSERT(hasBuffers_);
> +
> +	if (ownsBuffers_) {
> +		ASSERT(!bufferState_[Idle].empty());
> +		ASSERT(bufferState_[Idle].front() == buffer);
> +	}
> +
> +	LOG(BufferQueue, Debug) << name_ << ":Buffer prepare: "
> +				<< buffer;
> +	int correction = syncHelper_.correction();
> +	nextSequence_ += correction;
> +	expectedSequence_[buffer] = nextSequence_;
> +	if (ownsBuffers_)
> +		bufferState_[Idle].pop_front();
> +	bufferState_[Preparing].push_back(buffer);
> +	syncHelper_.pushCorrection(correction);
> +
> +	if (sequence)
> +		*sequence = nextSequence_;
> +
> +	nextSequence_++;
> +
> +	if (!(flags_ & PrepareStage))
> +		return internalPreparedBuffer();
> +
> +	return 0;
> +}
> +
> +int BufferQueue::internalPreparedBuffer()
> +{
> +	ASSERT(!bufferState_[Preparing].empty());
> +
> +	auto &srcQueue = bufferState_[Preparing];
> +	FrameBuffer *buffer = srcQueue.front();
> +
> +	srcQueue.pop_front();
> +	int ret = delegate_->queueBuffer(buffer);
> +	if (ret < 0) {
> +		LOG(BufferQueue, Error) << "Failed to queue buffer: "
> +					<< strerror(-ret);
> +		if (ownsBuffers_)
> +			bufferState_[Idle].push_back(buffer);

Shouldn't this do something with `expectedSequence_`,  `syncCorrection_`, `nextSequence_`, etc.?
What is the imagined way to (try to) recover from this?

On a related note, where are elements removed from `expectedSequence_` ?


> +		return ret;
> +	}
> +
> +	LOG(BufferQueue, Debug) << name_ << " Queued buffer: " << buffer;
> +
> +	bufferState_[Capturing].push_back(buffer);
> +	return 0;
> +}
> +
> +void BufferQueue::internalPostprocessedBuffer()
> +{
> +	ASSERT(!empty(Postprocessing));
> +
> +	FrameBuffer *buffer = bufferState_[Postprocessing].front();
> +	bufferState_[Postprocessing].pop_front();
> +	if (ownsBuffers_)
> +		bufferState_[Idle].push_back(buffer);

Sidenote, but it's a bit unfortunate how many allocations this list jumping produces.
There is `std::list::splice` to avoid some of that, but it's not the most convenient
function to use.


> +}
> +
> +} /* namespace libcamera */
> diff --git a/src/libcamera/meson.build b/src/libcamera/meson.build
> index 038d2dbf12cf..1ea97ae997bf 100644
> --- a/src/libcamera/meson.build
> +++ b/src/libcamera/meson.build
> @@ -18,6 +18,7 @@ libcamera_public_sources = files([
>   
>   libcamera_internal_sources = files([
>       'bayer_format.cpp',
> +    'buffer_queue.cpp',
>       'byte_stream_buffer.cpp',
>       'camera_controls.cpp',
>       'camera_lens.cpp',

Patch
diff mbox series

diff --git a/include/libcamera/internal/buffer_queue.h b/include/libcamera/internal/buffer_queue.h
new file mode 100644
index 000000000000..83a6a1f45be7
--- /dev/null
+++ b/include/libcamera/internal/buffer_queue.h
@@ -0,0 +1,136 @@ 
+/* SPDX-License-Identifier: LGPL-2.1-or-later */
+/*
+ * Copyright (C) 2025, Ideas on Board
+ *
+ * Sequence sync helper
+ */
+
+#pragma once
+
+#include <map>
+#include <memory>
+
+#include <libcamera/base/log.h>
+#include <libcamera/base/signal.h>
+
+#include <libcamera/framebuffer.h>
+
+#include "sequence_sync_helper.h"
+
+namespace libcamera {
+
+LOG_DECLARE_CATEGORY(RkISP1Schedule)
+
+struct BufferQueueDelegateBase {
+	virtual ~BufferQueueDelegateBase() = default;
+	virtual int allocateBuffers(unsigned int count,
+				    std::vector<std::unique_ptr<FrameBuffer>> *buffers) = 0;
+	virtual int importBuffers(unsigned int count) = 0;
+	virtual int releaseBuffers() = 0;
+
+	virtual int queueBuffer(FrameBuffer *buffer) = 0;
+
+	Signal<FrameBuffer *> bufferReady;
+};
+
+template<typename T>
+struct BufferQueueDelegate : public BufferQueueDelegateBase {
+	BufferQueueDelegate(T *target) : target_(target)
+	{
+		target_->bufferReady.connect(this, [this](FrameBuffer *buffer) {
+			this->bufferReady.emit(buffer);
+		});
+	}
+
+	~BufferQueueDelegate()
+	{
+		target_->bufferReady.disconnect(this);
+	}
+
+	int allocateBuffers(unsigned int count,
+			    std::vector<std::unique_ptr<FrameBuffer>> *buffers) override
+	{
+		return target_->allocateBuffers(count, buffers);
+	}
+
+	int importBuffers(unsigned int count) override
+	{
+		return target_->importBuffers(count);
+	}
+
+	int queueBuffer(FrameBuffer *buffer) override
+	{
+		return target_->queueBuffer(buffer);
+	}
+
+	int releaseBuffers() override
+	{
+		return target_->releaseBuffers();
+	}
+
+private:
+	T *target_;
+};
+
+class BufferQueue
+{
+public:
+	enum State {
+		Idle = 0,
+		Preparing,
+		Capturing,
+		Postprocessing
+	};
+
+	enum Flags {
+		PrepareStage = 1,
+		PostprocessStage = 2
+	};
+
+	BufferQueue(std::unique_ptr<BufferQueueDelegateBase> &&delegate, int flags = 0, std::string name = {});
+
+	int allocateBuffers(unsigned int count);
+	int importBuffers(unsigned int count);
+	int releaseBuffers();
+
+	int sequenceCorrection();
+	uint32_t nextSequence();
+
+	int prepareBuffer(uint32_t *sequence = nullptr);
+	int prepareBuffer(FrameBuffer *buffer, uint32_t *sequence = nullptr);
+	int preparedBuffer();
+
+	int queueBuffer(uint32_t *sequence = nullptr);
+	int queueBuffer(FrameBuffer *buffer, uint32_t *sequence = nullptr);
+
+	void postprocessedBuffer();
+
+	bool empty(State state);
+
+	FrameBuffer *front(State state);
+
+	unsigned int expectedSequence(FrameBuffer *buffer) const;
+	const std::vector<std::unique_ptr<FrameBuffer>> &buffers() const;
+
+	Signal<FrameBuffer *> bufferReady;
+
+private:
+	void onBufferReady(FrameBuffer *buffer);
+
+	int internalPrepareBuffer(FrameBuffer *buffer, uint32_t *sequence = nullptr);
+	int internalPreparedBuffer();
+	void internalPostprocessedBuffer();
+
+	std::map<State, std::list<FrameBuffer *>> bufferState_;
+	std::map<FrameBuffer *, unsigned int> expectedSequence_;
+	std::vector<std::unique_ptr<FrameBuffer>> buffers_;
+	SequenceSyncHelper syncHelper_;
+	uint32_t nextSequence_;
+	std::string name_;
+	bool ownsBuffers_;
+	bool hasBuffers_;
+	int flags_;
+	std::unique_ptr<BufferQueueDelegateBase> delegate_;
+};
+
+} /* namespace libcamera */
diff --git a/include/libcamera/internal/meson.build b/include/libcamera/internal/meson.build
index 73635e31bad9..c44c57949382 100644
--- a/include/libcamera/internal/meson.build
+++ b/include/libcamera/internal/meson.build
@@ -4,6 +4,7 @@  subdir('tracepoints')
 
 libcamera_internal_headers = files([
     'bayer_format.h',
+    'buffer_queue.h',
     'byte_stream_buffer.h',
     'camera.h',
     'camera_controls.h',
diff --git a/src/libcamera/buffer_queue.cpp b/src/libcamera/buffer_queue.cpp
new file mode 100644
index 000000000000..b92b633ca04d
--- /dev/null
+++ b/src/libcamera/buffer_queue.cpp
@@ -0,0 +1,541 @@ 
+/* SPDX-License-Identifier: LGPL-2.1-or-later */
+/*
+ * Copyright (C) 2025, Ideas on Board
+ *
+ * BufferQueue implementation
+ */
+
+#include "libcamera/internal/buffer_queue.h"
+
+#include <libcamera/base/log.h>
+
+namespace libcamera {
+
+LOG_DEFINE_CATEGORY(BufferQueue)
+
+/**
+ * \struct BufferQueueDelegateBase
+ * \brief Abstract class that defines the delegate interface for the BufferQueue
+ *
+ * The BufferQueue needs to have access to the buffer handling functions of the
+ * object that wraps the underlying device. This is implemented by means of a
+ * delegate that is passed into the buffer queue at construction time. In most
+ * cases the \a BufferQueueDelegate template is sufficient.
+ *
+ * \fn BufferQueueDelegateBase::allocateBuffers
+ * \copydoc V4L2VideoDevice::allocateBuffers
+ *
+ * \fn BufferQueueDelegateBase::importBuffers
+ * \copydoc V4L2VideoDevice::importBuffers
+ *
+ * \fn BufferQueueDelegateBase::releaseBuffers
+ * \copydoc V4L2VideoDevice::releaseBuffers
+ *
+ * \fn BufferQueueDelegateBase::queueBuffer
+ * \brief Queue a buffer to the device
+ * \param[in] buffer The buffer to be queued
+ * Call queueBuffer on the underlying device.
+ *
+ * \see V4L2VideoDevice::queueBuffer
+ *
+ * \var BufferQueueDelegateBase::bufferReady
+ * \copydoc V4L2VideoDevice::bufferReady
+ *
+ * \struct BufferQueueDelegate
+ * \brief Template implementation of the BufferQueueDelegateBase interface
+ *
+ * This class implements the BufferQueueDelegateBase interface and forwards it
+ * to an object that provides the same functions (like V4L2VideoDevice).
+ *
+ * A pointer to the target object is passed in at contruction time. The user is
+ * responsible to ensure that the target object exists for the whole lifetime of
+ * the delegate.
+ *
+ * \fn BufferQueueDelegate::BufferQueueDelegate
+ * \param target The target to delegate to
+ *
+ * \class BufferQueue
+ * \brief Helper class to handle buffer queues
+ *
+ * Handling buffer queues is a common task when dealing with V4L2 video device.
+ * Depending on the specific use case, additional functionalities are required:
+ * - Supporting internally allocated and imported buffers
+ * - Estimate the sequence number that a buffer will have after dequeuing
+ * - Correct for errors in that sequence scheme due to frames beeing dropped in
+ *   the kernel.
+ * - Optionally adding a preparation stage for a buffer before it get's queued
+ *   to the device
+ * - Optionally adding a postprocessing stage after dequeueing
+ *
+ * This class encapsulates these functionalities in one component. To support
+ * arbitrary V4L2VideoDevice like classes, the actual access to the buffer
+ * related functions is handled through the BufferQueueDelegateBase interface
+ * with BufferQueueDelegate providing a default implementation that works with
+ * the V4L2VideDevice class.
+ *
+ * On construction time, it must be specified if a preparation stage and/or a
+ * postprocessing stage shall be added.
+ *
+ * Internally there are 4 queues, a buffer can be in:
+ * Idle, Preparing, Capturing, Postprocessing.
+ *
+ * After creation of the BufferQueue it is important to route all buffer related
+ * actions through the queue. That is:
+ * 1. Allocation/import of buffers
+ * 2. Queueing buffers
+ * 3. Releasing buffers
+ * 4. Handling the bufferReady signal.
+ *
+ * The callable functions depend on the flags passed to the constructor and are
+ * shown on the following state diagram.
+ *
+ *  *------*
+ *  | Idle |
+ *  *------*
+ *      |                        +----------------+
+ *      +-- prepareBuffer() ---->|  Preparing     |
+ *      |                        +----------------+
+ *      |                                 |
+ *   (no prepare stage)           preparedBuffer()
+ *	|		                  |
+ *	|		         +----------------+
+ *      \-- queueBuffer() ------>| Capturing      |
+ *                               +----------------+
+ *     /--- buffer was cancelled ------/  |
+ *     |                                  |
+ *     |                         +----------------+
+ *     |                         | Postprocessing | --> bufferReady signal
+ *     |                         +----------------+
+ *     |/--(no postprocess stage)------/  |
+ *     |                                  | postprocessedBuffer()
+ *     | /--------------------------------/
+ *  *------*
+ *  | Idle |
+ *  *------*
+ *
+ * Notes:
+ * - If buffers are not allocated by the queue but imported they never end up in
+ *   the idle queue but are passed in by prepareBuffer()/queueBuffer() and leave
+ *   the Queue after postprocessing.
+ * - If a preparing stage is used, queueBuffer can not be called.
+ * - If the postprocessing stage is disabled, it will still be used while
+ *   emitting the bufferReady signal but the transition to idle happens
+ *   automatically afterwards.
+ */
+
+/**
+ * \enum BufferQueue::State
+ * \brief The states a buffer can be in
+ *
+ * \var BufferQueue::Idle
+ * \brief Buffer is not queued
+ *
+ * \var BufferQueue::Preparing
+ * \brief Buffer is beeing prepared for queueing
+ *
+ * \var BufferQueue::Capturing
+ * \brief Buffer is queued in
+ *
+ * \var BufferQueue::Postprocessing
+ * \brief Buffer beeing postprocessed
+ */
+
+/**
+ * \enum BufferQueue::Flags
+ * \brief Flags for a BufferQueue
+ *
+ * \var BufferQueue::PrepareStage
+ * \brief The queue has a prepare stage
+ *
+ * \var BufferQueue::PostprocessStage
+ * \brief The queue has a postprocess stage
+ */
+
+/**
+ * \brief Construct a BufferQueue
+ * \param[in] delegate The delegate
+ * \param[in] flags Optional flags
+ * \param[in] name Optional name
+ *
+ * Construct a buffer queue using the given delegate to forward the buffer
+ * handling to. The default queue has only Idle and queued states. This can be
+ * changed using the \a flags parameter, to either add a prepare stage, a
+ * postprocessing stage or both.
+ */
+BufferQueue::BufferQueue(std::unique_ptr<BufferQueueDelegateBase> &&delegate,
+			 int flags, std::string name)
+	: nextSequence_(0), name_(std::move(name)), ownsBuffers_(false),
+	  hasBuffers_(false), flags_(flags), delegate_(std::move(delegate))
+{
+	delegate_->bufferReady.connect(this, &BufferQueue::onBufferReady);
+}
+
+/**
+ * \brief Allocate buffers
+ * \param[in] count The number of buffers to allocate
+ *
+ * This function allocates \a count buffers by calling allocateBuffers() on the
+ * delegate and forwarding the return value. A non negative return code is
+ * treated as success. The buffers are owned by the BufferQueue.
+ *
+ * \return The value returned by the BufferQueueDelegateBase::allocateBuffers()
+ */
+int BufferQueue::allocateBuffers(unsigned int count)
+{
+	ASSERT(!hasBuffers_);
+	buffers_.clear();
+	int ret = delegate_->allocateBuffers(count, &buffers_);
+	if (ret < 0)
+		return ret;
+
+	for (const auto &buffer : buffers_)
+		bufferState_[Idle].push_back(buffer.get());
+
+	hasBuffers_ = true;
+	ownsBuffers_ = true;
+	return ret;
+}
+
+/**
+ * \brief Import buffers
+ * \param[in] count The number of buffers to import
+ *
+ * This function imports \a count buffers by calling importBuffers() on the
+ * delegate and forwarding the return value. A non negative return code is
+ * treated as success.
+ *
+ * \return The value returned by the BufferQueueDelegateBase::importBuffers()
+ */
+int BufferQueue::importBuffers(unsigned int count)
+{
+	int ret = delegate_->importBuffers(count);
+	if (ret < 0)
+		return ret;
+
+	hasBuffers_ = true;
+	ownsBuffers_ = false;
+	return 0;
+}
+
+/**
+ * \brief Get the necessary correction
+ *
+ * \return The offset to add to get in sync again
+ */
+int BufferQueue::sequenceCorrection()
+{
+	return syncHelper_.correction();
+}
+
+/**
+ * \brief Get the seqence of the next buffer
+ *
+ * \return The sequence including necessary corrections
+ */
+uint32_t BufferQueue::nextSequence()
+{
+	return nextSequence_ + syncHelper_.correction();
+}
+
+/**
+ * \brief Move the next buffer to prepare state
+ * \param[out] sequence The expected sequence of the buffer
+ *
+ * This function moves the front buffer from the idle queue to prepare queue. If
+ * \a sequence is provided it is set to the expected sequence number of that
+ * buffer.
+ *
+ * \note This function must only be called if the queue has a prepare state and
+ * owns the buffers.
+ *
+ * \return 0 on success, a negative error code otherwise
+ */
+int BufferQueue::prepareBuffer(uint32_t *sequence)
+{
+	ASSERT(hasBuffers_);
+	ASSERT(ownsBuffers_);
+	ASSERT(flags_ & PrepareStage);
+	ASSERT(!bufferState_[Idle].empty());
+
+	FrameBuffer *buffer = bufferState_[Idle].front();
+	return prepareBuffer(buffer, sequence);
+}
+
+/**
+ * \brief Move a buffer to prepare state
+ * \param[in] buffer The buffer
+ * \param[out] sequence The expected sequence of the buffer
+ *
+ * This function moves \a buffer to prepare queue. If
+ * \a sequence is provided it is set to the expected sequence number of that
+ * buffer.
+ *
+ * \note This function must only be called if the queue has a prepare state. If
+ * the queue owns the buffer, \a buffer must point to the front buffer of the
+ * idle queue.
+ *
+ * \return 0 on success, a negative error code otherwise
+ */
+int BufferQueue::prepareBuffer(FrameBuffer *buffer, uint32_t *sequence)
+{
+	ASSERT(flags_ & PrepareStage);
+
+	return internalPrepareBuffer(buffer, sequence);
+}
+
+/**
+ * \brief Exit prepare state
+ *
+ * This function pops the frontmost buffer from the prepare queue and queues it
+ * on the underlying device by calling queueBuffer() on the delegate.
+ *
+ * \note This function must only be called if the queue has a prepare state.
+ *
+ * \return 0 on success, a negative error code otherwise
+ */
+int BufferQueue::preparedBuffer()
+{
+	ASSERT(flags_ & PrepareStage);
+
+	return internalPreparedBuffer();
+}
+
+/**
+ * \brief Queue the next buffer
+ * \param[out] sequence The expected sequence of the buffer
+ *
+ * This function queues the front buffer from the idle queue to the underlying
+ * device ba calling queueBuffer() on the delegate. If \a sequence is provided
+ * it is set to the expected sequence number of that buffer.
+ *
+ * \note This function must only be called if the queue does not have a prepare
+ * state and owns the buffers.
+ *
+ * \return 0 on success, a negative error code otherwise
+ */
+int BufferQueue::queueBuffer(uint32_t *sequence)
+{
+	ASSERT(hasBuffers_);
+	ASSERT(ownsBuffers_);
+	ASSERT(!bufferState_[Idle].empty());
+
+	FrameBuffer *buffer = bufferState_[Idle].front();
+	return queueBuffer(buffer, sequence);
+}
+
+/**
+ * \brief Queue a buffer
+ * \param[in] buffer The buffer
+ * \param[out] sequence The expected sequence of the buffer
+ *
+ * This function queues \a buffer to the underlying device ba calling
+ * queueBuffer() on the delegate. If \a sequence is provided it is set to the
+ * expected sequence number of that buffer.
+ *
+ * \note This function must only be called if the queue does not have a prepare
+ * state. If the queue owns the buffers, \a buffer must point to the front
+ * buffer of the idle queue.
+ *
+ * \return 0 on success, a negative error code otherwise
+ */
+int BufferQueue::queueBuffer(FrameBuffer *buffer, uint32_t *sequence)
+{
+	ASSERT(!(flags_ & PrepareStage));
+	return internalPrepareBuffer(buffer, sequence);
+}
+
+/**
+ * \brief Exit postprocessed state
+ *
+ * This function pops the frontmost buffer from the postprocess queue and puts it
+ * back to the idle queue in case the buffers are owned by the queue
+ *
+ * \note This function must only be called if the queue has a prepare state.
+ */
+void BufferQueue::postprocessedBuffer()
+{
+	ASSERT(hasBuffers_);
+	ASSERT(flags_ & PostprocessStage);
+	return internalPostprocessedBuffer();
+}
+
+/**
+ * \brief Release buffers
+ *
+ * This function releases the allocated or imported buffers by calling
+ * releaseBuffers() on the delegate.
+ *
+ * \return 0 on success, a negative error code otherwise
+ */
+int BufferQueue::releaseBuffers()
+{
+	ASSERT(bufferState_[BufferQueue::Idle].size() == buffers_.size());
+	ASSERT(empty(Preparing) && empty(Capturing) && empty(Postprocessing));
+
+	bufferState_[BufferQueue::Idle] = {};
+	buffers_.clear();
+	hasBuffers_ = false;
+	syncHelper_.reset();
+	nextSequence_ = 0;
+
+	return delegate_->releaseBuffers();
+}
+
+/**
+ * \brief Check if queue is empty
+ * \param[in] state The state
+ *
+ * \return True if the queue for the given state is empty, false otherwise
+ */
+bool BufferQueue::empty(BufferQueue::State state)
+{
+	return bufferState_[state].empty();
+}
+
+/**
+ * \brief Get the front buffer of a queue
+ * \param[in] state The state
+ *
+ * \return The front buffer of the queue, or null otherwise
+ */
+FrameBuffer *BufferQueue::front(BufferQueue::State state)
+{
+	if (empty(state))
+		return nullptr;
+	return bufferState_[state].front();
+}
+
+/**
+ * \brief Get the expected sequence for a buffer
+ * \param[in] buffer The buffer
+ *
+ * \return The expected sequence
+ */
+unsigned int BufferQueue::expectedSequence(FrameBuffer *buffer) const
+{
+	auto it = expectedSequence_.find(buffer);
+	ASSERT(it != expectedSequence_.end());
+	return it->second;
+}
+
+/**
+ * \brief Get the allocated buffers
+ *
+ * \return The buffers owned by the queue
+ */
+const std::vector<std::unique_ptr<FrameBuffer>> &BufferQueue::buffers() const
+{
+	return buffers_;
+}
+
+/**
+ * \var BufferQueue::bufferReady
+ * \brief A Signal emitted when a framebuffer completes
+ *
+ * When this signal is emitted the buffer will be in Postprocessing state. If
+ * the queue was constructed without a postprocessing stage, the buffer will
+ * automatically move to the idle state after the signal was emitted. Otherwise
+ * it will stay in postprocessing state until postprocessedBuffer() is called.
+ */
+
+void BufferQueue::onBufferReady(FrameBuffer *buffer)
+{
+	ASSERT(!empty(Capturing));
+
+	auto &meta = buffer->metadata();
+	auto &queue = bufferState_[Capturing];
+
+	/*
+         * V4L2 does not guarantee that buffers are dequeued in order. We expect
+         * drivers to usually do so, and therefore warn, if a buffer is returned
+         * out of order. After streamoff V4L2VideoDevice returns the buffers in
+         * arbitrary order so there is no warning needed in that case.
+         */
+	auto it = std::find(queue.begin(), queue.end(), buffer);
+	ASSERT(it != queue.end());
+
+	if (it != queue.begin() &&
+	    meta.status != FrameMetadata::FrameCancelled)
+		LOG(BufferQueue, Warning) << name_ << ": Dequeued buffer out of order " << buffer;
+
+	queue.erase(it);
+	if (meta.status == FrameMetadata::FrameCancelled) {
+		syncHelper_.cancelFrame();
+		if (ownsBuffers_)
+			bufferState_[Idle].push_back(buffer);
+	} else {
+		syncHelper_.receivedFrame(expectedSequence_[buffer], meta.sequence);
+		bufferState_[Postprocessing].push_back(buffer);
+	}
+
+	bufferReady.emit(buffer);
+
+	if (!(flags_ & PostprocessStage) &&
+	    meta.status != FrameMetadata::FrameCancelled)
+		internalPostprocessedBuffer();
+}
+
+int BufferQueue::internalPrepareBuffer(FrameBuffer *buffer, uint32_t *sequence)
+{
+	ASSERT(hasBuffers_);
+
+	if (ownsBuffers_) {
+		ASSERT(!bufferState_[Idle].empty());
+		ASSERT(bufferState_[Idle].front() == buffer);
+	}
+
+	LOG(BufferQueue, Debug) << name_ << ":Buffer prepare: "
+				<< buffer;
+	int correction = syncHelper_.correction();
+	nextSequence_ += correction;
+	expectedSequence_[buffer] = nextSequence_;
+	if (ownsBuffers_)
+		bufferState_[Idle].pop_front();
+	bufferState_[Preparing].push_back(buffer);
+	syncHelper_.pushCorrection(correction);
+
+	if (sequence)
+		*sequence = nextSequence_;
+
+	nextSequence_++;
+
+	if (!(flags_ & PrepareStage))
+		return internalPreparedBuffer();
+
+	return 0;
+}
+
+int BufferQueue::internalPreparedBuffer()
+{
+	ASSERT(!bufferState_[Preparing].empty());
+
+	auto &srcQueue = bufferState_[Preparing];
+	FrameBuffer *buffer = srcQueue.front();
+
+	srcQueue.pop_front();
+	int ret = delegate_->queueBuffer(buffer);
+	if (ret < 0) {
+		LOG(BufferQueue, Error) << "Failed to queue buffer: "
+					<< strerror(-ret);
+		if (ownsBuffers_)
+			bufferState_[Idle].push_back(buffer);
+		return ret;
+	}
+
+	LOG(BufferQueue, Debug) << name_ << " Queued buffer: " << buffer;
+
+	bufferState_[Capturing].push_back(buffer);
+	return 0;
+}
+
+void BufferQueue::internalPostprocessedBuffer()
+{
+	ASSERT(!empty(Postprocessing));
+
+	FrameBuffer *buffer = bufferState_[Postprocessing].front();
+	bufferState_[Postprocessing].pop_front();
+	if (ownsBuffers_)
+		bufferState_[Idle].push_back(buffer);
+}
+
+} /* namespace libcamera */
diff --git a/src/libcamera/meson.build b/src/libcamera/meson.build
index 038d2dbf12cf..1ea97ae997bf 100644
--- a/src/libcamera/meson.build
+++ b/src/libcamera/meson.build
@@ -18,6 +18,7 @@  libcamera_public_sources = files([
 
 libcamera_internal_sources = files([
     'bayer_format.cpp',
+    'buffer_queue.cpp',
     'byte_stream_buffer.cpp',
     'camera_controls.cpp',
     'camera_lens.cpp',