Persistent message queue with at-least-once semantics, built on timeranger2. Used by the MQTT broker for durable subscriptions.
Source code:
trq_answer()¶
trq_answer() extracts and returns a new JSON object containing only the metadata from the given message.
json_t *trq_answer(
json_t *jn_message, // not owned, Gps message, to get only __MD_TRQ__
int result
);Parameters
| Key | Type | Description |
|---|---|---|
jn_message | json_t * | A JSON object representing the original message. The ownership of this object is not transferred. |
result | int | An integer result code that can be included in the returned metadata. |
Returns
A new JSON object containing only the metadata extracted from jn_message. The caller assumes ownership of the returned object.
Notes
The function is specifically designed to extract the __MD_TRQ__ metadata field from the input message.
trq_check_backup()¶
trq_check_backup() backs the queue’s topic up when it has grown to the
backup_queue_size given to trq_open(): the topic is moved to
a backup and re-created EMPTY, and first_rowid is reset. It is meant to be
called only when the queue holds no message (C_QIOGATE calls it every
timeout_backup seconds when trq_size() is 0 and no ack is pending).
int trq_check_backup(
tr_queue_t * trq
);Parameters
| Key | Type | Description |
|---|---|---|
trq | tr_queue_t * | The queue instance to check and perform backup if needed. |
Returns
0 when there was nothing to do or the backup was done. -1 when the backup
is REFUSED: the last trq_load() did not read every pending
message (load_failed is TRUE in the tr_queue_t). The queue is then empty
because its load failed, not because its messages were processed, and the
backup would take for good the pending messages the load could not read --
the ones trq_load() keeps first_rowid for. The refusal is said once, with
an ERROR, and lasts until a trq_load() reads the queue whole (a restart
after the repair):
ERROR trq_check_backup: Queue backup refused: its last load did not read every pending message
topic_name=emails topic_size=1000000 backup_queue_size=1000000Up to 7.25.4 the periodic backup of a
queue whose load had failed re-created the topic empty. tr2q_check_backup()
of the mqtt queues behaves the same after a failed tr2q_load().
-1 also when the backup itself FAILS
(tranger2_backup_topic() answers NULL: the
backup name is taken by a file, the rename() fails, the topic’s files do
not load, or the new topic cannot be created after the move -- no space, a
mkdir that fails -- and the backup is moved back). The queue goes on in its
topic, not backed up, which the backup opened again, and the next period tries
again:
ERROR tranger2_backup_topic: cannot backup topic errno=20 serrno="Not a directory"
WARN reopen_topic_not_backed_up: Backup of topic failed: the topic is opened again as it was, not backed up
ERROR trq_check_backup: Queue backup failed: the queue goes on in its topic, not backed upUp to 7.25.4 the queue was left with no topic and the call answered 0:
every read logged “What topic?”, every ack answered -1 (the messages were
sent again after a restart), and no backup happened again; a create that
failed after the move left the messages in the backup. tr2q_check_backup()
behaves the same.
When the topic cannot be opened again either -- the cause of the failure is
still there, as a topic_desc.json that cannot be read for a moment -- the
queue has no topic: “Queue backup failed, and the queue has no topic”. It
takes its topic again BY NAME as soon as it can be opened: at the next call,
and at the next trq_msg_json() or ack
(trq_set_hard_flag()) of a message, with an INFO,
“Queue topic taken again”. While it cannot, those calls fail (-1, NULL)
and the queue says it once, “Queue without topic, it cannot be opened”,
after the causes the open logs. A topic the tranger already has open again
(an append opens it by name) is taken as it is. Otherwise the queue asks the
disk quietly first, and tries the open again only when what stopped the last
one may be gone. After a failed open it asks what that was, the way the open
reads the topic, and without a log:
| What stopped the open | The open is tried again when |
|---|---|
topic_desc.json cannot be opened (not there, no permission, no descriptor left), or keys/ cannot be listed (EACCES, EMFILE, EIO, a key whose type cannot be asked) | the file can be opened AND keys/ can be listed |
topic_desc.json opens and is not json | the file changes (inode, mode, size, mtime, ctime, taken with stat() just before each open) |
| nothing the queue can see | topic_desc.json or keys/ changes |
The next calls log nothing. Asked through tranger2_topic() at every call, a
topic that cannot be opened would log three errors each time, and the MQTT
broker calls tr2q_check_backup() every second for each session with nothing
in flight:
CRITICAL load_persistent_json: Cannot open a json file (first call)
ERROR tranger2_open_topic: Cannot open topic: topic_desc.json does not load
ERROR tranger2_topic: Cannot open topic
ERROR take_queue_topic: Queue without topic, it cannot be opened (once)
(next calls: nothing)
INFO take_queue_topic: Queue topic taken again (the file is back)A topic_desc.json that can be read and does not load (a broken json) is
said the same way, once, with the parser’s cause first:
CRITICAL load_json_from_fd: Cannot load json file, bad json (first call)
ERROR tranger2_open_topic: Cannot open topic: topic_desc.json does not load
ERROR tranger2_topic: Cannot open topic
ERROR take_queue_topic: Queue without topic, it cannot be opened (once)
(next calls: nothing)Each CHANGE of the file is tried once: a new content that still does not load logs the three errors of the open again (the new cause), but not the queue’s line; the good content takes the topic again with the INFO.
A keys/ that cannot be listed (here of mode 0) is said the same way, and
the topic is taken again once keys/ can be listed, although
topic_desc.json never changed:
ERROR find_keys_in_disk: Cannot list the keys of the topic (first call)
ERROR tranger2_open_topic: Cannot open topic: its keys cannot be listed
ERROR tranger2_topic: Cannot open topic
ERROR take_queue_topic: Queue without topic, it cannot be opened (once)
(next calls: nothing)
INFO take_queue_topic: Queue topic taken again (keys/ listed again)Up to 7.25.4 the queue never took its topic again, whatever the cause: it stayed without topic until a restart, every ack and read failed, and the acked messages were sent again after the restart.
The mqtt queues (tr2q_check_backup(), tr2q_msg_json(),
tr2q_save_hard_mark()) do the same.
tests/c/tr_queue/test_tr_queue_backup_failed.
Example
tr_queue_t *trq = trq_open(tranger, "emails", "tm", 0, 1000000);
trq_load(trq); // -1: a pending message could not be read
if(trq_size(trq) == 0) {
trq_check_backup(trq); // -1: refused, the topic keeps its messages
}trq_check_pending_rowid()¶
trq_check_pending_rowid() checks the pending status of a message identified by its rowid in the queue.
int trq_check_pending_rowid(
tr_queue_t * trq,
uint64_t __t__,
uint64_t rowid
);Parameters
| Key | Type | Description |
|---|---|---|
trq | tr_queue_t * | The queue instance to check. |
__t__ | uint64_t | Time value of the message. |
rowid | uint64_t | The unique row identifier of the message. |
Returns
Returns -1 if the rowid does not exist, 1 if the message is pending, and 0 if it is not pending.
Notes
This function provides a low-level check for message status in the queue.
trq_close()¶
Closes the given tr_queue, releasing associated resources. After calling trq_close(), make sure that to invoke tranger2_shutdown() if no other queues are in use.
void trq_close(
tr_queue_t * trq
);Parameters
| Key | Type | Description |
|---|---|---|
trq | tr_queue_t * | The queue instance to be closed. |
Returns
This function does not return a value.
Notes
Make sure that trq_close() is called before shutting down the underlying TimeRanger instance with tranger2_shutdown().
trq_get_by_rowid()¶
trq_get_by_rowid() retrieves a message from the queue iterator using its row ID.
q_msg_t * trq_get_by_rowid(
tr_queue_t * trq,
uint64_t rowid
);Parameters
| Key | Type | Description |
|---|---|---|
trq | tr_queue_t * | The queue instance from which to retrieve the message. |
rowid | uint64_t | The row ID of the message to retrieve. |
Returns
Returns a q_msg_t * handle to the retrieved message, or NULL if the message is not found.
Notes
The returned message remains owned by the queue and must not be freed manually.
trq_get_metadata()¶
Retrieves the metadata associated with a given JSON object. The returned JSON object is not owned by the caller.
json_t *trq_get_metadata(
json_t *kw
);Parameters
| Key | Type | Description |
|---|---|---|
kw | json_t * | The JSON object containing metadata. |
Returns
A pointer to a JSON object containing the metadata. The returned JSON object is not owned by the caller and must not be modified or freed.
Notes
The returned JSON object is a reference and must not be altered or deallocated by the caller.
trq_load()¶
trq_load() loads the pending messages of the queue (the records flagged TRQ_MSG_PENDING) into memory, metadata only: the content is read when a message asks for it (trq_msg_json()). The load starts at the queue’s first_rowid, saved in topic_var.json by the previous load, and saves the rowid of the first pending message it finds as the new one (the size of the topic when none is pending).
int trq_load(
tr_queue_t * trq
);Parameters
| Key | Type | Description |
|---|---|---|
trq | tr_queue_t * | The queue instance from which pending messages will be loaded. |
Returns
0 when every pending message was read. -1 when the queue is NULL, or when the load could not read every pending message (a row of the queue’s md2 that cannot be read: the list says load_failed, and the log “Queue loaded without some of its messages: its first_rowid is not moved nor saved”). The messages it read ARE in the queue; first_rowid is kept as it was, so a load after the store is repaired finds the ones this one missed. Up to 7.25.4 such a load saved the size of the topic as first_rowid, and those messages were skipped for ever. The same holds for tr2q_load() of the mqtt queues.
Example
tr_queue_t *trq = trq_open(tranger, "emails", "tm", 0, 0);
if(trq_load(trq) < 0) {
/* some pending messages are not in memory (see the log); the rest are */
}
q_msg_t *msg;
qmsg_foreach_forward(trq, msg) {
/* deliver it, then trq_unload_msg(msg, 0) */
}Notes
Use trq_load_all() to load all messages, including non-pending ones.
trq_load_all()¶
trq_load_all() loads all messages from the queue within the specified rowid range, optionally filtering by key.
int trq_load_all(
tr_queue_t * trq,
int64_t from_rowid,
int64_t to_rowid
);Parameters
| Key | Type | Description |
|---|---|---|
trq | tr_queue_t * | The queue instance from which messages will be loaded. |
from_rowid | int64_t | The starting rowid for loading messages. |
to_rowid | int64_t | The ending rowid for loading messages. |
Returns
Returns an iterator over the loaded messages or an error code if the operation fails.
Notes
Use trq_load_all() to retrieve messages efficiently within a specific rowid range.
trq_msg_json()¶
trq_msg_json() retrieves the JSON representation of a queue message. The returned JSON object is not owned by the caller and must not be modified or freed.
json_t *trq_msg_json(
q_msg_t *msg
);Parameters
| Key | Type | Description |
|---|---|---|
msg | q_msg_t * | The queue message whose JSON representation is to be retrieved. |
Returns
A pointer to a json_t object representing the message. The returned JSON object is not owned by the caller.
Notes
The returned JSON object must not be modified or freed by the caller.
trq_open()¶
trq_open() initializes and opens a persistent queue using the specified tranger instance and topic configuration.
tr_queue_t *trq_open(
json_t *tranger,
const char *topic_name,
const char *tkey,
system_flag2_t system_flag,
size_t backup_queue_size
);Parameters
| Key | Type | Description |
|---|---|---|
tranger | json_t * | Pointer to the tranger instance managing the queue. |
topic_name | const char * | Name of the topic associated with the queue. |
tkey | const char * | Time key used for ordering messages in the queue. |
system_flag | system_flag2_t | System flags controlling queue behavior. |
backup_queue_size | size_t | Maximum number of messages to retain in the backup queue. |
Returns
Returns a tr_queue_t * handle representing the opened queue, or NULL on failure.
Notes
Make sure that tranger2_startup() is called before invoking trq_open().
trq_set_hard_flag()¶
trq_set_hard_flag() marks a message with a hard flag. This allows it to be recovered in the next queue open if the flag is used in trq_load().
int trq_set_hard_flag(
q_msg_t *msg,
uint16_t hard_mark,
BOOL set
);Parameters
| Key | Type | Description |
|---|---|---|
msg | q_msg_t * | The message to be marked. |
hard_mark | uint16_t | The hard flag to set on the message. |
set | BOOL | If TRUE, the flag is set. If FALSE, the flag is cleared. |
Returns
Returns 0 on success, or a negative value on failure.
Notes
A message must be flagged after being appended to the queue if it needs to be recovered in the next queue open using trq_load().
trq_set_metadata()¶
trq_set_metadata() sets a metadata key-value pair in the given JSON object.
int trq_set_metadata(
json_t *kw,
const char *key,
json_t *jn_value // owned
);Parameters
| Key | Type | Description |
|---|---|---|
kw | json_t * | The JSON object where the metadata will be stored. |
key | const char * | The key under which the metadata value will be stored. |
jn_value | json_t * | The JSON value to be stored as metadata. Ownership is transferred. |
Returns
Returns 0 on success, or a negative value on failure.
Notes
The caller must make sure that kw is a valid JSON object before calling trq_set_metadata().
trq_set_soft_mark()¶
trq_set_soft_mark() sets or clears a soft mark on a given queue message.
uint64_t trq_set_soft_mark(
q_msg_t *msg,
uint64_t soft_mark,
BOOL set
);Parameters
| Key | Type | Description |
|---|---|---|
msg | q_msg_t * | The queue message on which the soft mark is to be set or cleared. |
soft_mark | uint64_t | The soft mark value to be applied to the message. |
set | BOOL | If TRUE, the soft mark is set. If FALSE, the soft mark is cleared. |
Returns
Returns the updated soft mark value of the message.
Notes
Soft marks are used for temporary message state tracking and do not persist across queue restarts.
trq_unload_msg()¶
The trq_unload_msg() function unloads a message from the queue iterator, removing it from memory.
void trq_unload_msg(
q_msg_t *msg,
int32_t result
);Parameters
| Key | Type | Description |
|---|---|---|
msg | q_msg_t * | The message to be unloaded from the queue iterator. |
result | int32_t | The result code associated with the message unloading operation. |
Returns
This function does not return a value.
Notes
Use trq_unload_msg() to free a message from the queue iterator after processing it.
trq_append2()¶
Appends a new message to the queue with an explicit timestamp and optional user flags.
q_msg_t * trq_append2(
tr_queue_t * trq,
json_int_t t,
json_t *kw,
uint16_t user_flag
);Parameters
| Key | Type | Description |
|---|---|---|
trq | tr_queue_t * | The queue instance to append the message to. |
t | json_int_t | Timestamp for the message. Pass 0 to use the current time. |
kw | json_t * | The JSON payload of the message. Ownership is transferred to the queue. |
user_flag | uint16_t | Optional user-defined flags to associate with the message. |
Returns
Returns a q_msg_t * handle to the appended message, or NULL on failure.
trq_load_all_by_time()¶
Loads all messages from the queue within a specified time range.
int trq_load_all_by_time(
tr_queue_t * trq,
int64_t from_t,
int64_t to_t
);Parameters
| Key | Type | Description |
|---|---|---|
trq | tr_queue_t * | The queue instance from which messages will be loaded. |
from_t | int64_t | The starting timestamp for the time range. |
to_t | int64_t | The ending timestamp for the time range. |
Returns
Returns 0 on success, or a negative value on error.
trq_msg_md()¶
Retrieves the metadata record associated with a queue message.
md2_record_ex_t *trq_msg_md(
q_msg_t *msg
);Parameters
| Key | Type | Description |
|---|---|---|
msg | q_msg_t * | The queue message whose metadata is to be retrieved. |
Returns
Returns a md2_record_ex_t * pointer to the internal metadata record.