Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 19 additions & 2 deletions src/core/ddsi/include/dds/ddsi/ddsi_serdata.h
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,11 @@ typedef size_t (*ddsi_serdata_print_t) (const struct ddsi_sertype *type, const s
typedef void (*ddsi_serdata_get_keyhash_t) (const struct ddsi_serdata *d, struct ddsi_keyhash *buf, bool force_md5)
ddsrt_nonnull_all;

/* Return the related sample identity for a sample when one should be emitted in inline QoS.
Returns true when available, false otherwise. */
typedef bool (*ddsi_serdata_get_related_sample_identity_t) (const struct ddsi_serdata *d, ddsi_guid_t *writer_guid, ddsi_seqno_t *seq)
ddsrt_nonnull ((1)) ddsrt_attribute_warn_unused_result;

// Used for taking a loaned sample and constructing a serdata around this
// takes over ownership of loan on success (leaves it unchanged on failure)
typedef struct ddsi_serdata* (*ddsi_serdata_from_loan_t) (const struct ddsi_sertype *type, enum ddsi_serdata_kind kind, const char *sample, struct dds_loaned_sample *loaned_sample, bool will_require_cdr)
Expand Down Expand Up @@ -180,11 +185,17 @@ struct ddsi_serdata_ops {
ddsi_serdata_get_keyhash_t get_keyhash;
ddsi_serdata_from_loan_t from_loaned_sample;
ddsi_serdata_from_psmx_t from_psmx;
#ifdef __cplusplus
ddsi_serdata_get_related_sample_identity_t get_related_sample_identity = nullptr;
#else
ddsi_serdata_get_related_sample_identity_t get_related_sample_identity;
#endif
};

#define DDSI_SERDATA_HAS_PRINT 1
#define DDSI_SERDATA_HAS_FROM_SER_IOV 1
#define DDSI_SERDATA_HAS_GET_KEYHASH 1
#define DDSI_SERDATA_HAS_GET_RELATED_SAMPLE_IDENTITY 1

/** @component typesupport_if */
DDS_EXPORT void ddsi_serdata_init (struct ddsi_serdata *d, const struct ddsi_sertype *tp, enum ddsi_serdata_kind kind)
Expand Down Expand Up @@ -389,6 +400,14 @@ inline void ddsi_serdata_get_keyhash (const struct ddsi_serdata *d, struct ddsi_
d->ops->get_keyhash (d, buf, force_md5);
}

/** @component typesupport_if */
DDS_INLINE_EXPORT inline bool ddsi_serdata_get_related_sample_identity (const struct ddsi_serdata *d, ddsi_guid_t *writer_guid, ddsi_seqno_t *seq)
ddsrt_nonnull ((1)) ddsrt_attribute_warn_unused_result;

inline bool ddsi_serdata_get_related_sample_identity (const struct ddsi_serdata *d, ddsi_guid_t *writer_guid, ddsi_seqno_t *seq) {
return d->ops->get_related_sample_identity && d->ops->get_related_sample_identity (d, writer_guid, seq);
}

DDS_INLINE_EXPORT inline struct ddsi_serdata *ddsi_serdata_from_loaned_sample(const struct ddsi_sertype *type, enum ddsi_serdata_kind kind, const char *sample, struct dds_loaned_sample *loan, bool will_require_cdr) ddsrt_nonnull_all;

/** @component typesupport_if */
Expand All @@ -409,5 +428,3 @@ inline struct ddsi_serdata *ddsi_serdata_from_psmx(const struct ddsi_sertype *ty
#endif

#endif //DDSI_SERDATA_H


1 change: 1 addition & 0 deletions src/core/ddsi/src/ddsi__protocol.h
Original file line number Diff line number Diff line change
Expand Up @@ -287,6 +287,7 @@ typedef union ddsi_rtps_submessage {
#define DDSI_PID_ENTITY_NAME 0x62u
#define DDSI_PID_KEYHASH 0x70u
#define DDSI_PID_STATUSINFO 0x71u
#define DDSI_PID_RELATED_SAMPLE_IDENTITY 0x83u
#define DDSI_PID_CONTENT_FILTER_INFO 0x55u
#define DDSI_PID_COHERENT_SET 0x56u
#define DDSI_PID_DIRECTED_WRITE 0x57u
Expand Down
4 changes: 4 additions & 0 deletions src/core/ddsi/src/ddsi__xmsg.h
Original file line number Diff line number Diff line change
Expand Up @@ -287,6 +287,10 @@ void *ddsi_xmsg_addpar (struct ddsi_xmsg *m, ddsi_parameterid_t pid, size_t len)
void ddsi_xmsg_addpar_keyhash (struct ddsi_xmsg *m, const struct ddsi_serdata *serdata, bool force_md5)
ddsrt_nonnull_all;

/** @component rtps_submsg */
void ddsi_xmsg_addpar_related_sample_identity (struct ddsi_xmsg *m, const ddsi_guid_t *writer_guid, ddsi_seqno_t seq)
ddsrt_nonnull_all;

/** @component rtps_submsg */
void ddsi_xmsg_addpar_statusinfo (struct ddsi_xmsg *m, unsigned statusinfo)
ddsrt_nonnull_all;
Expand Down
2 changes: 1 addition & 1 deletion src/core/ddsi/src/ddsi_serdata.c
Original file line number Diff line number Diff line change
Expand Up @@ -11,7 +11,6 @@
#include <stddef.h>
#include <ctype.h>
#include <assert.h>
#include <string.h>

#include "dds/export.h"
#include "dds/ddsrt/md5.h"
Expand Down Expand Up @@ -80,5 +79,6 @@ DDS_EXPORT extern inline bool ddsi_serdata_eqkey (const struct ddsi_serdata *a,
DDS_EXPORT extern inline size_t ddsi_serdata_print (const struct ddsi_serdata *d, char *buf, size_t size);
DDS_EXPORT extern inline size_t ddsi_serdata_print_untyped (const struct ddsi_sertype *type, const struct ddsi_serdata *d, char *buf, size_t size);
DDS_EXPORT extern inline void ddsi_serdata_get_keyhash (const struct ddsi_serdata *d, struct ddsi_keyhash *buf, bool force_md5);
DDS_EXPORT extern inline bool ddsi_serdata_get_related_sample_identity (const struct ddsi_serdata *d, ddsi_guid_t *writer_guid, ddsi_seqno_t *seq);
DDS_EXPORT extern inline struct ddsi_serdata *ddsi_serdata_from_loaned_sample (const struct ddsi_sertype *type, enum ddsi_serdata_kind kind, const char *sample, struct dds_loaned_sample *loan, bool will_require_cdr);
DDS_EXPORT extern inline struct ddsi_serdata *ddsi_serdata_from_psmx (const struct ddsi_sertype *type, struct dds_loaned_sample *data);
22 changes: 16 additions & 6 deletions src/core/ddsi/src/ddsi_transmit.c
Original file line number Diff line number Diff line change
Expand Up @@ -57,9 +57,9 @@ static int have_reliable_subs (const struct ddsi_writer *wr)
static dds_return_t ddsi_create_fragment_message_simple (struct ddsi_writer *wr, ddsi_seqno_t seq, struct ddsi_serdata *serdata, struct ddsi_xmsg **pmsg)
{
#define TEST_KEYHASH 0
/* actual expected_inline_qos_size is typically 0, but always claiming 32 bytes won't make
/* actual expected_inline_qos_size is typically 0, but always claiming 60 bytes won't make
a difference, so no point in being precise */
const size_t expected_inline_qos_size = /* statusinfo */ 8 + /* keyhash */ 20 + /* sentinel */ 4;
const size_t expected_inline_qos_size = /* statusinfo */ 8 + /* keyhash */ 20 + /* related_sample_identity */ 28 + /* sentinel */ 4;
struct ddsi_domaingv const * const gv = wr->e.gv;
struct ddsi_xmsg_marker sm_marker;
unsigned char contentflag = 0;
Expand All @@ -83,7 +83,7 @@ static dds_return_t ddsi_create_fragment_message_simple (struct ddsi_writer *wr,

ASSERT_MUTEX_HELD (&wr->e.lock);

/* INFO_TS: 12 bytes, ddsi_rtps_data_t: 24 bytes, expected inline QoS: 32 => should be single chunk */
/* INFO_TS: 12 bytes, ddsi_rtps_data_t: 24 bytes, expected inline QoS: 60 => should be single chunk */
if ((*pmsg = ddsi_xmsg_new (gv->xmsgpool, &wr->e.guid, wr->c.pp, sizeof (ddsi_rtps_info_ts_t) + sizeof (ddsi_rtps_data_t) + expected_inline_qos_size, DDSI_XMSG_KIND_DATA)) == NULL)
return DDS_RETCODE_OUT_OF_RESOURCES;

Expand All @@ -104,10 +104,14 @@ static dds_return_t ddsi_create_fragment_message_simple (struct ddsi_writer *wr,
ddsi_xmsg_setwriterseq (*pmsg, &wr->e.guid, seq);

/* Adding parameters means potential reallocing, so sm, ddcmn now likely become invalid */
ddsi_guid_t related_writer_guid;
ddsi_seqno_t related_seq;
if (wr->num_readers_requesting_keyhash > 0)
ddsi_xmsg_addpar_keyhash (*pmsg, serdata, wr->force_md5_keyhash);
if (serdata->statusinfo)
ddsi_xmsg_addpar_statusinfo (*pmsg, serdata->statusinfo);
if (ddsi_serdata_get_related_sample_identity (serdata, &related_writer_guid, &related_seq))
ddsi_xmsg_addpar_related_sample_identity (*pmsg, &related_writer_guid, related_seq);
if (ddsi_xmsg_addpar_sentinel_ifparam (*pmsg) > 0)
{
data = ddsi_xmsg_submsg_from_marker (*pmsg, sm_marker);
Expand Down Expand Up @@ -137,9 +141,9 @@ dds_return_t ddsi_create_fragment_message (struct ddsi_writer *wr, ddsi_seqno_t
Note: fragnum is 0-based here, 1-based in DDSI. But 0-based is
much easier ...

actual expected_inline_qos_size is typically 0, but always claiming 32 bytes won't make
actual expected_inline_qos_size is typically 0, but always claiming 60 bytes won't make
a difference, so no point in being precise */
const size_t expected_inline_qos_size = /* statusinfo */ 8 + /* keyhash */ 20 + /* sentinel */ 4;
const size_t expected_inline_qos_size = /* statusinfo */ 8 + /* keyhash */ 20 + /* related_sample_identity */ 28 + /* sentinel */ 4;
struct ddsi_domaingv const * const gv = wr->e.gv;
struct ddsi_xmsg_marker sm_marker;
void *sm;
Expand All @@ -163,7 +167,7 @@ dds_return_t ddsi_create_fragment_message (struct ddsi_writer *wr, ddsi_seqno_t

fragging = (nfrags * (uint32_t) gv->config.fragment_size < size);

/* INFO_TS: 12 bytes, ddsi_rtps_datafrag_t: 36 bytes, expected inline QoS: 32 => should be single chunk */
/* INFO_TS: 12 bytes, ddsi_rtps_datafrag_t: 36 bytes, expected inline QoS: 60 => should be single chunk */
if ((*pmsg = ddsi_xmsg_new (gv->xmsgpool, &wr->e.guid, wr->c.pp, sizeof (ddsi_rtps_info_ts_t) + sizeof (ddsi_rtps_datafrag_t) + expected_inline_qos_size, xmsg_kind)) == NULL)
return DDS_RETCODE_OUT_OF_RESOURCES;

Expand Down Expand Up @@ -251,6 +255,8 @@ dds_return_t ddsi_create_fragment_message (struct ddsi_writer *wr, ddsi_seqno_t
if (fragnum == 0)
{
int rc;
ddsi_guid_t related_writer_guid;
ddsi_seqno_t related_seq;
/* Adding parameters means potential reallocing, so sm, ddcmn now likely become invalid */
if (wr->num_readers_requesting_keyhash > 0)
{
Expand All @@ -260,6 +266,10 @@ dds_return_t ddsi_create_fragment_message (struct ddsi_writer *wr, ddsi_seqno_t
{
ddsi_xmsg_addpar_statusinfo (*pmsg, serdata->statusinfo);
}
if (ddsi_serdata_get_related_sample_identity (serdata, &related_writer_guid, &related_seq))
{
ddsi_xmsg_addpar_related_sample_identity (*pmsg, &related_writer_guid, related_seq);
}
rc = ddsi_xmsg_addpar_sentinel_ifparam (*pmsg);
if (rc > 0)
{
Expand Down
9 changes: 9 additions & 0 deletions src/core/ddsi/src/ddsi_xmsg.c
Original file line number Diff line number Diff line change
Expand Up @@ -958,6 +958,15 @@ void ddsi_xmsg_addpar_keyhash (struct ddsi_xmsg *m, const struct ddsi_serdata *s
}
}

void ddsi_xmsg_addpar_related_sample_identity (struct ddsi_xmsg *m, const ddsi_guid_t *writer_guid, ddsi_seqno_t seq)
{
char *p = ddsi_xmsg_addpar (m, DDSI_PID_RELATED_SAMPLE_IDENTITY, sizeof (ddsi_guid_t) + sizeof (ddsi_sequence_number_t));
const ddsi_guid_t wire_guid = ddsi_hton_guid (*writer_guid);
const ddsi_sequence_number_t wire_seq = ddsi_to_seqno (seq);
memcpy (p, &wire_guid, sizeof (wire_guid));
memcpy (p + sizeof (wire_guid), &wire_seq, sizeof (wire_seq));
}

static void ddsi_xmsg_addpar_BE4u (struct ddsi_xmsg *m, ddsi_parameterid_t pid, uint32_t x)
{
unsigned *p = ddsi_xmsg_addpar (m, pid, sizeof (x));
Expand Down
1 change: 1 addition & 0 deletions src/core/ddsi/tests/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ set(ddsi_test_sources
"plist_generic.c"
"plist.c"
"plist_leasedur.c"
"related_sample_identity.c"
"pmd_message.c"
"radmin.c"
"receive_packet.c"
Expand Down
Loading