#define BT_COMP_LOG_SELF_COMP (muxer_comp->self_comp)
#define BT_LOG_OUTPUT_LEVEL (muxer_comp->log_level)
#define BT_LOG_TAG "PLUGIN/FLT.UTILS.MUXER"
-#include "plugins/comp-logging.h"
+#include "logging/comp-logging.h"
#include "common/macros.h"
#include "common/uuid.h"
#include <stdlib.h>
#include <string.h>
+#include "plugins/common/muxing/muxing.h"
+
#include "muxer.h"
#define ASSUME_ABSOLUTE_CLOCK_CLASSES_PARAM_NAME "assume-absolute-clock-classes"
struct muxer_msg_iter {
struct muxer_comp *muxer_comp;
+ /* Weak */
+ bt_self_message_iterator *self_msg_iter;
+
/*
* Array of struct muxer_upstream_msg_iter * (owned by this).
*
}
static
-bt_self_component_port_input_message_iterator *
+bt_self_component_port_input_message_iterator_create_from_message_iterator_status
create_msg_iter_on_input_port(struct muxer_comp *muxer_comp,
- bt_self_component_port_input *self_port)
+ struct muxer_msg_iter *muxer_msg_iter,
+ bt_self_component_port_input *self_port,
+ bt_self_component_port_input_message_iterator **msg_iter)
{
const bt_port *port = bt_self_component_port_as_port(
bt_self_component_port_input_as_self_component_port(
self_port));
- bt_self_component_port_input_message_iterator *msg_iter =
- NULL;
+ bt_self_component_port_input_message_iterator_create_from_message_iterator_status
+ status;
BT_ASSERT(port);
BT_ASSERT(bt_port_is_connected(port));
// TODO: Advance the iterator to >= the time of the latest
// returned message by the muxer message
// iterator which creates it.
- msg_iter = bt_self_component_port_input_message_iterator_create(
- self_port);
- if (!msg_iter) {
+ status = bt_self_component_port_input_message_iterator_create_from_message_iterator(
+ muxer_msg_iter->self_msg_iter, self_port, msg_iter);
+ if (status != BT_SELF_COMPONENT_PORT_INPUT_MESSAGE_ITERATOR_CREATE_FROM_MESSAGE_ITERATOR_STATUS_OK) {
BT_COMP_LOGE("Cannot create upstream message iterator on input port: "
"port-addr=%p, port-name=\"%s\"",
port, bt_port_get_name(port));
port, bt_port_get_name(port), msg_iter);
end:
- return msg_iter;
+ return status;
}
static
/* Expect no clock class */
muxer_msg_iter->clock_class_expectation =
MUXER_MSG_ITER_CLOCK_CLASS_EXPECTATION_NONE;
- } else {
- BT_COMP_LOGE("Expecting stream class with a default clock class: "
+ } else if (muxer_msg_iter->clock_class_expectation !=
+ MUXER_MSG_ITER_CLOCK_CLASS_EXPECTATION_NONE) {
+ BT_COMP_LOGE("Expecting stream class without a default clock class: "
"stream-class-addr=%p, stream-class-name=\"%s\", "
"stream-class-id=%" PRIu64,
stream_class, bt_stream_class_get_name(stream_class),
goto end;
}
- if (msg_ts_ns <= youngest_ts_ns) {
+ /*
+ * Update the current message iterator if it has not been set
+ * yet, or if its current message has a timestamp smaller than
+ * the previously selected youngest message.
+ */
+ if (G_UNLIKELY(*muxer_upstream_msg_iter == NULL) ||
+ msg_ts_ns < youngest_ts_ns) {
*muxer_upstream_msg_iter =
cur_muxer_upstream_msg_iter;
youngest_ts_ns = msg_ts_ns;
*ts_ns = youngest_ts_ns;
+ } else if (msg_ts_ns == youngest_ts_ns) {
+ /*
+ * The currently selected message to be sent downstream
+ * next has the exact same timestamp that of the
+ * current candidate message. We must break the tie
+ * in a predictable manner.
+ */
+ const bt_message *selected_msg = g_queue_peek_head(
+ (*muxer_upstream_msg_iter)->msgs);
+ BT_COMP_LOGD_STR("Two of the next message candidates have the same timestamps, pick one deterministically.");
+
+ /*
+ * Order the messages in an arbitrary but determinitic
+ * way.
+ */
+ ret = common_muxing_compare_messages(msg, selected_msg);
+ if (ret < 0) {
+ /*
+ * The `msg` should go first. Update the next
+ * iterator and the current timestamp.
+ */
+ *muxer_upstream_msg_iter =
+ cur_muxer_upstream_msg_iter;
+ youngest_ts_ns = msg_ts_ns;
+ *ts_ns = youngest_ts_ns;
+ } else if (ret == 0) {
+ /* Unable to pick which one should go first. */
+ BT_COMP_LOGW("Cannot deterministically pick next upstream message iterator because they have identical next messages: "
+ "muxer-upstream-msg-iter-wrap-addr=%p"
+ "cur-muxer-upstream-msg-iter-wrap-addr=%p",
+ *muxer_upstream_msg_iter,
+ cur_muxer_upstream_msg_iter);
+ }
}
}
}
static
-int muxer_msg_iter_init_upstream_iterators(struct muxer_comp *muxer_comp,
+bt_component_class_message_iterator_init_method_status
+muxer_msg_iter_init_upstream_iterators(struct muxer_comp *muxer_comp,
struct muxer_msg_iter *muxer_msg_iter)
{
int64_t count;
int64_t i;
- int ret = 0;
+ bt_component_class_message_iterator_init_method_status status;
count = bt_component_filter_get_input_port_count(
bt_self_component_filter_as_component_filter(
BT_COMP_LOGD("No input port to initialize for muxer component's message iterator: "
"muxer-comp-addr=%p, muxer-msg-iter-addr=%p",
muxer_comp, muxer_msg_iter);
+ status = BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_OK;
goto end;
}
bt_self_component_filter_borrow_input_port_by_index(
muxer_comp->self_comp_flt, i);
const bt_port *port;
+ bt_self_component_port_input_message_iterator_create_from_message_iterator_status
+ msg_iter_status;
+ int int_status;
BT_ASSERT(self_port);
port = bt_self_component_port_as_port(
continue;
}
- upstream_msg_iter = create_msg_iter_on_input_port(muxer_comp,
- self_port);
- if (!upstream_msg_iter) {
+ msg_iter_status = create_msg_iter_on_input_port(muxer_comp,
+ muxer_msg_iter, self_port, &upstream_msg_iter);
+ if (msg_iter_status != BT_SELF_COMPONENT_PORT_INPUT_MESSAGE_ITERATOR_CREATE_FROM_MESSAGE_ITERATOR_STATUS_OK) {
/* create_msg_iter_on_input_port() logs errors */
- BT_ASSERT(!upstream_msg_iter);
- ret = -1;
+ status = (int) msg_iter_status;
goto end;
}
- ret = muxer_msg_iter_add_upstream_msg_iter(muxer_msg_iter,
+ int_status = muxer_msg_iter_add_upstream_msg_iter(muxer_msg_iter,
upstream_msg_iter);
bt_self_component_port_input_message_iterator_put_ref(
upstream_msg_iter);
- if (ret) {
+ if (int_status) {
+ status = BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_ERROR;
/* muxer_msg_iter_add_upstream_msg_iter() logs errors */
goto end;
}
}
+ status = BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_OK;
+
end:
- return ret;
+ return status;
}
BT_HIDDEN
{
struct muxer_comp *muxer_comp = NULL;
struct muxer_msg_iter *muxer_msg_iter = NULL;
- bt_component_class_message_iterator_init_method_status status =
- BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_OK;
- int ret;
+ bt_component_class_message_iterator_init_method_status status;
muxer_comp = bt_self_component_get_data(
bt_self_component_filter_as_self_component(self_comp));
BT_COMP_LOGE("Recursive initialization of muxer component's message iterator: "
"comp-addr=%p, muxer-comp-addr=%p, msg-iter-addr=%p",
self_comp, muxer_comp, self_msg_iter);
+ status = BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_ERROR;
goto error;
}
muxer_msg_iter = g_new0(struct muxer_msg_iter, 1);
if (!muxer_msg_iter) {
BT_COMP_LOGE_STR("Failed to allocate one muxer component's message iterator.");
+ status = BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_MEMORY_ERROR;
goto error;
}
muxer_msg_iter->muxer_comp = muxer_comp;
+ muxer_msg_iter->self_msg_iter = self_msg_iter;
muxer_msg_iter->last_returned_ts_ns = INT64_MIN;
muxer_msg_iter->active_muxer_upstream_msg_iters =
g_ptr_array_new_with_free_func(
(GDestroyNotify) destroy_muxer_upstream_msg_iter);
if (!muxer_msg_iter->active_muxer_upstream_msg_iters) {
BT_COMP_LOGE_STR("Failed to allocate a GPtrArray.");
+ status = BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_MEMORY_ERROR;
goto error;
}
(GDestroyNotify) destroy_muxer_upstream_msg_iter);
if (!muxer_msg_iter->ended_muxer_upstream_msg_iters) {
BT_COMP_LOGE_STR("Failed to allocate a GPtrArray.");
+ status = BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_MEMORY_ERROR;
goto error;
}
- ret = muxer_msg_iter_init_upstream_iterators(muxer_comp,
+ status = muxer_msg_iter_init_upstream_iterators(muxer_comp,
muxer_msg_iter);
- if (ret) {
+ if (status) {
BT_COMP_LOGE("Cannot initialize connected input ports for muxer component's message iterator: "
"comp-addr=%p, muxer-comp-addr=%p, "
"muxer-msg-iter-addr=%p, msg-iter-addr=%p, ret=%d",
self_comp, muxer_comp, muxer_msg_iter,
- self_msg_iter, ret);
+ self_msg_iter, status);
goto error;
}
error:
destroy_muxer_msg_iter(muxer_msg_iter);
bt_self_message_iterator_set_data(self_msg_iter, NULL);
- status = BT_COMPONENT_CLASS_MESSAGE_ITERATOR_INIT_METHOD_STATUS_ERROR;
end:
muxer_comp->initializing_muxer_msg_iter = false;
}
static inline
-bt_bool muxer_upstream_msg_iters_can_all_seek_beginning(
- GPtrArray *muxer_upstream_msg_iters)
+bt_component_class_message_iterator_can_seek_beginning_method_status
+muxer_upstream_msg_iters_can_all_seek_beginning(
+ GPtrArray *muxer_upstream_msg_iters, bt_bool *can_seek)
{
+ bt_component_class_message_iterator_can_seek_beginning_method_status status;
uint64_t i;
- bt_bool ret = BT_TRUE;
for (i = 0; i < muxer_upstream_msg_iters->len; i++) {
struct muxer_upstream_msg_iter *upstream_msg_iter =
muxer_upstream_msg_iters->pdata[i];
+ status = (int) bt_self_component_port_input_message_iterator_can_seek_beginning(
+ upstream_msg_iter->msg_iter, can_seek);
+ if (status != BT_COMPONENT_CLASS_MESSAGE_ITERATOR_CAN_SEEK_BEGINNING_METHOD_STATUS_OK) {
+ goto end;
+ }
- if (!bt_self_component_port_input_message_iterator_can_seek_beginning(
- upstream_msg_iter->msg_iter)) {
- ret = BT_FALSE;
+ if (!*can_seek) {
goto end;
}
}
+ *can_seek = BT_TRUE;
+
end:
- return ret;
+ return status;
}
BT_HIDDEN
-bt_bool muxer_msg_iter_can_seek_beginning(
- bt_self_message_iterator *self_msg_iter)
+bt_component_class_message_iterator_can_seek_beginning_method_status
+muxer_msg_iter_can_seek_beginning(
+ bt_self_message_iterator *self_msg_iter, bt_bool *can_seek)
{
struct muxer_msg_iter *muxer_msg_iter =
bt_self_message_iterator_get_data(self_msg_iter);
- bt_bool ret = BT_TRUE;
+ bt_component_class_message_iterator_can_seek_beginning_method_status status;
- if (!muxer_upstream_msg_iters_can_all_seek_beginning(
- muxer_msg_iter->active_muxer_upstream_msg_iters)) {
- ret = BT_FALSE;
+ status = muxer_upstream_msg_iters_can_all_seek_beginning(
+ muxer_msg_iter->active_muxer_upstream_msg_iters, can_seek);
+ if (status != BT_COMPONENT_CLASS_MESSAGE_ITERATOR_CAN_SEEK_BEGINNING_METHOD_STATUS_OK) {
goto end;
}
- if (!muxer_upstream_msg_iters_can_all_seek_beginning(
- muxer_msg_iter->ended_muxer_upstream_msg_iters)) {
- ret = BT_FALSE;
+ if (!*can_seek) {
goto end;
}
+ status = muxer_upstream_msg_iters_can_all_seek_beginning(
+ muxer_msg_iter->ended_muxer_upstream_msg_iters, can_seek);
+
end:
- return ret;
+ return status;
}
BT_HIDDEN
muxer_msg_iter->ended_muxer_upstream_msg_iters->pdata[i] = NULL;
}
- g_ptr_array_remove_range(muxer_msg_iter->ended_muxer_upstream_msg_iters,
- 0, muxer_msg_iter->ended_muxer_upstream_msg_iters->len);
+ /*
+ * GLib < 2.48.0 asserts when g_ptr_array_remove_range() is
+ * called on an empty array.
+ */
+ if (muxer_msg_iter->ended_muxer_upstream_msg_iters->len > 0) {
+ g_ptr_array_remove_range(muxer_msg_iter->ended_muxer_upstream_msg_iters,
+ 0, muxer_msg_iter->ended_muxer_upstream_msg_iters->len);
+ }
muxer_msg_iter->last_returned_ts_ns = INT64_MIN;
muxer_msg_iter->clock_class_expectation =
MUXER_MSG_ITER_CLOCK_CLASS_EXPECTATION_ANY;