Work on sink

Remove _remove from properties, we can do the same with set of a NULL
value.
Add signals to the stream API to manage the buffers. Wrap those buffers
in a GstBuffer in the pinossrc and pinossink elements and pool them in a
bufferpool.
Remove SPA_EVENT_TYPE_PULL_INPUT, we can do the same with NEED_INPUT and
by using a ringbuffer.
Do more complete allocation of buffers in the link. Use the buffer
allocator if none of the nodes can allocate.
Follow the node state to trigger negotiation and allocation.
Remove offset and size when refering to buffers, we want to always deal
with the complete buffer and use a ringbuffer for ranges or change the
offset/size in the buffer data when needed.
Serialize port_info structures as part of the port_update
Print both the enum number and the name when debuging properties or
formats.
This commit is contained in:
Wim Taymans 2016-08-24 16:26:58 +02:00
parent a03352353f
commit ca7d08c406
45 changed files with 1614 additions and 570 deletions

View file

@ -239,7 +239,8 @@ libgstpinos_la_SOURCES = \
gst/gstpinosformat.c \
gst/gstpinosdeviceprovider.c \
gst/gstpinossrc.c \
gst/gstpinossink.c
gst/gstpinossink.c \
gst/gstpinospool.c
#gst/gstpinosdepay.c
#gst/gstpinosportsrc.c
#gst/gstpinosportsink.c

View file

@ -131,7 +131,8 @@ pinos_properties_free (PinosProperties *properties)
* @value: a value
*
* Set the property in @properties with @key to @value. Any previous value
* of @key will be overwritten.
* of @key will be overwritten. When @value is %NULL, the key will be
* removed.
*/
void
pinos_properties_set (PinosProperties *properties,
@ -140,9 +141,11 @@ pinos_properties_set (PinosProperties *properties,
{
g_return_if_fail (properties != NULL);
g_return_if_fail (key != NULL);
g_return_if_fail (value != NULL);
g_hash_table_replace (properties->hashtable, g_strdup (key), g_strdup (value));
if (value == NULL)
g_hash_table_remove (properties->hashtable, key);
else
g_hash_table_replace (properties->hashtable, g_strdup (key), g_strdup (value));
}
/**
@ -193,23 +196,6 @@ pinos_properties_get (PinosProperties *properties,
return g_hash_table_lookup (properties->hashtable, key);
}
/**
* pinos_properties_remove:
* @properties: a #PinosProperties
* @key: a key
*
* Remove the property in @properties with @key.
*/
void
pinos_properties_remove (PinosProperties *properties,
const gchar *key)
{
g_return_if_fail (properties != NULL);
g_return_if_fail (key != NULL);
g_hash_table_remove (properties->hashtable, key);
}
/**
* pinos_properties_iterate:
* @properties: a #PinosProperties

View file

@ -44,8 +44,6 @@ void pinos_properties_setf (PinosProperties *properties,
...) G_GNUC_PRINTF (3, 4);
const gchar * pinos_properties_get (PinosProperties *properties,
const gchar *key);
void pinos_properties_remove (PinosProperties *properties,
const gchar *key);
const gchar * pinos_properties_iterate (PinosProperties *properties,
gpointer *state);

View file

@ -71,8 +71,6 @@ enum
LAST_SIGNAL
};
static guint signals[LAST_SIGNAL] = { 0 };
static void
pinos_ringbuffer_get_property (GObject *_object,
guint prop_id,
@ -367,7 +365,8 @@ pinos_ringbuffer_read_advance (PinosRingbuffer *rbuf,
if (priv->mode == PINOS_RINGBUFFER_MODE_READ) {
val = 1;
write (priv->semaphore, &val, 8);
if (write (priv->semaphore, &val, 8) != 8)
g_warning ("error writing semaphore");
}
return TRUE;
@ -387,7 +386,8 @@ pinos_ringbuffer_write_advance (PinosRingbuffer *rbuf,
if (priv->mode == PINOS_RINGBUFFER_MODE_WRITE) {
val = 1;
write (priv->semaphore, &val, 8);
if (write (priv->semaphore, &val, 8) != 8)
g_warning ("error writing semaphore");
}
return TRUE;
}

View file

@ -37,7 +37,7 @@
#include "pinos/client/format.h"
#include "pinos/client/private.h"
#define MAX_BUFFER_SIZE 1024
#define MAX_BUFFER_SIZE 4096
#define MAX_FDS 16
typedef struct {
@ -46,6 +46,7 @@ typedef struct {
int fd;
off_t offset;
size_t size;
bool used;
SpaBuffer *buf;
} BufferId;
@ -75,6 +76,9 @@ struct _PinosStreamPrivate
GPtrArray *possible_formats;
SpaFormat *format;
SpaPortInfo port_info;
SpaAllocParam *port_params[2];
SpaAllocParamMetaEnable param_meta_enable;
SpaAllocParamBuffers param_buffers;
PinosStreamFlags flags;
@ -86,8 +90,6 @@ struct _PinosStreamPrivate
GSource *socket_source;
int fd;
SpaBuffer *buffer;
SpaControl *control;
SpaControl recv_control;
guint8 recv_data[MAX_BUFFER_SIZE];
@ -97,6 +99,7 @@ struct _PinosStreamPrivate
int send_fds[MAX_FDS];
GArray *buffer_ids;
gboolean in_order;
};
#define PINOS_STREAM_GET_PRIVATE(obj) \
@ -117,6 +120,8 @@ enum
enum
{
SIGNAL_ADD_BUFFER,
SIGNAL_REMOVE_BUFFER,
SIGNAL_NEW_BUFFER,
LAST_SIGNAL
};
@ -394,10 +399,41 @@ pinos_stream_class_init (PinosStreamClass * klass)
G_PARAM_STATIC_STRINGS));
/**
* PinosStream:new-buffer
* @id: the buffer id
*
* When doing pinos_stream_start() with #PINOS_STREAM_MODE_BUFFER, this signal
* will be fired whenever a new buffer can be obtained with
* pinos_stream_capture_buffer().
* this signal will be fired whenever a buffer is added to the pool of buffers.
*/
signals[SIGNAL_ADD_BUFFER] = g_signal_new ("add-buffer",
G_TYPE_FROM_CLASS (klass),
G_SIGNAL_RUN_LAST,
0,
NULL,
NULL,
g_cclosure_marshal_generic,
G_TYPE_NONE,
1,
G_TYPE_UINT);
/**
* PinosStream:remove-buffer
* @id: the buffer id
*
* this signal will be fired whenever a buffer is removed from the pool of buffers.
*/
signals[SIGNAL_REMOVE_BUFFER] = g_signal_new ("remove-buffer",
G_TYPE_FROM_CLASS (klass),
G_SIGNAL_RUN_LAST,
0,
NULL,
NULL,
g_cclosure_marshal_generic,
G_TYPE_NONE,
1,
G_TYPE_UINT);
/**
* PinosStream:new-buffer
* @id: the buffer id
*
* this signal will be fired whenever a buffer is ready to be processed.
*/
signals[SIGNAL_NEW_BUFFER] = g_signal_new ("new-buffer",
G_TYPE_FROM_CLASS (klass),
@ -407,8 +443,8 @@ pinos_stream_class_init (PinosStreamClass * klass)
NULL,
g_cclosure_marshal_generic,
G_TYPE_NONE,
0,
G_TYPE_NONE);
1,
G_TYPE_UINT);
}
static void
@ -421,6 +457,7 @@ pinos_stream_init (PinosStream * stream)
priv->state = PINOS_STREAM_STATE_UNCONNECTED;
priv->buffer_ids = g_array_sized_new (FALSE, FALSE, sizeof (BufferId), 64);
g_array_set_clear_func (priv->buffer_ids, (GDestroyNotify) clear_buffer_id);
priv->in_order = TRUE;
}
/**
@ -524,24 +561,110 @@ control_builder_init (PinosStream *stream, SpaControlBuilder *builder)
}
static void
send_need_input (PinosStream *stream, uint32_t port_id, uint32_t buffer_id)
add_node_update (PinosStream *stream, SpaControlBuilder *builder, uint32_t change_mask)
{
PinosStreamPrivate *priv = stream->priv;
SpaControlCmdNodeUpdate nu = { 0, };
nu.change_mask = change_mask;
if (change_mask & SPA_CONTROL_CMD_NODE_UPDATE_MAX_INPUTS)
nu.max_input_ports = priv->direction == PINOS_DIRECTION_INPUT ? 1 : 0;
if (change_mask & SPA_CONTROL_CMD_NODE_UPDATE_MAX_OUTPUTS)
nu.max_output_ports = priv->direction == PINOS_DIRECTION_OUTPUT ? 1 : 0;
nu.props = NULL;
spa_control_builder_add_cmd (builder, SPA_CONTROL_CMD_NODE_UPDATE, &nu);
}
static void
add_port_update (PinosStream *stream, SpaControlBuilder *builder, uint32_t change_mask)
{
PinosStreamPrivate *priv = stream->priv;
SpaControlCmdPortUpdate pu = { 0, };;
pu.port_id = 0;
pu.change_mask = change_mask;
if (change_mask & SPA_CONTROL_CMD_PORT_UPDATE_DIRECTION)
pu.direction = priv->direction;
if (change_mask & SPA_CONTROL_CMD_PORT_UPDATE_POSSIBLE_FORMATS) {
pu.n_possible_formats = priv->possible_formats->len;
pu.possible_formats = (SpaFormat **)priv->possible_formats->pdata;
}
pu.props = NULL;
if (change_mask & SPA_CONTROL_CMD_PORT_UPDATE_INFO)
pu.info = &priv->port_info;
spa_control_builder_add_cmd (builder, SPA_CONTROL_CMD_PORT_UPDATE, &pu);
}
static void
add_state_change (PinosStream *stream, SpaControlBuilder *builder, SpaNodeState state)
{
SpaControlCmdStateChange sc;
sc.state = state;
spa_control_builder_add_cmd (builder, SPA_CONTROL_CMD_STATE_CHANGE, &sc);
}
static void
add_need_input (PinosStream *stream, SpaControlBuilder *builder, uint32_t port_id)
{
SpaControlCmdNeedInput ni;
ni.port_id = port_id;
spa_control_builder_add_cmd (builder, SPA_CONTROL_CMD_NEED_INPUT, &ni);
}
static void
send_need_input (PinosStream *stream, uint32_t port_id)
{
PinosStreamPrivate *priv = stream->priv;
SpaControlBuilder builder;
SpaControl control;
control_builder_init (stream, &builder);
add_need_input (stream, &builder, port_id);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
spa_control_clear (&control);
}
static void
send_reuse_buffer (PinosStream *stream, uint32_t port_id, uint32_t buffer_id)
{
PinosStreamPrivate *priv = stream->priv;
SpaControlBuilder builder;
SpaControl control;
SpaControlCmdNeedInput ni;
SpaControlCmdReuseBuffer rb;
control_builder_init (stream, &builder);
if (buffer_id != SPA_ID_INVALID) {
rb.port_id = port_id;
rb.buffer_id = buffer_id;
rb.offset = 0;
rb.size = -1;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_REUSE_BUFFER, &rb);
}
ni.port_id = port_id;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_NEED_INPUT, &ni);
rb.port_id = port_id;
rb.buffer_id = buffer_id;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_REUSE_BUFFER, &rb);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
spa_control_clear (&control);
}
static void
send_process_buffer (PinosStream *stream, uint32_t port_id, uint32_t buffer_id)
{
PinosStreamPrivate *priv = stream->priv;
SpaControlBuilder builder;
SpaControl control;
SpaControlCmdProcessBuffer pb;
SpaControlCmdHaveOutput ho;
control_builder_init (stream, &builder);
pb.port_id = port_id;
pb.buffer_id = buffer_id;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_PROCESS_BUFFER, &pb);
ho.port_id = port_id;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_HAVE_OUTPUT, &ho);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
@ -555,10 +678,15 @@ find_buffer (PinosStream *stream, uint32_t id)
{
PinosStreamPrivate *priv = stream->priv;
guint i;
for (i = 0; i < priv->buffer_ids->len; i++) {
BufferId *bid = &g_array_index (priv->buffer_ids, BufferId, i);
if (bid->id == id)
return bid;
if (priv->in_order && id < priv->buffer_ids->len) {
return &g_array_index (priv->buffer_ids, BufferId, id);
} else {
for (i = 0; i < priv->buffer_ids->len; i++) {
BufferId *bid = &g_array_index (priv->buffer_ids, BufferId, i);
if (bid->id == id)
return bid;
}
}
return NULL;
}
@ -593,9 +721,6 @@ parse_control (PinosStream *stream,
case SPA_CONTROL_CMD_SET_FORMAT:
{
SpaControlCmdSetFormat p;
SpaControlBuilder builder;
SpaControl control;
SpaControlCmdStateChange sc;
if (spa_control_iter_parse_cmd (&it, &p) < 0)
break;
@ -607,19 +732,19 @@ parse_control (PinosStream *stream,
spa_debug_format (p.format);
g_object_notify (G_OBJECT (stream), "format");
control_builder_init (stream, &builder);
if (priv->port_info.n_params != 0) {
SpaControlBuilder builder;
SpaControl control;
/* FIXME send update port status */
control_builder_init (stream, &builder);
add_state_change (stream, &builder, SPA_NODE_STATE_READY);
spa_control_builder_end (&builder, &control);
/* send state-change */
sc.state = SPA_NODE_STATE_READY;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_STATE_CHANGE, &sc);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
spa_control_clear (&control);
spa_control_clear (&control);
}
break;
}
case SPA_CONTROL_CMD_SET_PROPERTY:
@ -628,21 +753,44 @@ parse_control (PinosStream *stream,
case SPA_CONTROL_CMD_START:
{
SpaControlBuilder builder;
SpaControl control;
g_debug ("stream %p: start", stream);
control_builder_init (stream, &builder);
if (priv->direction == PINOS_DIRECTION_INPUT)
send_need_input (stream, 0, SPA_ID_INVALID);
add_need_input (stream, &builder, 0);
add_state_change (stream, &builder, SPA_NODE_STATE_STREAMING);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
spa_control_clear (&control);
stream_set_state (stream, PINOS_STREAM_STATE_STREAMING, NULL);
break;
}
case SPA_CONTROL_CMD_STOP:
{
SpaControlBuilder builder;
SpaControl control;
g_debug ("stream %p: stop", stream);
control_builder_init (stream, &builder);
add_state_change (stream, &builder, SPA_NODE_STATE_PAUSED);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
spa_control_clear (&control);
stream_set_state (stream, PINOS_STREAM_STATE_READY, NULL);
break;
}
case SPA_CONTROL_CMD_ADD_MEM:
{
SpaControlCmdAddMem p;
@ -697,7 +845,12 @@ parse_control (PinosStream *stream,
bid.size = p.mem.size;
bid.buf = SPA_MEMBER (spa_memory_ensure_ptr (mem), p.mem.offset, SpaBuffer);
if (bid.id != priv->buffer_ids->len) {
g_warning ("unexpected id %u found, expected %u", bid.id, priv->buffer_ids->len);
priv->in_order = FALSE;
}
g_array_append_val (priv->buffer_ids, bid);
g_signal_emit (stream, signals[SIGNAL_ADD_BUFFER], 0, p.buffer_id);
break;
}
case SPA_CONTROL_CMD_REMOVE_BUFFER:
@ -709,30 +862,44 @@ parse_control (PinosStream *stream,
break;
g_debug ("remove buffer %d", p.buffer_id);
if ((bid = find_buffer (stream, p.buffer_id)))
bid->cleanup = true;
if ((bid = find_buffer (stream, p.buffer_id))) {
bid->cleanup = TRUE;
bid->used = TRUE;
g_signal_emit (stream, signals[SIGNAL_REMOVE_BUFFER], 0, p.buffer_id);
}
break;
}
case SPA_CONTROL_CMD_PROCESS_BUFFER:
{
SpaControlCmdProcessBuffer p;
BufferId *bid;
if (priv->direction != PINOS_DIRECTION_INPUT)
break;
if (spa_control_iter_parse_cmd (&it, &p) < 0)
break;
if ((bid = find_buffer (stream, p.buffer_id)))
priv->buffer = bid->buf;
g_signal_emit (stream, signals[SIGNAL_NEW_BUFFER], 0, p.buffer_id);
send_need_input (stream, 0);
break;
}
case SPA_CONTROL_CMD_REUSE_BUFFER:
{
SpaControlCmdReuseBuffer p;
BufferId *bid;
if (priv->direction != PINOS_DIRECTION_OUTPUT)
break;
if (spa_control_iter_parse_cmd (&it, &p) < 0)
break;
g_debug ("reuse buffer %d", p.buffer_id);
if ((bid = find_buffer (stream, p.buffer_id))) {
bid->used = FALSE;
g_signal_emit (stream, signals[SIGNAL_NEW_BUFFER], 0, p.buffer_id);
}
break;
}
@ -772,18 +939,17 @@ on_socket_condition (GSocket *socket,
parse_control (stream, control);
if (priv->buffer) {
g_signal_emit (stream, signals[SIGNAL_NEW_BUFFER], 0, NULL);
send_need_input (stream, 0, priv->buffer->id);
priv->buffer = NULL;
}
for (i = 0; i < priv->buffer_ids->len; i++) {
BufferId *bid = &g_array_index (priv->buffer_ids, BufferId, i);
if (bid->cleanup) {
g_array_remove_index_fast (priv->buffer_ids, i);
i--;
priv->in_order = FALSE;
}
}
if (!priv->in_order && priv->buffer_ids->len == 0)
priv->in_order = TRUE;
spa_control_clear (control);
break;
}
@ -835,45 +1001,6 @@ unhandle_socket (PinosStream *stream)
}
}
static void
do_node_init (PinosStream *stream)
{
PinosStreamPrivate *priv = stream->priv;
SpaControlCmdNodeUpdate nu;
SpaControlCmdPortUpdate pu;
SpaControlBuilder builder;
SpaControl control;
control_builder_init (stream, &builder);
nu.change_mask = SPA_CONTROL_CMD_NODE_UPDATE_MAX_INPUTS |
SPA_CONTROL_CMD_NODE_UPDATE_MAX_OUTPUTS;
nu.max_input_ports = priv->direction == PINOS_DIRECTION_INPUT ? 1 : 0;
nu.max_output_ports = priv->direction == PINOS_DIRECTION_OUTPUT ? 1 : 0;
nu.props = NULL;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_NODE_UPDATE, &nu);
pu.port_id = 0;
pu.change_mask = SPA_CONTROL_CMD_PORT_UPDATE_DIRECTION |
SPA_CONTROL_CMD_PORT_UPDATE_POSSIBLE_FORMATS |
SPA_CONTROL_CMD_PORT_UPDATE_INFO;
pu.direction = priv->direction;
pu.n_possible_formats = priv->possible_formats->len;
pu.possible_formats = (SpaFormat **)priv->possible_formats->pdata;
pu.props = NULL;
pu.info = &priv->port_info;
priv->port_info.flags = SPA_PORT_INFO_FLAG_NONE |
SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_PORT_UPDATE, &pu);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
spa_control_clear (&control);
}
static void
on_node_proxy (GObject *source_object,
GAsyncResult *res,
@ -882,6 +1009,9 @@ on_node_proxy (GObject *source_object,
PinosStream *stream = user_data;
PinosStreamPrivate *priv = stream->priv;
PinosContext *context = priv->context;
SpaControlBuilder builder;
SpaControl control;
GError *error = NULL;
priv->node = pinos_subscribe_get_proxy_finish (context->priv->subscribe,
@ -890,7 +1020,23 @@ on_node_proxy (GObject *source_object,
if (priv->node == NULL)
goto node_failed;
do_node_init (stream);
control_builder_init (stream, &builder);
add_node_update (stream, &builder, SPA_CONTROL_CMD_NODE_UPDATE_MAX_INPUTS |
SPA_CONTROL_CMD_NODE_UPDATE_MAX_OUTPUTS);
priv->port_info.flags = SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS;
add_port_update (stream, &builder, SPA_CONTROL_CMD_PORT_UPDATE_DIRECTION |
SPA_CONTROL_CMD_PORT_UPDATE_POSSIBLE_FORMATS |
SPA_CONTROL_CMD_PORT_UPDATE_INFO);
add_state_change (stream, &builder, SPA_NODE_STATE_CONFIGURE);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
spa_control_clear (&control);
stream_set_state (stream, PINOS_STREAM_STATE_READY, NULL);
g_object_unref (stream);
@ -978,8 +1124,9 @@ do_connect (PinosStream *stream)
if (priv->properties == NULL)
priv->properties = pinos_properties_new (NULL, NULL);
pinos_properties_set (priv->properties,
"pinos.target.node", priv->path);
if (priv->path)
pinos_properties_set (priv->properties,
"pinos.target.node", priv->path);
g_dbus_proxy_call (context->priv->daemon,
"CreateClientNode",
@ -1048,17 +1195,68 @@ pinos_stream_connect (PinosStream *stream,
return TRUE;
}
/**
* pinos_stream_start_allocation:
* @stream: a #PinosStream
* @props: a #PinosProperties
*
* Returns: %TRUE on success
*/
gboolean
pinos_stream_start_allocation (PinosStream *stream,
PinosProperties *props)
{
PinosStreamPrivate *priv;
PinosContext *context;
SpaControlBuilder builder;
SpaControl control;
g_return_val_if_fail (PINOS_IS_STREAM (stream), FALSE);
priv = stream->priv;
context = priv->context;
g_return_val_if_fail (pinos_context_get_state (context) == PINOS_CONTEXT_STATE_CONNECTED, FALSE);
control_builder_init (stream, &builder);
priv->port_info.params = priv->port_params;
priv->port_info.n_params = 1;
priv->port_params[0] = &priv->param_buffers.param;
priv->param_buffers.param.type = SPA_ALLOC_PARAM_TYPE_BUFFERS;
priv->param_buffers.param.size = sizeof (SpaAllocParamBuffers);
priv->param_buffers.minsize = 115200;
priv->param_buffers.stride = 640;
priv->param_buffers.min_buffers = 0;
priv->param_buffers.max_buffers = 0;
priv->param_buffers.align = 16;
/* send update port status */
add_port_update (stream, &builder, SPA_CONTROL_CMD_PORT_UPDATE_INFO);
/* send state-change */
if (priv->format)
add_state_change (stream, &builder, SPA_NODE_STATE_READY);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
g_warning ("stream %p: error writing control", stream);
spa_control_clear (&control);
return TRUE;
}
static gboolean
do_start (PinosStream *stream)
{
PinosStreamPrivate *priv = stream->priv;
SpaControlBuilder builder;
SpaControlCmdStateChange sc;
SpaControl control;
control_builder_init (stream, &builder);
sc.state = SPA_NODE_STATE_CONFIGURE;
spa_control_builder_add_cmd (&builder, SPA_CONTROL_CMD_STATE_CHANGE, &sc);
add_state_change (stream, &builder, SPA_NODE_STATE_CONFIGURE);
spa_control_builder_end (&builder, &control);
if (spa_control_write (&control, priv->fd) < 0)
@ -1223,46 +1421,111 @@ pinos_stream_disconnect (PinosStream *stream)
}
/**
* pinos_stream_peek_buffer:
* pinos_stream_get_empty_buffer:
* @stream: a #PinosStream
*
* Get the current buffer from @stream. This function should be called from
* Get the id of an empty buffer that can be filled
*
* Returns: the id of an empty buffer or #SPA_ID_INVALID when no buffer is
* available.
*/
guint
pinos_stream_get_empty_buffer (PinosStream *stream)
{
PinosStreamPrivate *priv;
guint i;
g_return_val_if_fail (PINOS_IS_STREAM (stream), FALSE);
priv = stream->priv;
g_return_val_if_fail (priv->direction == PINOS_DIRECTION_OUTPUT, FALSE);
for (i = 0; i < priv->buffer_ids->len; i++) {
BufferId *bid = &g_array_index (priv->buffer_ids, BufferId, i);
if (!bid->used)
return bid->id;
}
return SPA_ID_INVALID;
}
/**
* pinos_stream_recycle_buffer:
* @stream: a #PinosStream
* @id: a buffer id
*
* Recycle the buffer with @id.
*
* Returns: %TRUE on success.
*/
gboolean
pinos_stream_recycle_buffer (PinosStream *stream,
guint id)
{
PinosStreamPrivate *priv;
g_return_val_if_fail (PINOS_IS_STREAM (stream), FALSE);
g_return_val_if_fail (id != SPA_ID_INVALID, FALSE);
priv = stream->priv;
g_return_val_if_fail (priv->direction == PINOS_DIRECTION_INPUT, FALSE);
send_reuse_buffer (stream, 0, id);
return TRUE;
}
/**
* pinos_stream_peek_buffer:
* @stream: a #PinosStream
* @id: the buffer id
*
* Get the buffer with @id from @stream. This function should be called from
* the new-buffer signal callback.
*
* Returns: a #SpaBuffer or %NULL when there is no buffer.
*/
SpaBuffer *
pinos_stream_peek_buffer (PinosStream *stream)
pinos_stream_peek_buffer (PinosStream *stream, guint id)
{
PinosStreamPrivate *priv;
BufferId *bid;
g_return_val_if_fail (PINOS_IS_STREAM (stream), NULL);
priv = stream->priv;
return priv->buffer;
if ((bid = find_buffer (stream, id)))
return bid->buf;
return NULL;
}
/**
* pinos_stream_send_buffer:
* @stream: a #PinosStream
* @buffer: a #SpaBuffer
* @id: a buffer id
* @offset: the offset in the buffer
* @size: the size in the buffer
*
* Send a buffer to @stream.
* Send a buffer with @id to @stream.
*
* For provider streams, this function should be called whenever there is a new frame
* available.
*
* For capture streams, this functions should be called for each fd-payload that
* should be released.
*
* Returns: %TRUE when @buffer was handled
* Returns: %TRUE when @id was handled
*/
gboolean
pinos_stream_send_buffer (PinosStream *stream,
SpaBuffer *buffer)
pinos_stream_send_buffer (PinosStream *stream,
guint id)
{
g_return_val_if_fail (PINOS_IS_STREAM (stream), FALSE);
g_return_val_if_fail (buffer != NULL, FALSE);
PinosStreamPrivate *priv;
BufferId *bid;
return TRUE;
g_return_val_if_fail (PINOS_IS_STREAM (stream), FALSE);
g_return_val_if_fail (id != SPA_ID_INVALID, FALSE);
priv = stream->priv;
g_return_val_if_fail (priv->direction == PINOS_DIRECTION_OUTPUT, FALSE);
if ((bid = find_buffer (stream, id))) {
bid->used = TRUE;
send_process_buffer (stream, 0, id);
return TRUE;
} else {
return FALSE;
}
}

View file

@ -103,12 +103,19 @@ gboolean pinos_stream_connect (PinosStream *stream,
GPtrArray *possible_formats);
gboolean pinos_stream_disconnect (PinosStream *stream);
gboolean pinos_stream_start_allocation (PinosStream *stream,
PinosProperties *props);
gboolean pinos_stream_start (PinosStream *stream);
gboolean pinos_stream_stop (PinosStream *stream);
SpaBuffer * pinos_stream_peek_buffer (PinosStream *stream);
guint pinos_stream_get_empty_buffer (PinosStream *stream);
gboolean pinos_stream_recycle_buffer (PinosStream *stream,
guint id);
SpaBuffer * pinos_stream_peek_buffer (PinosStream *stream,
guint id);
gboolean pinos_stream_send_buffer (PinosStream *stream,
SpaBuffer *buffer);
guint id);
G_END_DECLS
#endif /* __PINOS_STREAM_H__ */

View file

@ -141,12 +141,5 @@
<method name='Remove'/>
<method name='GetRingbuffer'>
<arg type='a{sv}' name='properties' direction='in'/>
<arg type='h' name='buffermem' direction='out'/>
<arg type='u' name='buffersize' direction='out'/>
<arg type='h' name='fd' direction='out'/>
</method>
</interface>
</node>

149
pinos/gst/gstpinospool.c Normal file
View file

@ -0,0 +1,149 @@
/* GStreamer
* Copyright (C) 2016 Wim Taymans <wim.taymans@gmail.com>
*
* This library is free software; you can redistribute it and/or
* modify it under the terms of the GNU Library General Public
* License as published by the Free Software Foundation; either
* version 2 of the License, or (at your option) any later version.
*
* This library is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
* Library General Public License for more details.
*
* You should have received a copy of the GNU Library General Public
* License along with this library; if not, write to the
* Free Software Foundation, Inc., 51 Franklin Street, Suite 500,
* Boston, MA 02110-1335, USA.
*/
#ifdef HAVE_CONFIG_H
#include "config.h"
#endif
#include <gst/gst.h>
#include "gstpinospool.h"
GST_DEBUG_CATEGORY_STATIC (gst_pinos_pool_debug_category);
#define GST_CAT_DEFAULT gst_pinos_pool_debug_category
G_DEFINE_TYPE (GstPinosPool, gst_pinos_pool, GST_TYPE_BUFFER_POOL);
GstPinosPool *
gst_pinos_pool_new (void)
{
GstPinosPool *pool;
pool = g_object_new (GST_TYPE_PINOS_POOL, NULL);
return pool;
}
gboolean
gst_pinos_pool_add_buffer (GstPinosPool *pool, GstBuffer *buffer)
{
g_return_val_if_fail (GST_IS_PINOS_POOL (pool), FALSE);
g_return_val_if_fail (GST_IS_BUFFER (buffer), FALSE);
GST_OBJECT_LOCK (pool);
g_queue_push_tail (&pool->available, buffer);
g_cond_signal (&pool->cond);
GST_OBJECT_UNLOCK (pool);
return TRUE;
}
gboolean
gst_pinos_pool_remove_buffer (GstPinosPool *pool, GstBuffer *buffer)
{
g_return_val_if_fail (GST_IS_PINOS_POOL (pool), FALSE);
g_return_val_if_fail (GST_IS_BUFFER (buffer), FALSE);
GST_OBJECT_LOCK (pool);
g_queue_remove (&pool->available, buffer);
GST_OBJECT_UNLOCK (pool);
return TRUE;
}
static GstFlowReturn
acquire_buffer (GstBufferPool * pool, GstBuffer ** buffer,
GstBufferPoolAcquireParams * params)
{
GstPinosPool *p = GST_PINOS_POOL (pool);
GST_OBJECT_LOCK (pool);
while (p->available.length == 0) {
g_cond_wait (&p->cond, GST_OBJECT_GET_LOCK (pool));
}
*buffer = g_queue_pop_head (&p->available);
GST_OBJECT_UNLOCK (pool);
GST_DEBUG ("acquire buffer %p", *buffer);
return GST_FLOW_OK;
}
static void
release_buffer (GstBufferPool * pool, GstBuffer *buffer)
{
GstPinosPool *p = GST_PINOS_POOL (pool);
GST_DEBUG ("release buffer %p", buffer);
GST_OBJECT_LOCK (pool);
g_queue_push_tail (&p->available, buffer);
GST_OBJECT_UNLOCK (pool);
}
static gboolean
do_start (GstBufferPool * pool)
{
GstPinosPool *p = GST_PINOS_POOL (pool);
PinosProperties *props = NULL;
GstStructure *config;
GstCaps *caps;
guint size;
guint min_buffers;
guint max_buffers;
config = gst_buffer_pool_get_config (pool);
gst_buffer_pool_config_get_params (config, &caps, &size, &min_buffers, &max_buffers);
pinos_stream_start_allocation (p->stream, props);
return TRUE;
}
static void
gst_pinos_pool_finalize (GObject * object)
{
GstPinosPool *pool = GST_PINOS_POOL (object);
GST_DEBUG_OBJECT (pool, "finalize");
G_OBJECT_CLASS (gst_pinos_pool_parent_class)->finalize (object);
}
static void
gst_pinos_pool_class_init (GstPinosPoolClass * klass)
{
GObjectClass *gobject_class = G_OBJECT_CLASS (klass);
GstBufferPoolClass *bufferpool_class = GST_BUFFER_POOL_CLASS (klass);
gobject_class->finalize = gst_pinos_pool_finalize;
bufferpool_class->start = do_start;
bufferpool_class->acquire_buffer = acquire_buffer;
bufferpool_class->release_buffer = release_buffer;
GST_DEBUG_CATEGORY_INIT (gst_pinos_pool_debug_category, "pinospool", 0,
"debug category for pinospool object");
}
static void
gst_pinos_pool_init (GstPinosPool * pool)
{
g_cond_init (&pool->cond);
g_queue_init (&pool->available);
}

66
pinos/gst/gstpinospool.h Normal file
View file

@ -0,0 +1,66 @@
/* GStreamer
* Copyright (C) <2016> Wim Taymans <wim.taymans@gmail.com>
*
* This library is free software; you can redistribute it and/or
* modify it under the terms of the GNU Library General Public
* License as published by the Free Software Foundation; either
* version 2 of the License, or (at your option) any later version.
*
* This library is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
* Library General Public License for more details.
*
* You should have received a copy of the GNU Library General Public
* License along with this library; if not, write to the
* Free Software Foundation, Inc., 51 Franklin St, Fifth Floor,
* Boston, MA 02110-1301, USA.
*/
#ifndef __GST_PINOS_POOL_H__
#define __GST_PINOS_POOL_H__
#include <gst/gst.h>
#include <client/pinos.h>
G_BEGIN_DECLS
#define GST_TYPE_PINOS_POOL \
(gst_pinos_pool_get_type())
#define GST_PINOS_POOL(obj) \
(G_TYPE_CHECK_INSTANCE_CAST((obj),GST_TYPE_PINOS_POOL,GstPinosPool))
#define GST_PINOS_POOL_CLASS(klass) \
(G_TYPE_CHECK_CLASS_CAST((klass),GST_TYPE_PINOS_POOL,GstPinosPoolClass))
#define GST_IS_PINOS_POOL(obj) \
(G_TYPE_CHECK_INSTANCE_TYPE((obj),GST_TYPE_PINOS_POOL))
#define GST_IS_PINOS_POOL_CLASS(klass) \
(G_TYPE_CHECK_CLASS_TYPE((klass),GST_TYPE_PINOS_POOL))
#define GST_PINOS_POOL_GET_CLASS(klass) \
(G_TYPE_INSTANCE_GET_CLASS ((klass), GST_TYPE_PINOS_POOL, GstPinosPoolClass))
typedef struct _GstPinosPool GstPinosPool;
typedef struct _GstPinosPoolClass GstPinosPoolClass;
struct _GstPinosPool {
GstBufferPool parent;
PinosStream *stream;
GQueue available;
GCond cond;
};
struct _GstPinosPoolClass {
GstBufferPoolClass parent_class;
};
GType gst_pinos_pool_get_type (void);
GstPinosPool * gst_pinos_pool_new (void);
gboolean gst_pinos_pool_add_buffer (GstPinosPool *pool, GstBuffer *buffer);
gboolean gst_pinos_pool_remove_buffer (GstPinosPool *pool, GstBuffer *buffer);
G_END_DECLS
#endif /* __GST_PINOS_POOL_H__ */

View file

@ -48,6 +48,7 @@
#include "gsttmpfileallocator.h"
#include "gstpinosformat.h"
static GQuark process_mem_data_quark;
GST_DEBUG_CATEGORY_STATIC (pinos_sink_debug);
#define GST_CAT_DEFAULT pinos_sink_debug
@ -120,8 +121,9 @@ gst_pinos_sink_finalize (GObject * object)
if (pinossink->properties)
gst_structure_free (pinossink->properties);
g_hash_table_unref (pinossink->mem_ids);
g_object_unref (pinossink->allocator);
g_object_unref (pinossink->pool);
g_hash_table_unref (pinossink->buf_ids);
g_free (pinossink->path);
g_free (pinossink->client_name);
@ -133,7 +135,7 @@ gst_pinos_sink_propose_allocation (GstBaseSink * bsink, GstQuery * query)
{
GstPinosSink *pinossink = GST_PINOS_SINK (bsink);
gst_query_add_allocation_param (query, pinossink->allocator, NULL);
gst_query_add_allocation_pool (query, GST_BUFFER_POOL_CAST (pinossink->pool), 0, 0, 0);
return TRUE;
}
@ -207,18 +209,22 @@ gst_pinos_sink_class_init (GstPinosSinkClass * klass)
GST_DEBUG_CATEGORY_INIT (pinos_sink_debug, "pinossink", 0,
"Pinos Sink");
process_mem_data_quark = g_quark_from_static_string ("GstPinosSinkProcessMemQuark");
}
static void
gst_pinos_sink_init (GstPinosSink * sink)
{
sink->allocator = gst_tmpfile_allocator_new ();
sink->fdmanager = pinos_fd_manager_get (PINOS_FD_MANAGER_DEFAULT);
sink->pool = gst_pinos_pool_new ();
sink->client_name = pinos_client_name();
sink->mode = DEFAULT_PROP_MODE;
g_queue_init (&sink->empty);
g_queue_init (&sink->filled);
sink->mem_ids = g_hash_table_new_full (g_direct_hash, g_direct_equal, NULL,
(GDestroyNotify) gst_memory_unref);
sink->buf_ids = g_hash_table_new_full (g_direct_hash, g_direct_equal, NULL,
(GDestroyNotify) gst_buffer_unref);
}
static GstCaps *
@ -327,58 +333,128 @@ gst_pinos_sink_get_property (GObject * object, guint prop_id,
}
}
typedef struct {
GstPinosSink *sink;
guint id;
SpaMetaHeader *header;
guint flags;
} ProcessMemData;
static void
on_new_buffer (GObject *gobject,
process_mem_data_destroy (gpointer user_data)
{
ProcessMemData *data = user_data;
gst_object_unref (data->sink);
g_slice_free (ProcessMemData, data);
}
static void
on_add_buffer (GObject *gobject,
guint id,
gpointer user_data)
{
GstPinosSink *pinossink = user_data;
SpaBuffer *b;
GstBuffer *buf;
unsigned int i;
ProcessMemData data;
GST_LOG_OBJECT (pinossink, "add buffer");
if (!(b = pinos_stream_peek_buffer (pinossink->stream, id))) {
g_warning ("failed to peek buffer");
return;
}
buf = gst_buffer_new ();
data.sink = gst_object_ref (pinossink);
data.id = id;
data.header = NULL;
for (i = 0; i < b->n_metas; i++) {
SpaMeta *m = &SPA_BUFFER_METAS(b)[i];
switch (m->type) {
case SPA_META_TYPE_HEADER:
data.header = SPA_MEMBER (b, m->offset, SpaMetaHeader);
break;
default:
break;
}
}
for (i = 0; i < b->n_datas; i++) {
SpaData *d = &SPA_BUFFER_DATAS (b)[i];
SpaMemory *mem;
mem = spa_memory_find (&d->mem.mem);
if (mem->fd) {
GstMemory *fdmem = NULL;
fdmem = gst_fd_allocator_alloc (pinossink->allocator, dup (mem->fd),
d->mem.offset + d->mem.size, GST_FD_MEMORY_FLAG_NONE);
gst_memory_resize (fdmem, d->mem.offset, d->mem.size);
gst_buffer_append_memory (buf, fdmem);
} else {
gst_buffer_append_memory (buf,
gst_memory_new_wrapped (0, mem->ptr, mem->size, d->mem.offset,
d->mem.size, NULL, NULL));
}
}
data.flags = GST_BUFFER_FLAGS (buf);
gst_mini_object_set_qdata (GST_MINI_OBJECT_CAST (buf),
process_mem_data_quark,
g_slice_dup (ProcessMemData, &data),
process_mem_data_destroy);
gst_pinos_pool_add_buffer (GST_PINOS_POOL (pinossink->pool), buf);
g_hash_table_insert (pinossink->buf_ids, GINT_TO_POINTER (id), buf);
g_queue_push_tail (&pinossink->empty, buf);
pinos_main_loop_signal (pinossink->loop, FALSE);
}
static void
on_remove_buffer (GObject *gobject,
guint id,
gpointer user_data)
{
GstPinosSink *pinossink = user_data;
GstBuffer *buf;
GST_LOG_OBJECT (pinossink, "remove buffer");
buf = g_hash_table_lookup (pinossink->buf_ids, GINT_TO_POINTER (id));
GST_MINI_OBJECT_CAST (buf)->dispose = NULL;
gst_pinos_pool_remove_buffer (GST_PINOS_POOL (pinossink->pool), buf);
g_queue_remove (&pinossink->empty, buf);
g_queue_remove (&pinossink->filled, buf);
g_hash_table_remove (pinossink->buf_ids, GINT_TO_POINTER (id));
}
static void
on_new_buffer (GObject *gobject,
guint id,
gpointer user_data)
{
GstPinosSink *pinossink = user_data;
GstBuffer *buf;
GST_LOG_OBJECT (pinossink, "got new buffer");
if (pinossink->stream == NULL) {
GST_LOG_OBJECT (pinossink, "no stream");
return;
}
buf = g_hash_table_lookup (pinossink->buf_ids, GINT_TO_POINTER (id));
if (!(b = pinos_stream_peek_buffer (pinossink->stream))) {
g_warning ("failed to capture buffer");
return;
g_debug ("recycle buffer %d %p", id, buf);
if (buf) {
g_queue_remove (&pinossink->filled, buf);
g_queue_push_tail (&pinossink->empty, buf);
pinos_main_loop_signal (pinossink->loop, FALSE);
}
#if 0
pinos_buffer_iter_init (&it, pbuf);
while (pinos_buffer_iter_next (&it)) {
switch (pinos_buffer_iter_get_type (&it)) {
case PINOS_PACKET_TYPE_REUSE_MEM:
{
PinosPacketReuseMem p;
if (!pinos_buffer_iter_parse_reuse_mem (&it, &p))
continue;
GST_LOG ("mem index %d is reused", p.id);
g_hash_table_remove (pinossink->mem_ids, GINT_TO_POINTER (p.id));
break;
}
case PINOS_PACKET_TYPE_REFRESH_REQUEST:
{
PinosPacketRefreshRequest p;
if (!pinos_buffer_iter_parse_refresh_request (&it, &p))
continue;
GST_LOG ("refresh request");
gst_pad_push_event (GST_BASE_SINK_PAD (pinossink),
gst_video_event_new_upstream_force_key_unit (p.pts,
p.request_type == 1, 0));
break;
}
default:
break;
}
}
pinos_buffer_iter_end (&it);
#endif
}
static void
@ -409,6 +485,19 @@ on_stream_notify (GObject *gobject,
pinos_main_loop_signal (pinossink->loop, FALSE);
}
static void
on_format_notify (GObject *gobject,
GParamSpec *pspec,
gpointer user_data)
{
GstPinosSink *pinossink = user_data;
SpaFormat *format;
PinosProperties *props = NULL;
g_object_get (gobject, "format", &format, NULL);
}
static gboolean
gst_pinos_sink_setcaps (GstBaseSink * bsink, GstCaps * caps)
{
@ -452,7 +541,9 @@ gst_pinos_sink_setcaps (GstBaseSink * bsink, GstCaps * caps)
pinos_main_loop_wait (pinossink->loop);
}
}
res = TRUE;
#if 0
if (state != PINOS_STREAM_STATE_STREAMING) {
res = pinos_stream_start (pinossink->stream);
@ -468,6 +559,7 @@ gst_pinos_sink_setcaps (GstBaseSink * bsink, GstCaps * caps)
pinos_main_loop_wait (pinossink->loop);
}
}
#endif
pinos_main_loop_unlock (pinossink->loop);
pinossink->negotiated = res;
@ -483,115 +575,54 @@ start_error:
}
}
typedef struct {
SpaBuffer buffer;
SpaMeta metas[1];
SpaMetaHeader header;
SpaData datas[1];
GstMemory *mem;
GstPinosSink *pinossink;
int fd;
} SinkBuffer;
static GstFlowReturn
gst_pinos_sink_render (GstBaseSink * bsink, GstBuffer * buffer)
{
GstPinosSink *pinossink;
SinkBuffer *b;
GstMemory *mem = NULL;
GstClockTime pts, dts, base;
gsize size;
gboolean res;
ProcessMemData *data;
pinossink = GST_PINOS_SINK (bsink);
if (!pinossink->negotiated)
goto not_negotiated;
base = GST_ELEMENT_CAST (bsink)->base_time;
size = gst_buffer_get_size (buffer);
b = g_slice_new (SinkBuffer);
b->buffer.id = pinos_fd_manager_get_id (pinossink->fdmanager);
b->buffer.mem.mem.pool_id = SPA_ID_INVALID;
b->buffer.mem.mem.id = SPA_ID_INVALID;
b->buffer.mem.offset = 0;
b->buffer.mem.size = sizeof (SinkBuffer);
b->buffer.n_metas = 1;
b->buffer.metas = offsetof (SinkBuffer, metas);
b->buffer.n_datas = 1;
b->buffer.datas = offsetof (SinkBuffer, datas);
pts = GST_BUFFER_PTS (buffer);
dts = GST_BUFFER_DTS (buffer);
if (!GST_CLOCK_TIME_IS_VALID (pts))
pts = dts;
else if (!GST_CLOCK_TIME_IS_VALID (dts))
dts = pts;
b->header.flags = 0;
b->header.seq = GST_BUFFER_OFFSET (buffer);
b->header.pts = GST_CLOCK_TIME_IS_VALID (pts) ? pts + base : base;
b->header.dts_offset = GST_CLOCK_TIME_IS_VALID (dts) && GST_CLOCK_TIME_IS_VALID (pts) ? pts - dts : 0;
b->metas[0].type = SPA_META_TYPE_HEADER;
b->metas[0].offset = offsetof (SinkBuffer, header);
b->metas[0].size = sizeof (b->header);
if (gst_buffer_n_memory (buffer) == 1
&& gst_is_fd_memory (gst_buffer_peek_memory (buffer, 0))) {
mem = gst_buffer_get_memory (buffer, 0);
} else {
GstMapInfo minfo;
GstAllocationParams params = {0, 0, 0, 0, { NULL, }};
GST_INFO_OBJECT (bsink, "Buffer cannot be payloaded without copying");
mem = gst_allocator_alloc (pinossink->allocator, size, &params);
if (!gst_memory_map (mem, &minfo, GST_MAP_WRITE))
goto map_error;
gst_buffer_extract (buffer, 0, minfo.data, size);
gst_memory_unmap (mem, &minfo);
}
pinos_main_loop_lock (pinossink->loop);
if (pinos_stream_get_state (pinossink->stream) != PINOS_STREAM_STATE_STREAMING)
goto streaming_error;
b->mem = mem;
b->fd = gst_fd_memory_get_fd (mem);
if (buffer->pool != GST_BUFFER_POOL_CAST (pinossink->pool)) {
GstBuffer *b = NULL;
b->datas[0].mem.mem.pool_id = SPA_ID_INVALID;
b->datas[0].mem.mem.id = SPA_ID_INVALID;
b->datas[0].mem.offset = mem->offset;
b->datas[0].mem.size = mem->size;
b->datas[0].stride = 0;
while (TRUE) {
b = g_queue_peek_head (&pinossink->empty);
if (b)
break;
if (!(res = pinos_stream_send_buffer (pinossink->stream, &b->buffer)))
pinos_main_loop_wait (pinossink->loop);
}
g_queue_push_tail (&pinossink->filled, b);
buffer = b;
}
data = gst_mini_object_get_qdata (GST_MINI_OBJECT_CAST (buffer),
process_mem_data_quark);
if (!(res = pinos_stream_send_buffer (pinossink->stream, data->id)))
g_warning ("can't send buffer");
pinos_main_loop_unlock (pinossink->loop);
/* keep the memory around until we get the reuse mem message */
g_hash_table_insert (pinossink->mem_ids, GINT_TO_POINTER (b->buffer.id), b);
return GST_FLOW_OK;
not_negotiated:
{
return GST_FLOW_NOT_NEGOTIATED;
}
map_error:
{
GST_ELEMENT_ERROR (pinossink, RESOURCE, FAILED,
("failed to map buffer"), (NULL));
gst_memory_unref (mem);
return GST_FLOW_ERROR;
}
streaming_error:
{
pinos_main_loop_unlock (pinossink->loop);
gst_memory_unref (mem);
return GST_FLOW_ERROR;
}
}
@ -627,7 +658,11 @@ gst_pinos_sink_start (GstBaseSink * basesink)
pinos_main_loop_lock (pinossink->loop);
pinossink->stream = pinos_stream_new (pinossink->ctx, pinossink->client_name, props);
pinossink->pool->stream = pinossink->stream;
g_signal_connect (pinossink->stream, "notify::state", (GCallback) on_stream_notify, pinossink);
g_signal_connect (pinossink->stream, "notify::format", (GCallback) on_format_notify, pinossink);
g_signal_connect (pinossink->stream, "add-buffer", (GCallback) on_add_buffer, pinossink);
g_signal_connect (pinossink->stream, "remove-buffer", (GCallback) on_remove_buffer, pinossink);
g_signal_connect (pinossink->stream, "new-buffer", (GCallback) on_new_buffer, pinossink);
pinos_main_loop_unlock (pinossink->loop);
@ -644,6 +679,7 @@ gst_pinos_sink_stop (GstBaseSink * basesink)
pinos_stream_stop (pinossink->stream);
pinos_stream_disconnect (pinossink->stream);
g_clear_object (&pinossink->stream);
pinossink->pool->stream = NULL;
}
pinos_main_loop_unlock (pinossink->loop);
@ -787,10 +823,10 @@ gst_pinos_sink_change_state (GstElement * element, GstStateChange transition)
case GST_STATE_CHANGE_PLAYING_TO_PAUSED:
break;
case GST_STATE_CHANGE_PAUSED_TO_READY:
g_hash_table_remove_all (this->mem_ids);
g_hash_table_remove_all (this->buf_ids);
break;
case GST_STATE_CHANGE_READY_TO_NULL:
g_hash_table_remove_all (this->mem_ids);
g_hash_table_remove_all (this->buf_ids);
gst_pinos_sink_close (this);
break;
default:

View file

@ -24,6 +24,7 @@
#include <gst/base/gstbasesink.h>
#include <client/pinos.h>
#include <gst/gstpinospool.h>
G_BEGIN_DECLS
@ -84,8 +85,11 @@ struct _GstPinosSink {
GstStructure *properties;
GstPinosSinkMode mode;
PinosFdManager *fdmanager;
GHashTable *mem_ids;
GstPinosPool *pool;
GHashTable *buf_ids;
GQueue empty;
GQueue filled;
};
struct _GstPinosSinkClass {

View file

@ -187,7 +187,7 @@ gst_pinos_src_finalize (GObject * object)
gst_object_unref (pinossrc->clock);
g_free (pinossrc->path);
g_free (pinossrc->client_name);
g_hash_table_unref (pinossrc->mem_ids);
g_hash_table_unref (pinossrc->buf_ids);
G_OBJECT_CLASS (parent_class)->finalize (object);
}
@ -274,7 +274,7 @@ gst_pinos_src_init (GstPinosSrc * src)
src->fd_allocator = gst_fd_allocator_new ();
src->client_name = pinos_client_name ();
src->mem_ids = g_hash_table_new_full (g_direct_hash, g_direct_equal, NULL, (GDestroyNotify) gst_memory_unref);
src->buf_ids = g_hash_table_new_full (g_direct_hash, g_direct_equal, NULL, (GDestroyNotify) gst_buffer_unref);
}
static GstCaps *
@ -326,7 +326,9 @@ gst_pinos_src_src_fixate (GstBaseSrc * bsrc, GstCaps * caps)
typedef struct {
GstPinosSrc *src;
SpaBuffer *buffer;
guint id;
SpaMetaHeader *header;
guint flags;
} ProcessMemData;
static void
@ -338,8 +340,24 @@ process_mem_data_destroy (gpointer user_data)
g_slice_free (ProcessMemData, data);
}
static gboolean
buffer_recycle (GstMiniObject *obj)
{
ProcessMemData *data;
gst_mini_object_ref (obj);
data = gst_mini_object_get_qdata (obj,
process_mem_data_quark);
GST_BUFFER_FLAGS (obj) = data->flags;
pinos_stream_recycle_buffer (data->src->stream, data->id);
return FALSE;
}
static void
on_new_buffer (GObject *gobject,
on_add_buffer (GObject *gobject,
guint id,
gpointer user_data)
{
GstPinosSrc *pinossrc = user_data;
@ -348,39 +366,27 @@ on_new_buffer (GObject *gobject,
unsigned int i;
ProcessMemData data;
GST_LOG_OBJECT (pinossrc, "got new buffer");
if (!(b = pinos_stream_peek_buffer (pinossrc->stream))) {
g_warning ("failed to capture buffer");
GST_LOG_OBJECT (pinossrc, "add buffer");
if (!(b = pinos_stream_peek_buffer (pinossrc->stream, id))) {
g_warning ("failed to peek buffer");
return;
}
buf = gst_buffer_new ();
GST_MINI_OBJECT_CAST (buf)->dispose = buffer_recycle;
data.src = gst_object_ref (pinossrc);
data.buffer = b;
gst_mini_object_set_qdata (GST_MINI_OBJECT_CAST (buf),
process_mem_data_quark,
g_slice_dup (ProcessMemData, &data),
process_mem_data_destroy);
data.id = id;
data.header = NULL;
for (i = 0; i < b->n_metas; i++) {
SpaMeta *m = &SPA_BUFFER_METAS(b)[i];
switch (m->type) {
case SPA_META_TYPE_HEADER:
{
SpaMetaHeader *h = SPA_MEMBER (b, m->offset, SpaMetaHeader);
GST_INFO ("pts %" G_GUINT64_FORMAT ", dts_offset %"G_GUINT64_FORMAT, h->pts, h->dts_offset);
if (GST_CLOCK_TIME_IS_VALID (h->pts)) {
GST_BUFFER_PTS (buf) = h->pts;
if (GST_BUFFER_PTS (buf) + h->dts_offset > 0)
GST_BUFFER_DTS (buf) = GST_BUFFER_PTS (buf) + h->dts_offset;
}
GST_BUFFER_OFFSET (buf) = h->seq;
data.header = SPA_MEMBER (b, m->offset, SpaMetaHeader);
break;
}
default:
break;
}
@ -404,8 +410,59 @@ on_new_buffer (GObject *gobject,
d->mem.size, NULL, NULL));
}
}
data.flags = GST_BUFFER_FLAGS (buf);
gst_mini_object_set_qdata (GST_MINI_OBJECT_CAST (buf),
process_mem_data_quark,
g_slice_dup (ProcessMemData, &data),
process_mem_data_destroy);
g_hash_table_insert (pinossrc->buf_ids, GINT_TO_POINTER (id), buf);
}
static void
on_remove_buffer (GObject *gobject,
guint id,
gpointer user_data)
{
GstPinosSrc *pinossrc = user_data;
GstBuffer *buf;
GST_LOG_OBJECT (pinossrc, "remove buffer");
buf = g_hash_table_lookup (pinossrc->buf_ids, GINT_TO_POINTER (id));
GST_MINI_OBJECT_CAST (buf)->dispose = NULL;
g_hash_table_remove (pinossrc->buf_ids, GINT_TO_POINTER (id));
}
static void
on_new_buffer (GObject *gobject,
guint id,
gpointer user_data)
{
GstPinosSrc *pinossrc = user_data;
GstBuffer *buf;
GST_LOG_OBJECT (pinossrc, "got new buffer");
buf = g_hash_table_lookup (pinossrc->buf_ids, GINT_TO_POINTER (id));
if (buf) {
ProcessMemData *data;
SpaMetaHeader *h;
data = gst_mini_object_get_qdata (GST_MINI_OBJECT_CAST (buf),
process_mem_data_quark);
h = data->header;
if (h) {
GST_INFO ("pts %" G_GUINT64_FORMAT ", dts_offset %"G_GUINT64_FORMAT, h->pts, h->dts_offset);
if (GST_CLOCK_TIME_IS_VALID (h->pts)) {
GST_BUFFER_PTS (buf) = h->pts;
if (GST_BUFFER_PTS (buf) + h->dts_offset > 0)
GST_BUFFER_DTS (buf) = GST_BUFFER_PTS (buf) + h->dts_offset;
}
GST_BUFFER_OFFSET (buf) = h->seq;
}
g_queue_push_tail (&pinossrc->queue, buf);
pinos_main_loop_signal (pinossrc->loop, FALSE);
@ -486,7 +543,6 @@ parse_stream_properties (GstPinosSrc *pinossrc, PinosProperties *props)
static gboolean
gst_pinos_src_stream_start (GstPinosSrc *pinossrc)
{
SpaFormat *format;
gboolean res;
PinosProperties *props;
@ -505,17 +561,8 @@ gst_pinos_src_stream_start (GstPinosSrc *pinossrc)
}
g_object_get (pinossrc->stream, "properties", &props, NULL);
g_object_get (pinossrc->stream, "format", &format, NULL);
pinos_main_loop_unlock (pinossrc->loop);
if (format) {
GstCaps *caps = gst_caps_from_format (format);
gst_base_src_set_caps (GST_BASE_SRC (pinossrc), caps);
gst_caps_unref (caps);
spa_format_unref (format);
}
parse_stream_properties (pinossrc, props);
pinos_properties_free (props);
@ -675,6 +722,11 @@ on_format_notify (GObject *gobject,
caps = gst_caps_from_format (format);
gst_base_src_set_caps (GST_BASE_SRC (pinossrc), caps);
gst_caps_unref (caps);
pinos_stream_start_allocation (pinossrc->stream, NULL);
}
static gboolean
@ -944,6 +996,8 @@ gst_pinos_src_open (GstPinosSrc * pinossrc)
pinossrc->stream = pinos_stream_new (pinossrc->ctx, pinossrc->client_name, props);
g_signal_connect (pinossrc->stream, "notify::state", (GCallback) on_stream_notify, pinossrc);
g_signal_connect (pinossrc->stream, "notify::format", (GCallback) on_format_notify, pinossrc);
g_signal_connect (pinossrc->stream, "add-buffer", (GCallback) on_add_buffer, pinossrc);
g_signal_connect (pinossrc->stream, "remove-buffer", (GCallback) on_remove_buffer, pinossrc);
g_signal_connect (pinossrc->stream, "new-buffer", (GCallback) on_new_buffer, pinossrc);
pinos_main_loop_unlock (pinossrc->loop);

View file

@ -71,7 +71,7 @@ struct _GstPinosSrc {
GstAllocator *fd_allocator;
GstStructure *properties;
GHashTable *mem_ids;
GHashTable *buf_ids;
GQueue queue;
GstClock *clock;
};

View file

@ -110,7 +110,7 @@ bus_handler (GstBus *bus,
GST_INFO ("clock lost %s", GST_OBJECT_NAME (clock));
g_object_get (node, "properties", &props, NULL);
pinos_properties_remove (props, "gst.pipeline.clock");
pinos_properties_set (props, "gst.pipeline.clock", NULL);
g_object_set (node, "properties", props, NULL);
pinos_properties_free (props);

View file

@ -110,7 +110,7 @@ bus_handler (GstBus *bus,
GST_INFO ("clock lost %s", GST_OBJECT_NAME (clock));
g_object_get (node, "properties", &props, NULL);
pinos_properties_remove (props, "gst.pipeline.clock");
pinos_properties_set (props, "gst.pipeline.clock", NULL);
g_object_set (node, "properties", props, NULL);
pinos_properties_free (props);

View file

@ -159,20 +159,15 @@ on_sink_event (SpaNode *node, SpaEvent *event, void *user_data)
PinosSpaAlsaSinkPrivate *priv = this->priv;
switch (event->type) {
case SPA_EVENT_TYPE_PULL_INPUT:
case SPA_EVENT_TYPE_NEED_INPUT:
{
SpaInputInfo iinfo;
SpaResult res;
PinosRingbufferArea areas[2];
uint8_t *data;
size_t size, towrite, total;
SpaEventPullInput *pi;
pi = event->data;
g_debug ("pull ringbuffer %zd", pi->size);
size = pi->size;
size = 0;
data = NULL;
pinos_ringbuffer_get_read_areas (priv->ringbuffer, areas);
@ -194,8 +189,6 @@ on_sink_event (SpaNode *node, SpaEvent *event, void *user_data)
iinfo.port_id = event->port_id;
iinfo.flags = SPA_INPUT_FLAG_NONE;
iinfo.buffer_id = 0;
iinfo.offset = 0;
iinfo.size = total;
g_debug ("push sink %d", iinfo.buffer_id);
if ((res = spa_node_port_push_input (node, 1, &iinfo)) < 0)
@ -223,7 +216,13 @@ on_sink_event (SpaNode *node, SpaEvent *event, void *user_data)
}
break;
}
case SPA_EVENT_TYPE_STATE_CHANGE:
{
SpaEventStateChange *sc = event->data;
pinos_node_update_node_state (PINOS_NODE (this), sc->state);
break;
}
default:
g_debug ("got event %d", event->type);
break;
@ -244,7 +243,7 @@ setup_node (PinosSpaAlsaSink *this)
g_debug ("got get_props error %d", res);
value.type = SPA_PROP_TYPE_STRING;
value.value = "hw:0";
value.value = "hw:1";
value.size = strlen (value.value)+1;
spa_props_set_prop (props, spa_props_index_for_name (props, "device"), &value);
@ -376,7 +375,7 @@ on_received_buffer (PinosPort *port,
PinosSpaAlsaSink *this = user_data;
PinosSpaAlsaSinkPrivate *priv = this->priv;
unsigned int i;
SpaBuffer *buffer = NULL;
SpaBuffer *buffer = port->buffers[buffer_id];
for (i = 0; i < buffer->n_datas; i++) {
SpaData *d = SPA_BUFFER_DATAS (buffer);

View file

@ -139,7 +139,7 @@ on_source_event (SpaNode *node, SpaEvent *event, void *user_data)
PinosSpaV4l2SourcePrivate *priv = this->priv;
switch (event->type) {
case SPA_EVENT_TYPE_CAN_PULL_OUTPUT:
case SPA_EVENT_TYPE_HAVE_OUTPUT:
{
SpaOutputInfo info[1] = { 0, };
SpaResult res;
@ -188,6 +188,13 @@ on_source_event (SpaNode *node, SpaEvent *event, void *user_data)
}
break;
}
case SPA_EVENT_TYPE_STATE_CHANGE:
{
SpaEventStateChange *sc = event->data;
pinos_node_update_node_state (PINOS_NODE (this), sc->state);
break;
}
default:
g_debug ("got event %d", event->type);
break;
@ -208,7 +215,7 @@ setup_node (PinosSpaV4l2Source *this)
g_debug ("got get_props error %d", res);
value.type = SPA_PROP_TYPE_STRING;
value.value = "/dev/video0";
value.value = "/dev/video1";
value.size = strlen (value.value)+1;
spa_props_set_prop (props, spa_props_index_for_name (props, "device"), &value);
@ -363,9 +370,7 @@ on_received_event (PinosPort *port, SpaEvent *event, GError **error, gpointer us
if ((res = spa_node_port_reuse_buffer (node->node,
event->port_id,
rb->buffer_id,
rb->offset,
rb->size)) < 0)
rb->buffer_id)) < 0)
g_warning ("client-node %p: error reuse buffer: %d", node, res);
break;
}

View file

@ -180,8 +180,6 @@ on_received_buffer (PinosPort *port, uint32_t buffer_id, GError **error, gpointe
info[0].port_id = port->id;
info[0].buffer_id = buffer_id;
info[0].flags = SPA_INPUT_FLAG_NONE;
info[0].offset = 0;
info[0].size = -1;
if ((res = spa_node_port_push_input (node->node, 1, info)) < 0)
g_warning ("client-node %p: error pushing buffer: %d, %d", node, res, info[0].status);
@ -296,6 +294,8 @@ on_node_event (SpaNode *node, SpaEvent *event, void *user_data)
{
SpaEventStateChange *sc = event->data;
pinos_node_update_node_state (PINOS_NODE (this), sc->state);
switch (sc->state) {
case SPA_NODE_STATE_CONFIGURE:
{
@ -329,6 +329,24 @@ on_node_event (SpaNode *node, SpaEvent *event, void *user_data)
stop_thread (this);
break;
}
case SPA_EVENT_TYPE_HAVE_OUTPUT:
{
PinosPort *port;
SpaOutputInfo info[1] = { 0, };
SpaResult res;
GError *error = NULL;
if ((res = spa_node_port_pull_output (node, 1, info)) < 0)
g_debug ("client-node %p: got pull error %d, %d", this, res, info[0].status);
port = pinos_node_find_port (PINOS_NODE (this), info[0].port_id);
if (!pinos_port_send_buffer (port, info[0].buffer_id, &error)) {
g_debug ("send failed: %s", error->message);
g_clear_error (&error);
}
break;
}
case SPA_EVENT_TYPE_REUSE_BUFFER:
{
PinosPort *port;

View file

@ -56,8 +56,13 @@ struct _PinosLinkPrivate
SpaNode *input_node;
uint32_t input_port;
SpaBuffer *buffers[16];
unsigned int n_buffers;
SpaNodeState input_state;
SpaNodeState output_state;
SpaBuffer *in_buffers[16];
unsigned int n_in_buffers;
SpaBuffer *out_buffers[16];
unsigned int n_out_buffers;
};
G_DEFINE_TYPE (PinosLink, pinos_link, G_TYPE_OBJECT);
@ -268,8 +273,9 @@ do_allocation (PinosLink *this)
PinosLinkPrivate *priv = this->priv;
SpaResult res;
const SpaPortInfo *iinfo, *oinfo;
SpaPortInfoFlags in_flags, out_flags;
g_debug ("link %p: doing alloc buffers", this);
g_debug ("link %p: doing alloc buffers %p %p", this, priv->output_node, priv->input_node);
/* find out what's possible */
if ((res = spa_node_port_get_info (priv->output_node, priv->output_port, &oinfo)) < 0) {
g_warning ("error get port info: %d", res);
@ -280,18 +286,79 @@ do_allocation (PinosLink *this)
return res;
}
priv->n_buffers = 16;
if ((res = spa_node_port_alloc_buffers (priv->output_node, priv->output_port,
iinfo->params, iinfo->n_params,
priv->buffers, &priv->n_buffers)) < 0) {
g_warning ("error alloc buffers: %d", res);
return res;
spa_debug_port_info (oinfo);
spa_debug_port_info (iinfo);
priv->n_in_buffers = 16;
priv->n_out_buffers = 16;
if ((oinfo->flags & SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS) &&
(iinfo->flags & SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS)) {
out_flags = SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS;
in_flags = SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS;
} else if ((oinfo->flags & SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS) &&
(iinfo->flags & SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS)) {
out_flags = SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS;
in_flags = SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS;
} else if ((oinfo->flags & SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS) &&
(iinfo->flags & SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS)) {
out_flags = SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS;
in_flags = SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS;
if ((res = spa_buffer_alloc (oinfo->params, oinfo->n_params,
priv->in_buffers,
&priv->n_in_buffers)) < 0) {
g_warning ("error alloc buffers: %d", res);
return res;
}
memcpy (priv->out_buffers, priv->in_buffers, priv->n_in_buffers * sizeof (SpaBuffer*));
priv->n_out_buffers = priv->n_in_buffers;
} else if ((oinfo->flags & SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS) &&
(iinfo->flags & SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS)) {
out_flags = SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS;
in_flags = SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS;
} else {
g_warning ("error no common allocation found");
return SPA_RESULT_ERROR;
}
if ((res = spa_node_port_use_buffers (priv->input_node, priv->input_port,
priv->buffers, priv->n_buffers)) < 0) {
g_warning ("error alloc buffers: %d", res);
return res;
if (in_flags & SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS) {
if ((res = spa_node_port_alloc_buffers (priv->input_node, priv->input_port,
oinfo->params, oinfo->n_params,
priv->in_buffers, &priv->n_in_buffers)) < 0) {
g_warning ("error alloc buffers: %d", res);
return res;
}
priv->input->n_buffers = priv->n_in_buffers;
priv->input->buffers = priv->in_buffers;
}
if (out_flags & SPA_PORT_INFO_FLAG_CAN_ALLOC_BUFFERS) {
if ((res = spa_node_port_alloc_buffers (priv->output_node, priv->output_port,
iinfo->params, iinfo->n_params,
priv->out_buffers, &priv->n_out_buffers)) < 0) {
g_warning ("error alloc buffers: %d", res);
return res;
}
priv->output->n_buffers = priv->n_out_buffers;
priv->output->buffers = priv->out_buffers;
}
if (in_flags & SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS) {
if ((res = spa_node_port_use_buffers (priv->input_node, priv->input_port,
priv->out_buffers, priv->n_out_buffers)) < 0) {
g_warning ("error use buffers: %d", res);
return res;
}
priv->input->n_buffers = priv->n_out_buffers;
priv->input->buffers = priv->out_buffers;
}
if (out_flags & SPA_PORT_INFO_FLAG_CAN_USE_BUFFERS) {
if ((res = spa_node_port_use_buffers (priv->output_node, priv->output_port,
priv->in_buffers, priv->n_in_buffers)) < 0) {
g_warning ("error use buffers: %d", res);
return res;
}
priv->output->n_buffers = priv->n_in_buffers;
priv->output->buffers = priv->in_buffers;
}
priv->allocated = TRUE;
@ -299,13 +366,86 @@ do_allocation (PinosLink *this)
return SPA_RESULT_OK;
}
static SpaResult
do_start (PinosLink *this)
{
PinosLinkPrivate *priv = this->priv;
SpaCommand cmd;
SpaResult res;
cmd.type = SPA_COMMAND_START;
if ((res = spa_node_send_command (priv->input_node, &cmd)) < 0)
g_warning ("got error %d", res);
if ((res = spa_node_send_command (priv->output_node, &cmd)) < 0)
g_warning ("got error %d", res);
return res;
}
static SpaResult
do_stop (PinosLink *this)
{
PinosLinkPrivate *priv = this->priv;
SpaCommand cmd;
SpaResult res;
cmd.type = SPA_COMMAND_STOP;
if ((res = spa_node_send_command (priv->input_node, &cmd)) < 0)
g_warning ("got error %d", res);
if ((res = spa_node_send_command (priv->output_node, &cmd)) < 0)
g_warning ("got error %d", res);
return res;
}
static SpaResult
check_states (PinosLink *this)
{
PinosLinkPrivate *priv = this->priv;
SpaResult res;
g_debug ("link %p: input %d, output %d", this, priv->input_state, priv->output_state);
if (priv->input_state == SPA_NODE_STATE_CONFIGURE &&
priv->output_state == SPA_NODE_STATE_CONFIGURE &&
!priv->negotiated) {
if ((res = do_negotiate (this)) < 0)
return res;
}
if (priv->input_state == SPA_NODE_STATE_READY &&
priv->output_state == SPA_NODE_STATE_READY &&
!priv->allocated) {
if ((res = do_allocation (this)) < 0)
return res;
if ((res = do_start (this)) < 0)
return res;
}
return SPA_RESULT_OK;
}
static void
on_node_state_notify (GObject *obj,
GParamSpec *pspec,
gpointer user_data)
{
PinosLink *this = user_data;
PinosLinkPrivate *priv = this->priv;
g_debug ("link %p: node %p state change", this, obj);
if (obj == G_OBJECT (priv->input->node))
priv->input_state = priv->input->node->node_state;
else
priv->output_state = priv->output->node->node_state;
check_states (this);
}
static gboolean
on_activate (PinosPort *port, gpointer user_data)
{
PinosLink *this = user_data;
PinosLinkPrivate *priv = this->priv;
SpaCommand cmd;
SpaResult res;
if (priv->active)
return TRUE;
@ -316,18 +456,7 @@ on_activate (PinosPort *port, gpointer user_data)
else
pinos_port_activate (priv->input);
if (!priv->negotiated)
do_negotiate (this);
/* negotiate allocation */
if (!priv->allocated)
do_allocation (this);
cmd.type = SPA_COMMAND_START;
if ((res = spa_node_send_command (priv->input_node, &cmd)) < 0)
g_warning ("got error %d", res);
if ((res = spa_node_send_command (priv->output_node, &cmd)) < 0)
g_warning ("got error %d", res);
check_states (this);
return TRUE;
}
@ -335,10 +464,8 @@ on_activate (PinosPort *port, gpointer user_data)
static gboolean
on_deactivate (PinosPort *port, gpointer user_data)
{
PinosLink *link = user_data;
PinosLinkPrivate *priv = link->priv;
SpaCommand cmd;
SpaResult res;
PinosLink *this = user_data;
PinosLinkPrivate *priv = this->priv;
if (!priv->active)
return TRUE;
@ -349,11 +476,7 @@ on_deactivate (PinosPort *port, gpointer user_data)
else
pinos_port_deactivate (priv->input);
cmd.type = SPA_COMMAND_STOP;
if ((res = spa_node_send_command (priv->input_node, &cmd)) < 0)
g_warning ("got error %d", res);
if ((res = spa_node_send_command (priv->output_node, &cmd)) < 0)
g_warning ("got error %d", res);
do_stop (this);
return TRUE;
}
@ -363,8 +486,8 @@ on_property_notify (GObject *obj,
GParamSpec *pspec,
gpointer user_data)
{
PinosLink *link = user_data;
PinosLinkPrivate *priv = link->priv;
PinosLink *this = user_data;
PinosLinkPrivate *priv = this->priv;
if (pspec == NULL || strcmp (g_param_spec_get_name (pspec), "output") == 0) {
gchar *port = g_strdup_printf ("%s:%d", pinos_node_get_object_path (priv->output->node),
@ -384,32 +507,38 @@ on_property_notify (GObject *obj,
static void
pinos_link_constructed (GObject * object)
{
PinosLink *link = PINOS_LINK (object);
PinosLinkPrivate *priv = link->priv;
PinosLink *this = PINOS_LINK (object);
PinosLinkPrivate *priv = this->priv;
priv->output_id = pinos_port_add_send_cb (priv->output,
on_output_buffer,
on_output_event,
link,
this,
NULL);
priv->input_id = pinos_port_add_send_cb (priv->input,
on_input_buffer,
on_input_event,
link,
this,
NULL);
g_signal_connect (priv->input, "activate", (GCallback) on_activate, link);
g_signal_connect (priv->input, "deactivate", (GCallback) on_deactivate, link);
g_signal_connect (priv->output, "activate", (GCallback) on_activate, link);
g_signal_connect (priv->output, "deactivate", (GCallback) on_deactivate, link);
priv->input_state = priv->input->node->node_state;
priv->output_state = priv->output->node->node_state;
g_signal_connect (link, "notify", (GCallback) on_property_notify, link);
g_signal_connect (priv->input->node, "notify::node-state", (GCallback) on_node_state_notify, this);
g_signal_connect (priv->output->node, "notify::node-state", (GCallback) on_node_state_notify, this);
g_signal_connect (priv->input, "activate", (GCallback) on_activate, this);
g_signal_connect (priv->input, "deactivate", (GCallback) on_deactivate, this);
g_signal_connect (priv->output, "activate", (GCallback) on_activate, this);
g_signal_connect (priv->output, "deactivate", (GCallback) on_deactivate, this);
g_signal_connect (this, "notify", (GCallback) on_property_notify, this);
G_OBJECT_CLASS (pinos_link_parent_class)->constructed (object);
on_property_notify (G_OBJECT (link), NULL, link);
g_debug ("link %p: constructed", link);
link_register_object (link);
on_property_notify (G_OBJECT (this), NULL, this);
g_debug ("link %p: constructed", this);
link_register_object (this);
}
static void
@ -427,8 +556,8 @@ pinos_link_dispose (GObject * object)
pinos_port_deactivate (priv->input);
pinos_port_deactivate (priv->output);
}
g_clear_object (&priv->input);
g_clear_object (&priv->output);
priv->input = NULL;
priv->output = NULL;
link_unregister_object (link);
G_OBJECT_CLASS (pinos_link_parent_class)->dispose (object);

View file

@ -64,6 +64,7 @@ enum
PROP_STATE,
PROP_PROPERTIES,
PROP_NODE,
PROP_NODE_STATE,
};
enum
@ -266,6 +267,10 @@ pinos_node_get_property (GObject *_object,
g_value_set_pointer (value, node->node);
break;
case PROP_NODE_STATE:
g_value_set_uint (value, node->node_state);
break;
default:
G_OBJECT_WARN_INVALID_PROPERTY_ID (node, prop_id, pspec);
break;
@ -496,6 +501,17 @@ pinos_node_class_init (PinosNodeClass * klass)
G_PARAM_CONSTRUCT_ONLY |
G_PARAM_STATIC_STRINGS));
g_object_class_install_property (gobject_class,
PROP_NODE_STATE,
g_param_spec_uint ("node-state",
"Node State",
"The state of the SPA node",
0,
G_MAXUINT,
SPA_NODE_STATE_INIT,
G_PARAM_READABLE |
G_PARAM_STATIC_STRINGS));
signals[SIGNAL_REMOVE] = g_signal_new ("remove",
G_TYPE_FROM_CLASS (klass),
G_SIGNAL_RUN_LAST,
@ -954,3 +970,24 @@ pinos_node_report_busy (PinosNode *node)
g_debug ("node %p: report busy", node);
pinos_node_set_state (node, PINOS_NODE_STATE_RUNNING);
}
/**
* pinos_node_update_node_state:
* @node: a #PinosNode
* @state: a #SpaNodeState
*
* Update the state of a SPA node. This method is used from
* inside @node itself.
*/
void
pinos_node_update_node_state (PinosNode *node,
SpaNodeState state)
{
g_return_if_fail (PINOS_IS_NODE (node));
if (node->node_state != state) {
g_debug ("node %p: update SPA state to %d", node, state);
node->node_state = state;
g_object_notify (G_OBJECT (node), "node-state");
}
}

View file

@ -52,6 +52,7 @@ struct _PinosNode {
GObject object;
SpaNode *node;
SpaNodeState node_state;
PinosNodePrivate *priv;
};
@ -111,6 +112,8 @@ void pinos_node_report_error (PinosNode *node, GError
void pinos_node_report_idle (PinosNode *node);
void pinos_node_report_busy (PinosNode *node);
void pinos_node_update_node_state (PinosNode *node, SpaNodeState state);
G_END_DECLS
#endif /* __PINOS_NODE_H__ */

View file

@ -55,6 +55,9 @@ struct _PinosPort {
uint32_t id;
PinosNode *node;
SpaBuffer **buffers;
guint n_buffers;
PinosPortPrivate *priv;
};