Skip to article frontmatterSkip to article content
Site not loading correctly?

This may be due to an incorrect BASE_URL configuration. See the MyST Documentation for reference.

tr_queue

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

KeyTypeDescription
jn_messagejson_t *A JSON object representing the original message. The ownership of this object is not transferred.
resultintAn 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

KeyTypeDescription
trqtr_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=1000000

Up 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 up

Up 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 openThe 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 jsonthe file changes (inode, mode, size, mtime, ctime, taken with stat() just before each open)
nothing the queue can seetopic_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

KeyTypeDescription
trqtr_queue_t *The queue instance to check.
__t__uint64_tTime value of the message.
rowiduint64_tThe 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

KeyTypeDescription
trqtr_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

KeyTypeDescription
trqtr_queue_t *The queue instance from which to retrieve the message.
rowiduint64_tThe 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

KeyTypeDescription
kwjson_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

KeyTypeDescription
trqtr_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

KeyTypeDescription
trqtr_queue_t *The queue instance from which messages will be loaded.
from_rowidint64_tThe starting rowid for loading messages.
to_rowidint64_tThe 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

KeyTypeDescription
msgq_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

KeyTypeDescription
trangerjson_t *Pointer to the tranger instance managing the queue.
topic_nameconst char *Name of the topic associated with the queue.
tkeyconst char *Time key used for ordering messages in the queue.
system_flagsystem_flag2_tSystem flags controlling queue behavior.
backup_queue_sizesize_tMaximum 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

KeyTypeDescription
msgq_msg_t *The message to be marked.
hard_markuint16_tThe hard flag to set on the message.
setBOOLIf 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

KeyTypeDescription
kwjson_t *The JSON object where the metadata will be stored.
keyconst char *The key under which the metadata value will be stored.
jn_valuejson_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

KeyTypeDescription
msgq_msg_t *The queue message on which the soft mark is to be set or cleared.
soft_markuint64_tThe soft mark value to be applied to the message.
setBOOLIf 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

KeyTypeDescription
msgq_msg_t *The message to be unloaded from the queue iterator.
resultint32_tThe 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

KeyTypeDescription
trqtr_queue_t *The queue instance to append the message to.
tjson_int_tTimestamp for the message. Pass 0 to use the current time.
kwjson_t *The JSON payload of the message. Ownership is transferred to the queue.
user_flaguint16_tOptional 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

KeyTypeDescription
trqtr_queue_t *The queue instance from which messages will be loaded.
from_tint64_tThe starting timestamp for the time range.
to_tint64_tThe 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

KeyTypeDescription
msgq_msg_t *The queue message whose metadata is to be retrieved.

Returns

Returns a md2_record_ex_t * pointer to the internal metadata record.