MCPcopy Create free account
hub / github.com/GCWing/BitFun / enqueue

Method enqueue

src/crates/execution/agent-runtime/src/event_queue.rs:76–120  ·  view source on GitHub ↗

Enqueue event

(
        &self,
        event: AgenticEvent,
        priority: Option<EventPriority>,
    )

Source from the content-addressed store, hash-verified

74
75 /// Enqueue event
76 pub async fn enqueue(
77 &self,
78 event: AgenticEvent,
79 priority: Option<EventPriority>,
80 ) -> EventBusResult<String> {
81 let priority = priority.unwrap_or_else(|| event.default_priority());
82 let envelope = EventEnvelope::new(event, priority);
83 let event_id = envelope.id.clone();
84
85 // Check queue size
86 {
87 let queue = self.queue.lock().await;
88 if queue.len() >= self.config.max_queue_size {
89 warn!("Event queue full, dropping event: event_id={}", event_id);
90 return Ok(event_id);
91 }
92 }
93
94 // Add to queue
95 {
96 let mut queue = self.queue.lock().await;
97 queue.push(std::cmp::Reverse(envelope.clone()));
98 }
99
100 let _ = self.broadcast_tx.send(envelope);
101
102 // Update statistics: get queue size first, then update statistics (avoid getting queue lock while holding stats lock)
103 let queue_len = self.queue.lock().await.len();
104 {
105 let mut stats = self.stats.lock().await;
106 stats.total_enqueued += 1;
107 stats.pending_events = queue_len;
108 }
109
110 // Notify waiting consumers
111 self.notify.notify_one();
112
113 trace!(
114 "Event enqueued: event_id={}, priority={:?}",
115 event_id,
116 priority
117 );
118
119 Ok(event_id)
120 }
121
122 /// Dequeue batch of events
123 pub async fn dequeue_batch(&self, max_size: usize) -> Vec<EventEnvelope> {

Callers

nothing calls this directly

Calls 5

default_priorityMethod · 0.80
cloneMethod · 0.45
lenMethod · 0.45
pushMethod · 0.45
sendMethod · 0.45

Tested by

no test coverage detected