Skip to content
Merged
67 changes: 65 additions & 2 deletions src/client.c
Original file line number Diff line number Diff line change
Expand Up @@ -1711,13 +1711,62 @@ client_set_default_metadata (Client * client,
"Setting new default metadata: %" GST_PTR_FORMAT, client->metadata);
}

/* A re-announced onMetaData only carries what the muxer knows about, so it may
* update values but must not drop the fields it omits (avcprofile, avclevel,
* ...) nor change the type of one we already publish. */
static gboolean
client_merge_metadata_field (const GstIdStr * field, const GValue * value,
gpointer user_data)
{
GstStructure *metadata = user_data;
const GValue *existing = gst_structure_id_str_get_value (metadata, field);
Comment thread
camilo-celis marked this conversation as resolved.

if (existing == NULL || G_VALUE_TYPE (existing) == G_VALUE_TYPE (value))
gst_structure_id_str_set_value (metadata, field, value);

return TRUE;
}

/* An FLV script-data tag ("onMetaData") carries the same AMF payload as the
* RTMP "@setDataFrame"/"onMetaData" MSG_NOTIFY, so decode and apply it. */
static void
client_handle_flv_script_data (Client * client)
{
AmfDec *dec = amf_dec_new (client->buf, 0);
Comment thread
gregadams marked this conversation as resolved.
Comment thread
gregadams marked this conversation as resolved.
gchar *type = amf_dec_load_string (dec);

if (g_strcmp0 (type, "onMetaData") == 0) {
GstStructure *metadata = amf_dec_load_object (dec);
if (amf_dec_payload_fully_decoded (dec)) {
if (client->metadata == NULL)
client->metadata = gst_structure_new_empty ("object");
gst_structure_foreach_id_str (metadata, client_merge_metadata_field,
client->metadata);
client->new_metadata = TRUE;
GST_DEBUG_OBJECT (client->server, "(%s) FLV METADATA %" GST_PTR_FORMAT,
client->path, client->metadata);
} else {
GST_WARNING_OBJECT (client->server,
"(%s) ignoring malformed FLV onMetaData, decoded %" G_GSIZE_FORMAT
" of %u bytes", client->path, dec->pos, client->buf->len);
}
Comment thread
Copilot marked this conversation as resolved.
gst_structure_free (metadata);
} else {
GST_DEBUG_OBJECT (client->server, "ignoring FLV script data: %s",
type ? type : "(unknown)");
}

g_free (type);
amf_dec_free (dec);
}

static PexRtmpServerStatus
client_handle_flv_buffer (Client * client, GstBuffer * buf)
{
RTMPMessage msg;
GstMapInfo map;
guint payload_size;
PexRtmpServerStatus ret = PEX_RTMP_SERVER_STATUS_BAD;
PexRtmpServerStatus ret = PEX_RTMP_SERVER_STATUS_OK;
guint total_parsed = 0;

gst_buffer_map (buf, &map, GST_MAP_READ);
Expand All @@ -1744,6 +1793,7 @@ client_handle_flv_buffer (Client * client, GstBuffer * buf)
if (!(parsed = flv_parse_tag (data, map.size - total_parsed,
&msg.type, &payload_size, &msg.abs_timestamp))) {
GST_WARNING_OBJECT (client->server, "Could not parse header!");
ret = PEX_RTMP_SERVER_STATUS_PARSE_FAILED;
goto done;
}

Expand All @@ -1756,7 +1806,20 @@ client_handle_flv_buffer (Client * client, GstBuffer * buf)
data + parsed, payload_size);
msg.len = payload_size;
msg.buf = client->buf;
ret = client_handle_message (client, &msg);
/* this largely reports how forwarding to the subscribers went, which
must not abandon the rest of the publisher's input */
PexRtmpServerStatus msg_ret = client_handle_message (client, &msg);
if (msg_ret != PEX_RTMP_SERVER_STATUS_OK) {
GST_WARNING_OBJECT (client->server,
Comment thread
gregadams marked this conversation as resolved.
"(%s) failed to handle FLV tag 0x%x (ret=%d), continuing",
client->path, msg.type, msg_ret);
}
client->buf =
g_byte_array_remove_range (client->buf, 0, client->buf->len);
} else if (msg.type == MSG_NOTIFY) {
client->buf = g_byte_array_append (client->buf,
data + parsed, payload_size);
client_handle_flv_script_data (client);
client->buf =
g_byte_array_remove_range (client->buf, 0, client->buf->len);
}
Expand Down
27 changes: 23 additions & 4 deletions utils/amf.c
Original file line number Diff line number Diff line change
Expand Up @@ -154,11 +154,14 @@ amf_enc_write_string (AmfEnc * enc, const gchar * str)
static void
amf_enc_write_int (AmfEnc * enc, gint i)
{
if (enc->version == AMF3_VERSION)
amf_enc_add_char (enc, AMF3_INTEGER);
else
g_assert_not_reached ();
/* AMF0 has no integer type, so widen to a number instead of dying on
values that came in over an AMF3 decode */
if (enc->version != AMF3_VERSION) {
amf_enc_write_double (enc, (gdouble) i);
return;
}

amf_enc_add_char (enc, AMF3_INTEGER);
amf_enc_add_int (enc, i);
}

Expand Down Expand Up @@ -570,6 +573,22 @@ amf_dec_load_object (AmfDec * dec)
return amf_dec_load_object_with_depth (dec, 0);
}

/* amf_dec_load_object() always hands back a (possibly partially filled)
* structure, so a truncated payload has to be spotted by the caller: a good
* one is consumed in full and ends with the AMF0 object-end marker. */
gboolean
amf_dec_payload_fully_decoded (const AmfDec * dec)
{
const guint8 object_end[] = { 0x00, 0x00, AMF0_OBJECT_END };

if (dec->pos != dec->buf->len)
return FALSE;

return dec->buf->len >= sizeof (object_end) &&
memcmp (&dec->buf->data[dec->buf->len - sizeof (object_end)],
object_end, sizeof (object_end)) == 0;
}

GValue *
amf_dec_load (AmfDec * dec)
{
Expand Down
1 change: 1 addition & 0 deletions utils/amf.h
Original file line number Diff line number Diff line change
Expand Up @@ -129,5 +129,6 @@ GValue * amf_dec_load (AmfDec * dec);
gboolean amf_dec_load_number (AmfDec * dec, gdouble * ret);
gboolean amf_dec_load_integer (AmfDec * dec, gint * ret);
gboolean amf_dec_load_boolean (AmfDec * dec, gboolean * ret);
gboolean amf_dec_payload_fully_decoded (const AmfDec * dec);

#endif /* __AMF_H__ */
2 changes: 1 addition & 1 deletion utils/parse.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@

#include <gst/gst.h>

#ifdef G_OS_WIN32
#if defined(G_OS_WIN32) && !defined(PEX_RTMPSERVER_STATIC_BUILD)
# ifdef PEX_RTMPSERVER_EXPORTS
# define PEX_RTMPSERVER_EXPORT __declspec(dllexport)
# else
Expand Down
2 changes: 1 addition & 1 deletion utils/tcp.h
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@

#include <gst/gst.h>

#ifdef G_OS_WIN32
#if defined(G_OS_WIN32) && !defined(PEX_RTMPSERVER_STATIC_BUILD)
# ifdef PEX_RTMPSERVER_EXPORTS
# define PEX_RTMPSERVER_EXPORT __declspec(dllexport)
# else
Expand Down