- 1
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 2
- 3
use std::collections::VecDeque; - 4
use std::sync::atomic::{AtomicUsize, Ordering}; - 5
use std::sync::{Arc, Mutex}; - 6
- 7
use tokio_util::sync::CancellationToken; - 8
- 9
use vak_core::Core; - 10
use vak_llm::stream; - 11
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, Usage}; - 12
use vak_llm::{EventStream, LlmError, Provider}; - 13
use vak_server::surfaces::telegram::TelegramBridge; - 14
- 15
struct Scripted { - 16
responses: Mutex<VecDeque<AssistantMessage>>, - 17
} - 18
- 19
#[async_trait::async_trait] - 20
impl Provider for Scripted { - 21
fn name(&self) -> &str { - 22
"scripted" - 23
} - 24
- 25
async fn stream( - 26
&self, - 27
_request: ChatRequest, - 28
_cancel: CancellationToken, - 29
) -> Result<EventStream, LlmError> { - 30
let next = self.responses.lock().unwrap().pop_front(); - 31
let (mut sink, rx) = stream::channel(64); - 32
match next { - 33
Some(m) => { - 34
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 35
sink.close_message(m).await; - 36
} - 37
None => sink.close_error(LlmError::Parse("exhausted".into())).await, - 38
} - 39
Ok(rx) - 40
} - 41
} - 42
- 43
fn text(t: &str) -> AssistantMessage { - 44
AssistantMessage { - 45
content: vec![ContentBlock::text(t)], - 46
stop_reason: vak_llm::types::StopReason::EndTurn, - 47
usage: Usage { - 48
input_tokens: 7, - 49
output_tokens: 3, - 50
..Default::default() - 51
}, - 52
model: "test-model".into(), - 53
response_id: None, - 54
} - 55
} - 56
- 57
/// Minimal Bot API double: serves one scripted update, then empty polls; - 58
/// records every sendMessage payload. - 59
async fn spawn_mock_telegram() -> (String, Arc<Mutex<Vec<serde_json::Value>>>, Arc<AtomicUsize>) { - 60
let updates_left = Arc::new(Mutex::new(1u32)); - 61
let sent: Arc<Mutex<Vec<serde_json::Value>>> = Arc::new(Mutex::new(Vec::new())); - 62
let get_calls = Arc::new(AtomicUsize::new(0)); - 63
- 64
let ul = updates_left.clone(); - 65
let gc = get_calls.clone(); - 66
async fn ok_json() -> serde_json::Value { - 67
serde_json::json!({ "ok": true }) - 68
} - 69
let app = axum::Router::new() - 70
.route( - 71
"/botbottok/getUpdates", - 72
axum::routing::get(move || async move { - 73
gc.fetch_add(1, Ordering::SeqCst); - 74
let mut left = ul.lock().unwrap(); - 75
if *left > 0 { - 76
*left -= 1; - 77
// First scripted update is a PHOTO message. - 78
return axum::Json(serde_json::json!({ - 79
"ok": true, - 80
"result": [{ - 81
"update_id": 777, - 82
"message": { - 83
"chat": {"id": 4242}, - 84
"photo": [ - 85
{"file_id": "small", "width": 90, "height": 90}, - 86
{"file_id": "large", "width": 640, "height": 640} - 87
] - 88
} - 89
}] - 90
})); - 91
} - 92
axum::Json(serde_json::json!({ "ok": true, "result": [] })) - 93
}), - 94
) - 95
.route( - 96
"/botbottok/getFile", - 97
axum::routing::get(|| async { - 98
axum::Json(serde_json::json!({ - 99
"ok": true, - 100
"result": {"file_path": "photos/img.jpg"} - 101
})) - 102
}), - 103
) - 104
.route( - 105
"/file/botbottok/photos/img.jpg", - 106
axum::routing::get(|| async { [0x89u8, b'P', b'N', b'G'].to_vec() }), - 107
) - 108
.route( - 109
"/botbottok/sendMessage", - 110
axum::routing::post({ - 111
let sent = sent.clone(); - 112
move |axum::Json(body): axum::Json<serde_json::Value>| async move { - 113
sent.lock().unwrap().push(body); - 114
axum::Json(ok_json().await) - 115
} - 116
}), - 117
); - 118
- 119
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 120
let addr = listener.local_addr().unwrap(); - 121
tokio::spawn(async move { - 122
axum::serve(listener, app).await.unwrap(); - 123
}); - 124
(format!("http://{addr}"), sent, get_calls) - 125
} - 126
- 127
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 128
async fn telegram_bridge_routes_message_and_delivers_reply() { - 129
let (tg_base, sent, get_calls) = spawn_mock_telegram().await; - 130
- 131
// The gateway side of the contract. - 132
let dir = tempfile::tempdir().unwrap(); - 133
let cwd = dir.path().to_path_buf(); - 134
// Hermetic against the developer's global config (e.g. reflection=true): - 135
// pin learning flags off for deterministic scripted flows. - 136
let _ = std::fs::create_dir_all(cwd.join(".vak")); - 137
let _ = std::fs::write( - 138
cwd.join(".vak/config.toml"), - 139
"[memory]\nreflection = false\n\n[gateway]\nchat_allowlist = [\"telegram:4242\"]\n", - 140
); - 141
vak_config::paths::isolate_home_for_tests(); - 142
let core = Core::new_with_trust(cwd.clone(), true).unwrap(); - 143
core.set_sessions_home(dir.path().join("home")); - 144
core.set_provider_instance(Arc::new(Scripted { - 145
responses: Mutex::new(VecDeque::from(vec![text("pong from agent")])), - 146
})); - 147
std::mem::forget(dir); - 148
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 149
let gw_addr = listener.local_addr().unwrap(); - 150
tokio::spawn(async move { - 151
axum::serve(listener, vak_server::gateway_router(core)) - 152
.await - 153
.unwrap(); - 154
}); - 155
- 156
let bridge = TelegramBridge { - 157
token_env: String::new(), - 158
locks_dir: None, - 159
api_base: tg_base, - 160
bot_token: "bottok".into(), - 161
gateway_url: format!("http://{gw_addr}"), - 162
gateway_token: "vk_test".into(), - 163
bot_id: None, - 164
}; - 165
- 166
let next = bridge.tick(0).await.unwrap(); - 167
assert_eq!(next, 778, "offset advances past the handled update"); - 168
- 169
{ - 170
let delivered = sent.lock().unwrap(); - 171
assert_eq!(delivered.len(), 1); - 172
assert_eq!(delivered[0]["chat_id"], 4242); - 173
assert_eq!(delivered[0]["text"], "pong from agent"); - 174
} - 175
- 176
// A quiet poll delivers nothing new. - 177
bridge.tick(next).await.unwrap(); - 178
assert_eq!(sent.lock().unwrap().len(), 1); - 179
assert!( - 180
get_calls.load(Ordering::SeqCst) >= 2, - 181
"long-poll must keep polling" - 182
); - 183
} - 184
- 185
// Regression (docs/design/31-network-resilience.md): the bridge must ride - 186
// out an outage window — refused polls, network switch, sleep/wake — and - 187
// resume consumption gap-free. It used to exit after 10 consecutive - 188
// failures, killing the channel until a human restarted it. - 189
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 190
async fn bridge_survives_outage_window_and_resumes_cursor() { - 191
let down = Arc::new(std::sync::atomic::AtomicBool::new(true)); - 192
let served = Arc::new(std::sync::atomic::AtomicUsize::new(0)); - 193
let sent: Arc<Mutex<Vec<serde_json::Value>>> = Arc::new(Mutex::new(Vec::new())); - 194
let get_calls = Arc::new(std::sync::atomic::AtomicUsize::new(0)); - 195
- 196
let d = down.clone(); - 197
let s = served.clone(); - 198
let gc = get_calls.clone(); - 199
async fn ok_json() -> serde_json::Value { - 200
serde_json::json!({ "ok": true }) - 201
} - 202
async fn updates_handler( - 203
d: Arc<std::sync::atomic::AtomicBool>, - 204
s: Arc<std::sync::atomic::AtomicUsize>, - 205
gc: Arc<std::sync::atomic::AtomicUsize>, - 206
) -> axum::response::Response { - 207
use axum::response::IntoResponse; - 208
gc.fetch_add(1, Ordering::SeqCst); - 209
if d.load(Ordering::SeqCst) { - 210
return ( - 211
axum::http::StatusCode::SERVICE_UNAVAILABLE, - 212
axum::Json(serde_json::json!({ "ok": false })), - 213
) - 214
.into_response(); - 215
} - 216
if s.fetch_add(1, Ordering::SeqCst) == 0 { - 217
return axum::Json(serde_json::json!({ - 218
"ok": true, - 219
"result": [{ - 220
"update_id": 900, - 221
"message": { "chat": {"id": 1}, "text": "ping during recovery" } - 222
}] - 223
})) - 224
.into_response(); - 225
} - 226
axum::Json(serde_json::json!({ "ok": true, "result": [] })).into_response() - 227
} - 228
let app = axum::Router::new() - 229
.route( - 230
"/botbottok/getUpdates", - 231
axum::routing::get(move || updates_handler(d.clone(), s.clone(), gc.clone())), - 232
) - 233
.route( - 234
"/botbottok/sendMessage", - 235
axum::routing::post({ - 236
let sent = sent.clone(); - 237
move |axum::Json(body): axum::Json<serde_json::Value>| async move { - 238
sent.lock().unwrap().push(body); - 239
axum::Json(ok_json().await) - 240
} - 241
}), - 242
); - 243
- 244
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 245
let addr = listener.local_addr().unwrap(); - 246
tokio::spawn(async move { - 247
axum::serve(listener, app).await.unwrap(); - 248
}); - 249
- 250
// Gateway side accepts the inbound message. - 251
let dir = tempfile::tempdir().unwrap(); - 252
let cwd = dir.path().to_path_buf(); - 253
let _ = std::fs::create_dir_all(cwd.join(".vak")); - 254
let _ = std::fs::write( - 255
cwd.join(".vak/config.toml"), - 256
"[memory]\nreflection = false\n\n[gateway]\nchat_allowlist = [\"telegram:1\"]\n", - 257
); - 258
vak_config::paths::isolate_home_for_tests(); - 259
let core = Core::new_with_trust(cwd.clone(), true).unwrap(); - 260
core.set_sessions_home(dir.path().join("home")); - 261
core.set_provider_instance(Arc::new(Scripted { - 262
responses: Mutex::new(VecDeque::from(vec![text("pong after outage")])), - 263
})); - 264
std::mem::forget(dir); - 265
let listener2 = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); - 266
let gw_addr = listener2.local_addr().unwrap(); - 267
tokio::spawn(async move { - 268
axum::serve(listener2, vak_server::gateway_router(core)) - 269
.await - 270
.unwrap(); - 271
}); - 272
- 273
let bridge = TelegramBridge { - 274
token_env: String::new(), - 275
locks_dir: None, - 276
api_base: format!("http://{addr}"), - 277
bot_token: "bottok".into(), - 278
gateway_url: format!("http://{gw_addr}"), - 279
gateway_token: "vk_test".into(), - 280
bot_id: None, - 281
}; - 282
- 283
// Ownership probe runs while "down" — must classify as transient, not - 284
// conflict, and fall through to the loop. - 285
let handle = tokio::spawn(async move { bridge.run().await }); - 286
- 287
// Let the bridge hit several refused polls, then heal the network. - 288
tokio::time::sleep(std::time::Duration::from_millis(1500)).await; - 289
let refused = get_calls.load(Ordering::SeqCst); - 290
assert!(refused >= 1, "bridge should have attempted polls"); - 291
down.store(false, Ordering::SeqCst); - 292
- 293
// Recovery must deliver the queued update end-to-end. - 294
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(20); - 295
while sent.lock().unwrap().is_empty() { - 296
assert!( - 297
std::time::Instant::now() < deadline, - 298
"bridge never recovered after outage window" - 299
); - 300
tokio::time::sleep(std::time::Duration::from_millis(100)).await; - 301
} - 302
assert_eq!(sent.lock().unwrap()[0]["text"], "pong after outage"); - 303
assert!( - 304
!handle.is_finished(), - 305
"run() must not exit on transients (old give-up-after-10 behavior)" - 306
); - 307
} - 308
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.