MCPcopy Create free account

hub / github.com/densumesh/broccoli / functions

Functions186 in github.com/densumesh/broccoli

Methodacknowledge
Acknowledges the processing of a message, removing it from the processing queue. # Arguments `queue_name` - The name of the queue. `message` - The me
src/brokers/surrealdb/broker.rs:346
Methodacknowledge
( &self, _queue_name: &str, message: InternalBrokerMessage, )
src/brokers/rabbitmq/broker.rs:302
Methodacknowledge
Acknowledges the processing of a message, removing it from the processing queue. # Arguments `queue_name` - The name of the queue. `message` - The me
src/brokers/redis/broker.rs:257
Functionbenchmark_broccoli_batch_handler_throughput
( queue: &BroccoliQueue, options: Option<ConsumeOptions>, message_count: usize, )
benches/surrealdb_benchmark.rs:373
Functionbenchmark_raw_surrealdb_throughput
(db: &Surreal<Any>, message_count: usize)
benches/surrealdb_benchmark.rs:139
Methodbuilder
()
src/queue.rs:247
Methodbuilder_with
( db: surrealdb::Surreal<surrealdb::engine::any::Any>, )
src/queue.rs:448
Methodcancel
Cancels a message, removing it from the processing queue. # Arguments `queue_name` - The name of the queue. `message_id` - The ID of the message to b
src/brokers/surrealdb/broker.rs:463
Methodcancel
(&self, _queue_name: &str, _message_id: String)
src/brokers/rabbitmq/broker.rs:376
Methodcancel
Cancels a message, removing it from the queue. # Arguments `queue_name` - The name of the queue. `message_id` - The ID of the message to be canceled.
src/brokers/redis/broker.rs:379
Methodclient_from_url
we create a surreadlb connection from the url configuration URL parameters after ? are username, password, ns, database if unspecified, will default b
src/brokers/surrealdb/utils.rs:80
Methodconnect
Connects to the broker using the provided URL. # Arguments `broker_url` - The URL of the broker, in URL format with params, namely: `<protocol>://<ho
src/brokers/surrealdb/broker.rs:71
Methodconnect
(&mut self, broker_url: &str)
src/brokers/rabbitmq/broker.rs:37
Methodconsume
Consumes a message from the specified queue, blocking until a message is available. Uses live querying, so if there are no messages yet, it will block
src/brokers/surrealdb/broker.rs:243
Methodconsume
( &self, queue_name: &str, options: Option<ConsumeOptions>, )
src/brokers/rabbitmq/broker.rs:230
Methodconsume
Consumes a message from the specified queue, blocking until a message is available. # Arguments `queue_name` - The name of the queue. # Returns A `R
src/brokers/redis/broker.rs:222
Methodconsume_batch
Consumes a batch of messages from the specified topic. This method will block until the specified number of messages are consumed. This will not ackno
src/queue.rs:565
Functioncriterion_benchmark
(c: &mut Criterion)
benches/surrealdb_benchmark.rs:436
Functioncriterion_benchmark
(c: &mut Criterion)
benches/redis_benchmark.rs:109
Functioncriterion_benchmark
(c: &mut Criterion)
benches/amqp_benchmark.rs:101
Methoddefault
Creates a default retry strategy with 3 attempts and retries enabled.
src/queue.rs:32
Methoddefault
()
src/brokers/broker.rs:162
Methoddefault
()
src/brokers/surrealdb/utils.rs:33
Methodfrom
(msg: BrokerMessage<T>)
src/brokers/broker.rs:248
Methodfrom
(val: InternalSurrealDBBrokerMessage)
src/brokers/surrealdb/utils.rs:174
Methodfrom_redis_value
(v: &redis::Value)
src/brokers/redis/utils.rs:237
Methodget_queue_status
( &self, queue_name: String, disambiguator: Option<String>, )
src/brokers/rabbitmq/management.rs:11
Methodhandler_ack
(mut self, followup: bool)
src/queue.rs:296
Functionmain
()
examples/publisher.rs:24
Functionmain
()
examples/consumer.rs:40
Methodnew
()
src/queue.rs:46
Methodnew
Creates a new `BrokerMessage` with the provided payload.
src/brokers/broker.rs:195
Methodnew
()
src/brokers/surrealdb/utils.rs:41
Methodnew
()
src/brokers/rabbitmq/utils.rs:12
Methodnew
()
src/brokers/redis/utils.rs:19
Methodnew_with_config
(config: BrokerConfig)
src/brokers/surrealdb/utils.rs:51
Methodnew_with_config
(config: BrokerConfig)
src/brokers/rabbitmq/utils.rs:16
Methodnew_with_config
(config: BrokerConfig)
src/brokers/redis/utils.rs:32
Methodnew_with_surrealdb
(db: surrealdb::Surreal<surrealdb::engine::any::Any>)
src/queue.rs:134
Functionprocess_job
(m: BenchmarkMessage)
benches/surrealdb_benchmark.rs:338
Methodpublish
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/surrealdb/broker.rs:93
Methodreject
Rejects a message, re-queuing it or moving it to a failed queue if the retry limit is reached. # Arguments `queue_name` - The name of the queue. `mes
src/brokers/surrealdb/broker.rs:371
Methodreject
( &self, queue_name: &str, message: InternalBrokerMessage, )
src/brokers/rabbitmq/broker.rs:337
Methodreject
Rejects a message, re-queuing it or moving it to a failed queue if the retry limit is reached. # Arguments `queue_name` - The name of the queue. `mes
src/brokers/redis/broker.rs:288
Methodretry_failed
(mut self, retry_failed: bool)
src/queue.rs:75
Methodsize
(&self, _queue_name: &str)
src/brokers/surrealdb/broker.rs:493
Methodsize
(&self, _queue_name: &str)
src/brokers/rabbitmq/broker.rs:380
Methodsize
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/brokers/redis/broker.rs:401
Functiontest_batch_publish_and_consume
()
tests/happy_path.rs:90
Functiontest_concurrent_consume
()
tests/edge_cases.rs:130
Functiontest_delayed_message
()
tests/happy_path.rs:228
Functiontest_empty_payload
()
tests/edge_cases.rs:27
Functiontest_fairness_round_robin
()
tests/fairness.rs:22
Functiontest_fairness_with_delayed_messages
()
tests/fairness.rs:198
Functiontest_fairness_with_priorities
()
tests/fairness.rs:105
Functiontest_fairness_with_retries
()
tests/fairness.rs:310
Functiontest_invalid_broker_url
()
tests/edge_cases.rs:21
Functiontest_message_acknowledgment
()
tests/happy_path.rs:448
Functiontest_message_auto_ack
()
tests/happy_path.rs:538
Functiontest_message_cancellation
()
tests/happy_path.rs:604
Functiontest_message_ordering
()
tests/edge_cases.rs:288
Functiontest_message_priority
()
tests/happy_path.rs:674
Functiontest_message_retry
()
tests/happy_path.rs:372
Functiontest_multiple_batch_publish_and_consume
()
tests/edge_cases.rs:425
Functiontest_multiple_batch_publish_and_handler
()
tests/edge_cases.rs:501
Functiontest_process_messages
()
tests/happy_path.rs:953
Functiontest_process_messages_with_handlers
()
tests/happy_path.rs:1015
Functiontest_publish_and_consume
()
tests/happy_path.rs:22
Functiontest_queue_size
()
tests/happy_path.rs:813
Functiontest_queue_status_empty_name_error
()
tests/management.rs:403
Functiontest_queue_status_fairness_queue
()
tests/management.rs:103
Functiontest_queue_status_main_queue
()
tests/management.rs:15
Functiontest_queue_status_processing_and_failed
()
tests/management.rs:238
Functiontest_queue_status_specific_queue_lookup
()
tests/management.rs:307
Functiontest_redis_specific_queue_structure
()
tests/edge_cases.rs:578
Functiontest_scheduled_message
()
tests/happy_path.rs:302
Functiontest_try_consume_batch
()
tests/happy_path.rs:159
Functiontest_ttl_not_implemented
()
tests/edge_cases.rs:194
Functiontest_very_large_payload
()
tests/edge_cases.rs:78
Functiontest_zero_ttl
()
tests/edge_cases.rs:213
Methodtry_consume
Attempts to consume a message from the specified queue. Cost is O(NlogN) with N being the number of pending messages to be consumed # Arguments `queu
src/brokers/surrealdb/broker.rs:187
Methodtry_consume
( &self, queue_name: &str, options: Option<ConsumeOptions>, )
src/brokers/rabbitmq/broker.rs:168
Methodtry_consume_batch
Attempts to consume up to a number of messages from the specified queue. Does not block if not enough messages are available, and returns immediately.
src/queue.rs:624
Methodtry_consume_batch
Attempts to consume up to a number of messages from the specified queue. Does not block if not enough messages are available, and returns immediately.
src/brokers/broker.rs:60
Methodtry_consume_batch
Attempts to consume up to a number of messages from the specified queue. Does not block if not enough messages are available, and returns immmediately
src/brokers/surrealdb/broker.rs:217
Methodwith_attempts
(mut self, attempts: u8)
src/queue.rs:61
← previous101–186 of 186, ranked by callers