- 1
use std::collections::VecDeque; - 2
use std::sync::atomic::{AtomicBool, Ordering}; - 3
use std::sync::{Mutex, PoisonError}; - 4
use tokio::sync::Notify; - 5
- 6
use vak_llm::Message; - 7
- 8
#[derive(Debug, Clone, Copy, PartialEq, Eq)] - 9
pub enum DrainMode { - 10
OneAtATime, - 11
All, - 12
} - 13
- 14
/// Cap on each queue's length. Without one, a sender who fires messages - 15
/// (each potentially carrying base64 image blocks) faster than a - 16
/// long-running turn drains them grows memory without bound. Once full, the - 17
/// oldest entry is dropped to make room for the newest — for steering in - 18
/// particular, the most recently expressed intent is the one that should - 19
/// win when the sender is over-sending, not the stalest queued message. - 20
const MAX_QUEUE_LEN: usize = 200; - 21
- 22
#[derive(Debug, Default)] - 23
struct Queues { - 24
steering: VecDeque<Message>, - 25
follow_up: VecDeque<Message>, - 26
outcome_updates: VecDeque<vak_session::types::IntentRecord>, - 27
} - 28
- 29
/// Two queues with two polling sites: steering interrupts the current run, - 30
/// follow-up continues after a natural stop. Interior-mutable so the UI can - 31
/// push while an agent run holds only a shared reference. Entries carry full - 32
/// user messages (text + image blocks) so queued input is never degraded — - 33
/// model-visible content must match what the sender supplied (invariant 1). - 34
#[derive(Debug, Default)] - 35
pub struct SteeringQueues { - 36
q: Mutex<Queues>, - 37
paused: AtomicBool, - 38
resumed: Notify, - 39
} - 40
- 41
impl SteeringQueues { - 42
pub fn new() -> Self { - 43
Self::default() - 44
} - 45
- 46
/// Pause at the next agent-turn boundary. Queued input remains durable. - 47
pub fn pause(&self) { - 48
self.paused.store(true, Ordering::Release); - 49
} - 50
- 51
pub fn resume(&self) { - 52
self.paused.store(false, Ordering::Release); - 53
self.resumed.notify_waiters(); - 54
} - 55
- 56
pub fn is_paused(&self) -> bool { - 57
self.paused.load(Ordering::Acquire) - 58
} - 59
- 60
pub async fn wait_if_paused(&self, cancel: &tokio_util::sync::CancellationToken) -> bool { - 61
while self.is_paused() { - 62
// Register the notification before rechecking the flag. A resume - 63
// between the flag check and waiter registration otherwise loses - 64
// the wakeup and can strand the run indefinitely. - 65
let resumed = self.resumed.notified(); - 66
if !self.is_paused() { - 67
break; - 68
} - 69
tokio::select! { - 70
_ = resumed => {} - 71
_ = cancel.cancelled() => return false, - 72
} - 73
} - 74
true - 75
} - 76
- 77
pub fn push_steering(&self, text: impl Into<String>) { - 78
self.push_steering_message(Message::user_text(text)); - 79
} - 80
- 81
pub fn push_steering_message(&self, message: Message) { - 82
let mut q = self.lock(); - 83
q.steering.push_back(message); - 84
while q.steering.len() > MAX_QUEUE_LEN { - 85
q.steering.pop_front(); - 86
} - 87
} - 88
- 89
pub fn push_follow_up(&self, text: impl Into<String>) { - 90
let mut q = self.lock(); - 91
q.follow_up.push_back(Message::user_text(text)); - 92
while q.follow_up.len() > MAX_QUEUE_LEN { - 93
q.follow_up.pop_front(); - 94
} - 95
} - 96
- 97
pub fn push_outcome_update(&self, update: vak_session::types::IntentRecord) { - 98
let mut q = self.lock(); - 99
q.outcome_updates.push_back(update); - 100
while q.outcome_updates.len() > MAX_QUEUE_LEN { - 101
q.outcome_updates.pop_front(); - 102
} - 103
} - 104
- 105
pub fn take_outcome_update(&self) -> Option<vak_session::types::IntentRecord> { - 106
self.lock().outcome_updates.pop_front() - 107
} - 108
- 109
pub fn drain(&self, mode: DrainMode) -> Vec<Message> { - 110
let take = match mode { - 111
DrainMode::OneAtATime => 1, - 112
DrainMode::All => usize::MAX, - 113
}; - 114
let mut q = self.lock(); - 115
let mut out = Vec::new(); - 116
while out.len() < take { - 117
match q.steering.pop_front() { - 118
Some(m) => out.push(m), - 119
None => break, - 120
} - 121
} - 122
out - 123
} - 124
- 125
pub fn take_follow_up(&self) -> Option<Message> { - 126
self.lock().follow_up.pop_front() - 127
} - 128
- 129
pub fn has_pending(&self) -> bool { - 130
let q = self.lock(); - 131
!q.steering.is_empty() || !q.follow_up.is_empty() || !q.outcome_updates.is_empty() - 132
} - 133
- 134
/// Merge drained messages into one prompt turn. Same-role block - 135
/// concatenation keeps text AND image blocks; the merged message is - 136
/// exactly what gets logged and dispatched. - 137
pub fn merge_prompt(messages: Vec<Message>) -> Option<Message> { - 138
if messages.is_empty() { - 139
return None; - 140
} - 141
if messages.len() == 1 { - 142
return messages.into_iter().next(); - 143
} - 144
let mut content = Vec::new(); - 145
for m in messages { - 146
content.extend(m.content); - 147
} - 148
Some(Message { - 149
role: vak_llm::Role::User, - 150
content, - 151
}) - 152
} - 153
- 154
fn lock(&self) -> std::sync::MutexGuard<'_, Queues> { - 155
self.q.lock().unwrap_or_else(PoisonError::into_inner) - 156
} - 157
} - 158
- 159
#[cfg(test)] - 160
mod tests { - 161
#![allow(clippy::unwrap_used, clippy::expect_used)] - 162
use super::*; - 163
use vak_llm::{ContentBlock, Role}; - 164
- 165
#[test] - 166
fn steering_round_trips_images_through_merge() { - 167
let q = SteeringQueues::new(); - 168
q.push_steering("plain"); - 169
let with_image = Message { - 170
role: Role::User, - 171
content: vec![ - 172
ContentBlock::text("look"), - 173
ContentBlock::image_base64("image/png", "QUJD"), - 174
], - 175
}; - 176
q.push_steering_message(with_image); - 177
- 178
let drained = q.drain(DrainMode::All); - 179
assert_eq!(drained.len(), 2); - 180
let merged = SteeringQueues::merge_prompt(drained).unwrap(); - 181
assert_eq!(merged.role, Role::User); - 182
assert_eq!(merged.content.len(), 3, "text + image + text survive"); - 183
} - 184
- 185
#[test] - 186
fn empty_drain_merges_to_none() { - 187
assert!(SteeringQueues::merge_prompt(Vec::new()).is_none()); - 188
} - 189
- 190
#[test] - 191
fn steering_queue_drops_oldest_past_cap() { - 192
let q = SteeringQueues::new(); - 193
for i in 0..MAX_QUEUE_LEN + 10 { - 194
q.push_steering(format!("msg {i}")); - 195
} - 196
let drained = q.drain(DrainMode::All); - 197
assert_eq!(drained.len(), MAX_QUEUE_LEN); - 198
// Oldest entries (0..10) were evicted; the newest survive in order. - 199
let vak_llm::ContentBlock::Text { text } = &drained[0].content[0] else { - 200
unreachable!("expected text block"); - 201
}; - 202
assert_eq!(text, "msg 10"); - 203
} - 204
} - 205
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.