MCPcopy Create free account

hub / github.com/densumesh/broccoli / functions

Functions186 in github.com/densumesh/broccoli

↓ 116 callersMethodclone
(&self)
src/queue.rs:417
↓ 48 callersMethodpublish
Publishes a message to the specified topic. # Arguments `topic` - The name of the topic. `message` - The message to be published. # Returns A `Resul
src/queue.rs:465
↓ 44 callersMethodbuild
Builds the `BroccoliQueue` with the specified configuration. # Returns A `Result` containing the `BroccoliQueue` on success, or a `BroccoliError` on
src/queue.rs:196
↓ 33 callersFunctionsetup_queue
()
tests/common/mod.rs:3
↓ 27 callersMethodacknowledge
Acknowledges the processing of a message, removing it from the processing queue. # Arguments `topic` - The name of the topic. `message` - The message
src/queue.rs:654
↓ 23 callersMethodfairness
(mut self, fairness: bool)
src/queue.rs:282
↓ 21 callersFunctionget_redis_client
()
tests/common/mod.rs:24
↓ 18 callersMethodpublish_batch
Publishes a batch of messages to the specified topic. # Arguments `topic` - The name of the topic. `messages` - An iterator over the messages to be p
src/queue.rs:496
↓ 13 callersMethodqueue_status
( &self, queue_name: String, disambiguator: Option<String>, )
src/queue.rs:1097
↓ 11 callersMethodinto_message
Converts the internal message to a `BrokerMessage`. # Returns A `Result` containing the `BrokerMessage` or a `BroccoliError` on failure. # Errors If
src/brokers/broker.rs:279
↓ 9 callersFunctiontransaction_error
(e: &surrealdb::Error, msg: String)
src/brokers/surrealdb/utils.rs:731
↓ 7 callersMethodcheck_connected
check and return current active connection
src/brokers/surrealdb/utils.rs:60
↓ 7 callersMethodget_redis_connection
Retrieves a Redis connection from the pool, retrying with exponential backoff if necessary. # Arguments `redis_pool` - A reference to the Redis conne
src/brokers/redis/utils.rs:162
↓ 7 callersMethodpool_connections
(mut self, connections: u8)
src/queue.rs:168
↓ 7 callersMethodreject
Rejects the processing of a message, moving it to the failed queue. # Arguments `topic` - The name of the topic. `message` - The message to be reject
src/queue.rs:680
↓ 4 callersFunctionclient_from_url
helper: public so we can call it from testing
src/brokers/surrealdb/utils.rs:89
↓ 4 callersMethodconsume
Consumes a message from the specified topic. This method will block until a message is available. This will not acknowledge the message, use `acknowle
src/queue.rs:537
↓ 4 callersMethodensure_pool
(&self)
src/brokers/rabbitmq/utils.rs:25
↓ 4 callersFunctiongenerate_test_messages
(queue_name: &str, n: usize)
benches/surrealdb_benchmark.rs:88
↓ 4 callersFunctionget_param_value
helper: given a url get a named parameter
src/brokers/surrealdb/utils.rs:157
↓ 4 callersFunctionmessage_record_id
`queue_name` + task id : namely `queue_name:<uuid>task_id`
src/brokers/surrealdb/utils.rs:235
↓ 4 callersMethodpriority
(mut self, priority: u8)
src/queue.rs:381
↓ 4 callersFunctionprocessing_table
table holds messages in process
src/brokers/surrealdb/utils.rs:218
↓ 4 callersFunctionqueue_table
(queue_name: &str)
src/brokers/surrealdb/utils.rs:213
↓ 4 callersFunctionread_param
(s: &str)
benches/surrealdb_benchmark.rs:32
↓ 3 callersFunctionadd_to_queue_scheduled
add to the timeseries at a scheduled time, can be in the past and it will be triggered immediately
src/brokers/surrealdb/utils.rs:319
↓ 3 callersMethodauto_ack
(mut self, auto_ack: bool)
src/queue.rs:275
↓ 3 callersFunctionindex_table
this is an index to go from <queue_name>___index:[u<taskid>,<queuename>] to the queue table in O(k) time
src/brokers/surrealdb/utils.rs:224
↓ 3 callersMethodis_done
we either have a result or we have exhausted the number of retries
src/brokers/surrealdb/utils.rs:1062
↓ 3 callersMethodprocess_messages
Processes messages from the specified topic with the provided handler function. # Example ```no_run use broccoli_queue::queue::BroccoliQueue; use br
src/queue.rs:755
↓ 3 callersMethodsetup_queue
( &self, channel: &Channel, queue_name: &str, )
src/brokers/rabbitmq/utils.rs:70
↓ 3 callersMethodsize
Returns the size of the queue(s). For fairness queues, returns a map with each disambiguator queue and its size. For unfair queues, returns a map wit
src/queue.rs:1076
↓ 3 callersMethodstep
take a surrealdb result and process it, updating internal state if the result is a retriable transaction, sleep for a small number of ms
src/brokers/surrealdb/utils.rs:1068
↓ 3 callersMethodwrapup
wrapup, either get the original result, wrapped in a broccoli error or a too many retries error
src/brokers/surrealdb/utils.rs:1089
↓ 2 callersFunctionadd_to_queue
add to the end of the timeseries queue without scheduling
src/brokers/surrealdb/utils.rs:284
↓ 2 callersMethodcancel
Cancels the processing of a message, deleting it from the queue. # Arguments `topic` - The name of the topic. `message_id` - The ID of the message to
src/queue.rs:704
↓ 2 callersMethodconnect
Connects to the Redis broker using the provided URL. # Arguments `broker_url` - The URL of the Redis broker. # Returns A `Result` indicating success
src/brokers/redis/broker.rs:33
↓ 2 callersFunctionconsume_loop
( queue: &BroccoliQueue, queue_name: &str, message_count: usize, )
benches/surrealdb_benchmark.rs:99
↓ 2 callersMethoddelay
(mut self, duration: Duration)
src/queue.rs:364
↓ 2 callersMethodensure_pool
(&self)
src/brokers/redis/utils.rs:40
↓ 2 callersMethodfailed_message_retry_strategy
(mut self, strategy: RetryStrategy)
src/queue.rs:154
↓ 2 callersFunctionget_queued_transaction
get first queued message payload if any, non-blocking, as part of a transaciton 1) get queued 2
src/brokers/surrealdb/utils.rs:456
↓ 2 callersFunctionget_surrealdb_client
()
tests/common/mod.rs:31
↓ 2 callersFunctionindex_record_id
`index_table`:[<uuid>`task_id`, `queue_name`] could simplify to index only by `task_id` but the queue name is useful for observability purposes
src/brokers/surrealdb/utils.rs:271
↓ 2 callersFunctionprocess_job
(job: JobPayload)
examples/consumer.rs:20
↓ 2 callersMethodprocess_messages_with_handlers
Processes messages from the specified topic with the provided handler functions for message processing, success, and error handling. # Example ```no
src/queue.rs:931
↓ 2 callersFunctionremove_from_queue
remove from ordered queue `queued_message_id` must be: queue:[priority, timestamp, `task_id`]
src/brokers/surrealdb/utils.rs:589
↓ 2 callersFunctionremove_message
remove actual message (the one with the payload) we also remove it from the internal index `message_id`: `<queue_table>`:[<`task_id`>]
src/brokers/surrealdb/utils.rs:868
↓ 2 callersFunctionsetup_broccoli
(url: String)
benches/surrealdb_benchmark.rs:69
↓ 2 callersFunctionsetup_surrealdb
setup the connection to the database if `url` is informed or create an in-memory instance otherwise
benches/surrealdb_benchmark.rs:37
↓ 2 callersFunctiontarget_counter
(message_count: usize)
benches/surrealdb_benchmark.rs:77
↓ 2 callersMethodtry_consume
Attempts to consume a message from the specified topic. This method will not block, returning immediately if no message is available. This will not ac
src/queue.rs:596
↓ 2 callersMethodtry_consume
Attempts to consume a message from the specified queue. Will not block if no message is available. This will check for scheduled messages first and th
src/brokers/redis/broker.rs:189
↓ 2 callersMethodttl
(mut self, duration: Duration)
src/queue.rs:357
↓ 1 callersFunctionadd_message
add the message itself with it's payload
src/brokers/surrealdb/utils.rs:761
↓ 1 callersFunctionadd_record_to_queue
internal implementation to add the record to the timeseries queue + index 1) we add the index first 2) we add the queue entry last, as that is what co
src/brokers/surrealdb/utils.rs:334
↓ 1 callersFunctionadd_to_failed
add to the failed queue, will also remove from index
src/brokers/surrealdb/utils.rs:1002
↓ 1 callersFunctionadd_to_queue_delayed
add to the end of the timeseries queue with a delay duration
src/brokers/surrealdb/utils.rs:297
↓ 1 callersFunctionadd_to_queue_index
we add an entry into a queue index, index:[messageid,queue_name] {queue_id} where queue_id is basically: queue:[timestamp,messageid]. We can use this
src/brokers/surrealdb/utils.rs:377
↓ 1 callersFunctionbenchmark_broccoli_batch_publish_consume_throughput
( queue: &BroccoliQueue, message_count: usize, )
benches/surrealdb_benchmark.rs:267
↓ 1 callersFunctionbenchmark_broccoli_consume_loop_throughput
( queue: &BroccoliQueue, message_count: usize, )
benches/surrealdb_benchmark.rs:304
↓ 1 callersFunctionbenchmark_broccoli_throughput
(queue: &BroccoliQueue, message_count: usize)
benches/redis_benchmark.rs:75
↓ 1 callersFunctionbenchmark_broccoli_throughput
(queue: &BroccoliQueue, message_count: usize)
benches/amqp_benchmark.rs:67
↓ 1 callersFunctionbenchmark_raw_amqp_throughput
(conn: &mut Connection, message_count: usize)
benches/amqp_benchmark.rs:28
↓ 1 callersFunctionbenchmark_raw_redis_throughput
( conn: &mut redis::aio::MultiplexedConnection, message_count: usize, )
benches/redis_benchmark.rs:28
↓ 1 callersFunctionconnect_to_broker
( broker_url: &str, config: Option<BrokerConfig>, )
src/brokers/connect.rs:25
↓ 1 callersMethodconsume_wait
(mut self, consume_wait: std::time::Duration)
src/queue.rs:289
↓ 1 callersMethodenable_scheduling
(mut self, enable_scheduling: bool)
src/queue.rs:184
↓ 1 callersFunctionerror_handler
(_: TestMessage, err: BroccoliError)
tests/happy_path.rs:947
↓ 1 callersFunctionerror_handler
(msg: JobPayload, err: BroccoliError)
examples/consumer.rs:34
↓ 1 callersFunctionfailed_table
this is the failed messages table
src/brokers/surrealdb/utils.rs:229
↓ 1 callersFunctionget_message
get the actual message `message_id` <`queue_table>`:[<`task_id`>]
src/brokers/surrealdb/utils.rs:845
↓ 1 callersFunctionget_message_from
get the message payload given the queue record
src/brokers/surrealdb/utils.rs:833
↓ 1 callersFunctionget_queue_index
get the index message given a queue name and task id, in O(k) time None is a valid return value if the message was never in the system in the first pl
src/brokers/surrealdb/utils.rs:406
↓ 1 callersMethodget_queue_status
( &self, queue_name: String, disambiguator: Option<String>, )
src/brokers/redis/management.rs:14
↓ 1 callersFunctionget_queued_transaction_impl
( db: &Surreal<Any>, queue_name: &str, auto_ack: bool, batch_size: usize, err_msg: &'stati
src/brokers/surrealdb/utils.rs:485
↓ 1 callersMethodget_task_id
( &self, queue_name: &str, redis_connection: &mut MultiplexedConnection, optio
src/brokers/redis/utils.rs:54
↓ 1 callersFunctionprocess_handler
(m: TestMessage)
tests/happy_path.rs:933
↓ 1 callersFunctionprocess_job
(m: TestMessage)
tests/edge_cases.rs:489
↓ 1 callersFunctionprocess_job
(m: TestMessage)
tests/happy_path.rs:925
↓ 1 callersMethodpublish
( &self, queue_name: &str, _disambiguator: Option<String>, messages: &[Interna
src/brokers/rabbitmq/broker.rs:78
↓ 1 callersMethodpublish
Publishes a message to the specified queue. # Arguments `queue_name` - The name of the queue. `message` - The message to be published. # Returns A `
src/brokers/redis/broker.rs:62
↓ 1 callersFunctionqueue_record_id
time+id range record id, namely `queue_table:[priority, when,<uuid>task_id]`
src/brokers/surrealdb/utils.rs:253
↓ 1 callersFunctionremove_from_processing
remove from the processing queue
src/brokers/surrealdb/utils.rs:892
↓ 1 callersFunctionremove_from_queue_add_to_processed_transaction
remove from ordered queue and add to in process list within the same transaction (unused at the moment as it allows for concurrent reads) `queued_mess
src/brokers/surrealdb/utils.rs:705
↓ 1 callersFunctionremove_from_queue_add_to_processed_transaction_impl
( db: &Surreal<Any>, queue_name: &str, queued_message: InternalSurrealDBBrokerMessageEntry, er
src/brokers/surrealdb/utils.rs:609
↓ 1 callersFunctionremove_from_queue_index
clear the index table
src/brokers/surrealdb/utils.rs:423
↓ 1 callersFunctionremove_message_and_from_processing_transaction
( db: &Surreal<Any>, queue_name: &str, task_id: &str, err_msg: &'static str, )
src/brokers/surrealdb/utils.rs:923
↓ 1 callersFunctionremove_queued_from_index
given the user facing task id, remove from the queue
src/brokers/surrealdb/utils.rs:743
↓ 1 callersMethodschedule_at
(mut self, time: OffsetDateTime)
src/queue.rs:371
↓ 1 callersFunctionsetup_amqp
()
benches/amqp_benchmark.rs:16
↓ 1 callersFunctionsetup_broccoli
()
benches/redis_benchmark.rs:20
↓ 1 callersFunctionsetup_broccoli
()
benches/amqp_benchmark.rs:20
↓ 1 callersMethodsetup_exchange
( &self, channel: &Channel, exchange_name: &str, )
src/brokers/rabbitmq/utils.rs:36
↓ 1 callersFunctionsetup_queue_with_url
( url: &str, )
tests/common/mod.rs:14
↓ 1 callersFunctionsetup_redis
()
benches/redis_benchmark.rs:15
↓ 1 callersFunctionsuccess_handler
(m: TestMessage)
tests/happy_path.rs:940
↓ 1 callersFunctionsuccess_handler
(msg: JobPayload)
examples/consumer.rs:29
↓ 1 callersFunctionto_rfc3339
used to parse dates in surrealdb format
benches/surrealdb_benchmark.rs:131
↓ 1 callersFunctionupdate_message
update the message, done to update the number of attempts, leaving rest unchanged
src/brokers/surrealdb/utils.rs:809
next →1–100 of 186, ranked by callers