Code
Hub
Workspaces
Following
Trending
Connect
MCP
copy
Create free account
hub
/
github.com/densumesh/broccoli
/ functions
Functions
186 in github.com/densumesh/broccoli
⨍
Functions
186
◇
Types & classes
41
↓ 116 callers
Method
clone
(&self)
src/queue.rs:417
↓ 48 callers
Method
publish
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 callers
Method
build
Builds the `BroccoliQueue` with the specified configuration. # Returns A `Result` containing the `BroccoliQueue` on success, or a `BroccoliError` on
src/queue.rs:196
↓ 33 callers
Function
setup_queue
()
tests/common/mod.rs:3
↓ 27 callers
Method
acknowledge
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 callers
Method
fairness
(mut self, fairness: bool)
src/queue.rs:282
↓ 21 callers
Function
get_redis_client
()
tests/common/mod.rs:24
↓ 18 callers
Method
publish_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 callers
Method
queue_status
( &self, queue_name: String, disambiguator: Option<String>, )
src/queue.rs:1097
↓ 11 callers
Method
into_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 callers
Function
transaction_error
(e: &surrealdb::Error, msg: String)
src/brokers/surrealdb/utils.rs:731
↓ 7 callers
Method
check_connected
check and return current active connection
src/brokers/surrealdb/utils.rs:60
↓ 7 callers
Method
get_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 callers
Method
pool_connections
(mut self, connections: u8)
src/queue.rs:168
↓ 7 callers
Method
reject
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 callers
Function
client_from_url
helper: public so we can call it from testing
src/brokers/surrealdb/utils.rs:89
↓ 4 callers
Method
consume
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 callers
Method
ensure_pool
(&self)
src/brokers/rabbitmq/utils.rs:25
↓ 4 callers
Function
generate_test_messages
(queue_name: &str, n: usize)
benches/surrealdb_benchmark.rs:88
↓ 4 callers
Function
get_param_value
helper: given a url get a named parameter
src/brokers/surrealdb/utils.rs:157
↓ 4 callers
Function
message_record_id
`queue_name` + task id : namely `queue_name:<uuid>task_id`
src/brokers/surrealdb/utils.rs:235
↓ 4 callers
Method
priority
(mut self, priority: u8)
src/queue.rs:381
↓ 4 callers
Function
processing_table
table holds messages in process
src/brokers/surrealdb/utils.rs:218
↓ 4 callers
Function
queue_table
(queue_name: &str)
src/brokers/surrealdb/utils.rs:213
↓ 4 callers
Function
read_param
(s: &str)
benches/surrealdb_benchmark.rs:32
↓ 3 callers
Function
add_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 callers
Method
auto_ack
(mut self, auto_ack: bool)
src/queue.rs:275
↓ 3 callers
Function
index_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 callers
Method
is_done
we either have a result or we have exhausted the number of retries
src/brokers/surrealdb/utils.rs:1062
↓ 3 callers
Method
process_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 callers
Method
setup_queue
( &self, channel: &Channel, queue_name: &str, )
src/brokers/rabbitmq/utils.rs:70
↓ 3 callers
Method
size
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 callers
Method
step
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 callers
Method
wrapup
wrapup, either get the original result, wrapped in a broccoli error or a too many retries error
src/brokers/surrealdb/utils.rs:1089
↓ 2 callers
Function
add_to_queue
add to the end of the timeseries queue without scheduling
src/brokers/surrealdb/utils.rs:284
↓ 2 callers
Method
cancel
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 callers
Method
connect
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 callers
Function
consume_loop
( queue: &BroccoliQueue, queue_name: &str, message_count: usize, )
benches/surrealdb_benchmark.rs:99
↓ 2 callers
Method
delay
(mut self, duration: Duration)
src/queue.rs:364
↓ 2 callers
Method
ensure_pool
(&self)
src/brokers/redis/utils.rs:40
↓ 2 callers
Method
failed_message_retry_strategy
(mut self, strategy: RetryStrategy)
src/queue.rs:154
↓ 2 callers
Function
get_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 callers
Function
get_surrealdb_client
()
tests/common/mod.rs:31
↓ 2 callers
Function
index_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 callers
Function
process_job
(job: JobPayload)
examples/consumer.rs:20
↓ 2 callers
Method
process_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 callers
Function
remove_from_queue
remove from ordered queue `queued_message_id` must be: queue:[priority, timestamp, `task_id`]
src/brokers/surrealdb/utils.rs:589
↓ 2 callers
Function
remove_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 callers
Function
setup_broccoli
(url: String)
benches/surrealdb_benchmark.rs:69
↓ 2 callers
Function
setup_surrealdb
setup the connection to the database if `url` is informed or create an in-memory instance otherwise
benches/surrealdb_benchmark.rs:37
↓ 2 callers
Function
target_counter
(message_count: usize)
benches/surrealdb_benchmark.rs:77
↓ 2 callers
Method
try_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 callers
Method
try_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 callers
Method
ttl
(mut self, duration: Duration)
src/queue.rs:357
↓ 1 callers
Function
add_message
add the message itself with it's payload
src/brokers/surrealdb/utils.rs:761
↓ 1 callers
Function
add_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 callers
Function
add_to_failed
add to the failed queue, will also remove from index
src/brokers/surrealdb/utils.rs:1002
↓ 1 callers
Function
add_to_queue_delayed
add to the end of the timeseries queue with a delay duration
src/brokers/surrealdb/utils.rs:297
↓ 1 callers
Function
add_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 callers
Function
benchmark_broccoli_batch_publish_consume_throughput
( queue: &BroccoliQueue, message_count: usize, )
benches/surrealdb_benchmark.rs:267
↓ 1 callers
Function
benchmark_broccoli_consume_loop_throughput
( queue: &BroccoliQueue, message_count: usize, )
benches/surrealdb_benchmark.rs:304
↓ 1 callers
Function
benchmark_broccoli_throughput
(queue: &BroccoliQueue, message_count: usize)
benches/redis_benchmark.rs:75
↓ 1 callers
Function
benchmark_broccoli_throughput
(queue: &BroccoliQueue, message_count: usize)
benches/amqp_benchmark.rs:67
↓ 1 callers
Function
benchmark_raw_amqp_throughput
(conn: &mut Connection, message_count: usize)
benches/amqp_benchmark.rs:28
↓ 1 callers
Function
benchmark_raw_redis_throughput
( conn: &mut redis::aio::MultiplexedConnection, message_count: usize, )
benches/redis_benchmark.rs:28
↓ 1 callers
Function
connect_to_broker
( broker_url: &str, config: Option<BrokerConfig>, )
src/brokers/connect.rs:25
↓ 1 callers
Method
consume_wait
(mut self, consume_wait: std::time::Duration)
src/queue.rs:289
↓ 1 callers
Method
enable_scheduling
(mut self, enable_scheduling: bool)
src/queue.rs:184
↓ 1 callers
Function
error_handler
(_: TestMessage, err: BroccoliError)
tests/happy_path.rs:947
↓ 1 callers
Function
error_handler
(msg: JobPayload, err: BroccoliError)
examples/consumer.rs:34
↓ 1 callers
Function
failed_table
this is the failed messages table
src/brokers/surrealdb/utils.rs:229
↓ 1 callers
Function
get_message
get the actual message `message_id` <`queue_table>`:[<`task_id`>]
src/brokers/surrealdb/utils.rs:845
↓ 1 callers
Function
get_message_from
get the message payload given the queue record
src/brokers/surrealdb/utils.rs:833
↓ 1 callers
Function
get_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 callers
Method
get_queue_status
( &self, queue_name: String, disambiguator: Option<String>, )
src/brokers/redis/management.rs:14
↓ 1 callers
Function
get_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 callers
Method
get_task_id
( &self, queue_name: &str, redis_connection: &mut MultiplexedConnection, optio
src/brokers/redis/utils.rs:54
↓ 1 callers
Function
process_handler
(m: TestMessage)
tests/happy_path.rs:933
↓ 1 callers
Function
process_job
(m: TestMessage)
tests/edge_cases.rs:489
↓ 1 callers
Function
process_job
(m: TestMessage)
tests/happy_path.rs:925
↓ 1 callers
Method
publish
( &self, queue_name: &str, _disambiguator: Option<String>, messages: &[Interna
src/brokers/rabbitmq/broker.rs:78
↓ 1 callers
Method
publish
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 callers
Function
queue_record_id
time+id range record id, namely `queue_table:[priority, when,<uuid>task_id]`
src/brokers/surrealdb/utils.rs:253
↓ 1 callers
Function
remove_from_processing
remove from the processing queue
src/brokers/surrealdb/utils.rs:892
↓ 1 callers
Function
remove_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 callers
Function
remove_from_queue_add_to_processed_transaction_impl
( db: &Surreal<Any>, queue_name: &str, queued_message: InternalSurrealDBBrokerMessageEntry, er
src/brokers/surrealdb/utils.rs:609
↓ 1 callers
Function
remove_from_queue_index
clear the index table
src/brokers/surrealdb/utils.rs:423
↓ 1 callers
Function
remove_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 callers
Function
remove_queued_from_index
given the user facing task id, remove from the queue
src/brokers/surrealdb/utils.rs:743
↓ 1 callers
Method
schedule_at
(mut self, time: OffsetDateTime)
src/queue.rs:371
↓ 1 callers
Function
setup_amqp
()
benches/amqp_benchmark.rs:16
↓ 1 callers
Function
setup_broccoli
()
benches/redis_benchmark.rs:20
↓ 1 callers
Function
setup_broccoli
()
benches/amqp_benchmark.rs:20
↓ 1 callers
Method
setup_exchange
( &self, channel: &Channel, exchange_name: &str, )
src/brokers/rabbitmq/utils.rs:36
↓ 1 callers
Function
setup_queue_with_url
( url: &str, )
tests/common/mod.rs:14
↓ 1 callers
Function
setup_redis
()
benches/redis_benchmark.rs:15
↓ 1 callers
Function
success_handler
(m: TestMessage)
tests/happy_path.rs:940
↓ 1 callers
Function
success_handler
(msg: JobPayload)
examples/consumer.rs:29
↓ 1 callers
Function
to_rfc3339
used to parse dates in surrealdb format
benches/surrealdb_benchmark.rs:131
↓ 1 callers
Function
update_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