- 1
#![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)] - 2
- 3
use std::collections::VecDeque; - 4
use std::sync::{Arc, Mutex}; - 5
- 6
use tokio_util::sync::CancellationToken; - 7
- 8
use vak_core::Core; - 9
use vak_llm as _; - 10
use vak_llm::stream; - 11
use vak_llm::types::{AssistantMessage, ChatRequest, ContentBlock, Usage}; - 12
use vak_llm::{EventStream, LlmError, Provider}; // keep types import path stable for future edits - 13
- 14
struct Scripted { - 15
responses: Mutex<VecDeque<AssistantMessage>>, - 16
} - 17
- 18
#[async_trait::async_trait] - 19
impl Provider for Scripted { - 20
fn name(&self) -> &str { - 21
"scripted" - 22
} - 23
- 24
async fn stream( - 25
&self, - 26
_request: ChatRequest, - 27
_cancel: CancellationToken, - 28
) -> Result<EventStream, LlmError> { - 29
let next = self.responses.lock().unwrap().pop_front(); - 30
let (mut sink, rx) = stream::channel(64); - 31
match next { - 32
Some(m) => { - 33
sink.push(stream::StreamEvent::Start { partial: m.clone() }); - 34
sink.close_message(m).await; - 35
} - 36
None => sink.close_error(LlmError::Parse("exhausted".into())).await, - 37
} - 38
Ok(rx) - 39
} - 40
} - 41
- 42
fn text(t: &str) -> AssistantMessage { - 43
AssistantMessage { - 44
content: vec![ContentBlock::text(t)], - 45
stop_reason: vak_llm::types::StopReason::EndTurn, - 46
usage: Usage { - 47
input_tokens: 7, - 48
output_tokens: 3, - 49
..Default::default() - 50
}, - 51
model: "test-model".into(), - 52
response_id: None, - 53
} - 54
} - 55
- 56
fn tool_call(id: &str, name: &str, input: serde_json::Value) -> AssistantMessage { - 57
AssistantMessage { - 58
content: vec![ContentBlock::ToolUse { - 59
id: id.into(), - 60
name: name.into(), - 61
input, - 62
}], - 63
stop_reason: vak_llm::types::StopReason::ToolUse, - 64
usage: Usage::default(), - 65
model: "test-model".into(), - 66
response_id: None, - 67
} - 68
} - 69
- 70
#[tokio::test(flavor = "multi_thread", worker_threads = 4)] - 71
async fn remember_propose_recall_promote_loop() { - 72
let dir = tempfile::tempdir().unwrap(); - 73
let home = dir.path().join("home"); - 74
let cwd = dir.path().to_path_buf(); - 75
- 76
vak_config::paths::isolate_home_for_tests(); - 77
let core = Core::new(cwd.clone()).unwrap(); - 78
core.set_sessions_home(home.clone()); - 79
core.set_permission_mode(vak_config::PermissionMode::WorkspaceWrite); - 80
core.set_provider_instance(Arc::new(Scripted { - 81
responses: Mutex::new(VecDeque::from(vec![ - 82
// Turn 1 of run A: the model journals a decision and drafts a skill. - 83
tool_call( - 84
"r1", - 85
"remember", - 86
serde_json::json!({ - 87
"note": "the deploy script must pause before rollback windows", - 88
"kind": "decision", - 89
"tag": "deploy-rollbacks" - 90
}), - 91
), - 92
tool_call( - 93
"p1", - 94
"propose_skill", - 95
serde_json::json!({ - 96
"name": "deploy-safely", - 97
"description": "Run deploys with rollback pauses", - 98
"instructions": "1. read scripts/deploy.sh\n2. pause before rollback windows" - 99
}), - 100
), - 101
text("Noted both."), - 102
])), - 103
})); - 104
- 105
// ---- Run A: remember + propose -------------------------------------- - 106
let s_a = core.start_session().await.unwrap(); - 107
let sid_a = s_a.header().unwrap().session_id.clone(); - 108
let (tx_a, _rx_a) = tokio::sync::mpsc::channel(256); - 109
let (outcome_a, log_a) = core - 110
.run_turn_with( - 111
s_a, - 112
"save what we learned", - 113
CancellationToken::new(), - 114
None, - 115
None, - 116
None, - 117
tx_a, - 118
) - 119
.await - 120
.unwrap(); - 121
assert!(matches!( - 122
outcome_a, - 123
vak_agent::TurnOutcome::Completed { .. } - 124
)); - 125
- 126
// Notes landed in this Agent's own memory, with provenance... - 127
let notes = vak_core::memory::list_notes(&core.sessions_home(), &cwd); - 128
assert_eq!(notes.len(), 1); - 129
assert_eq!(notes[0].kind, "decision"); - 130
assert_eq!(notes[0].tag, "deploy-rollbacks"); - 131
assert_eq!(notes[0].session_id, sid_a); - 132
- 133
// ...and the confirmations are logged on the ledger (invariant 1). - 134
let logged = log_a - 135
.message_chain() - 136
.iter() - 137
.map(|(_, m)| { - 138
m.content - 139
.iter() - 140
.filter_map(|b| match b { - 141
ContentBlock::ToolResult { content, .. } => Some(content.clone()), - 142
_ => None, - 143
}) - 144
.collect::<Vec<_>>() - 145
.join("\n") - 146
}) - 147
.collect::<Vec<_>>() - 148
.join("\n"); - 149
assert!(logged.contains("remembered"), "{logged}"); - 150
assert!(logged.contains("queued for review"), "{logged}"); - 151
- 152
// Proposal is pending, not installed. - 153
let pending = vak_core::learning::list_proposals(&home, &cwd); - 154
assert_eq!(pending.len(), 1); - 155
assert_eq!(pending[0].name, "deploy-safely"); - 156
assert!(!home.join("skills/deploy-safely/SKILL.md").exists()); - 157
- 158
// ---- Run B: recall ranks memory first -------------------------------- - 159
core.set_provider_instance(Arc::new(Scripted { - 160
responses: Mutex::new(VecDeque::from(vec![ - 161
tool_call( - 162
"s1", - 163
"session_search", - 164
serde_json::json!({"query": "deploy rollback"}), - 165
), - 166
text("Recalled from memory."), - 167
])), - 168
})); - 169
let s_b = core.start_session().await.unwrap(); - 170
let (_tx_b, _rx_b) = tokio::sync::mpsc::channel::<vak_agent::AgentEvent>(256); - 171
let (tx_b2, _rx_b2) = tokio::sync::mpsc::channel(256); - 172
let _ = (&_tx_b, &_rx_b); - 173
let (outcome_b, log_b) = core - 174
.run_turn_with( - 175
s_b, - 176
"what do we know about deploys?", - 177
CancellationToken::new(), - 178
None, - 179
None, - 180
None, - 181
tx_b2, - 182
) - 183
.await - 184
.unwrap(); - 185
assert!(matches!( - 186
outcome_b, - 187
vak_agent::TurnOutcome::Completed { .. } - 188
)); - 189
let search_output = log_b - 190
.message_chain() - 191
.iter() - 192
.flat_map(|(_, m)| m.content.iter()) - 193
.filter_map(|b| match b { - 194
ContentBlock::ToolResult { content, .. } => Some(content.clone()), - 195
_ => None, - 196
}) - 197
.find(|c| c.contains("hit(s)")) - 198
.expect("search result logged"); - 199
assert!( - 200
search_output.contains("deploy-rollbacks"), - 201
"memory hit present: {search_output}" - 202
); - 203
assert!( - 204
search_output.contains("memory"), - 205
"role surfaces as memory: {search_output}" - 206
); - 207
- 208
// ---- Promotion closes the loop --------------------------------------- - 209
let id = pending[0].id.clone(); - 210
let promoted = vak_core::learning::promote(&home, &cwd, &id).unwrap(); - 211
assert_eq!(promoted, "deploy-safely"); - 212
assert!(home.join("skills/deploy-safely/SKILL.md").exists()); - 213
assert!(vak_core::learning::list_proposals(&home, &cwd).is_empty()); - 214
let discovered = vak_core::skills::discover(&cwd, &home); - 215
assert!( - 216
discovered - 217
.iter() - 218
.any(|s| s.name == "deploy-safely" && s.description.contains("rollback")), - 219
"promoted skill enters discovery: {discovered:?}" - 220
); - 221
- 222
// Promoting twice or promoting after rejection fails cleanly. - 223
assert!(vak_core::learning::promote(&home, &cwd, &id).is_err()); - 224
} - 225
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.