1use std::{
2 error::Error,
3 fmt,
4 path::{Path, PathBuf},
5 time::{SystemTime, UNIX_EPOCH},
6};
7
8use serde::Serialize;
9use tokio::fs;
10use tokio::time::{Duration, Instant, sleep};
11
12use crate::{
13 api::{ApiMetadata, RequestContext},
14 domain::{
15 KnowledgeMap, KnowledgeMapChange, KnowledgeMapRoute, KnowledgeMapSource,
16 KnowledgeMapSourceKind,
17 },
18 project::{AGENT_CONTRACT_DIR_NAME, KNOWLEDGE_MAP_FILE_NAME, KNOWLEDGE_MAP_RELATIVE_PATH},
19};
20
21#[derive(Debug, Clone, PartialEq, Eq)]
23pub struct KnowledgeMapSourceAddRequest {
24 pub id: String,
25 pub topic: String,
26 pub kind: KnowledgeMapSourceKind,
27 pub uri: String,
28 pub source_scope: Option<String>,
29 pub description: Option<String>,
30}
31
32#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
34pub struct KnowledgeMapMutationResponse {
35 pub metadata: ApiMetadata,
36 pub path: String,
37 pub map_version: u64,
38 pub summary: String,
39}
40
41#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
43pub struct KnowledgeMapShowResponse {
44 pub metadata: ApiMetadata,
45 pub path: String,
46 pub map: KnowledgeMap,
47}
48
49#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
51pub struct KnowledgeMapRouteResponse {
52 pub metadata: ApiMetadata,
53 pub path: String,
54 pub topic: String,
55 pub route: Option<KnowledgeMapRoute>,
56 pub sources: Vec<KnowledgeMapSource>,
57}
58
59#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
61pub struct KnowledgeMapValidationResponse {
62 pub metadata: ApiMetadata,
63 pub path: String,
64 pub valid: bool,
65 pub diagnostics: Vec<String>,
66}
67
68#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
70pub struct KnowledgeMapAgentSnippetResponse {
71 pub metadata: ApiMetadata,
72 pub snippet: String,
73}
74
75pub struct KnowledgeMapService {
77 repository_root: PathBuf,
78}
79
80impl KnowledgeMapService {
81 pub fn new(repository_root: PathBuf) -> Self {
82 Self { repository_root }
83 }
84
85 pub async fn init(
86 &self,
87 context: &RequestContext,
88 ) -> Result<KnowledgeMapMutationResponse, KnowledgeMapServiceError> {
89 let _lock = self.acquire_write_lock().await?;
90 let path = self.map_path();
91 if fs::try_exists(&path).await? {
92 let map = self.load_map().await?;
93 map.validate()?;
94 return Ok(self.mutation_response(
95 context,
96 map.map_version,
97 "knowledge map already exists".to_owned(),
98 ));
99 }
100
101 let map = KnowledgeMap::initial(now_stamp());
102 self.write_map(&map).await?;
103 Ok(self.mutation_response(context, map.map_version, "created knowledge map".to_owned()))
104 }
105
106 pub async fn show(
107 &self,
108 context: &RequestContext,
109 topic: Option<String>,
110 ) -> Result<KnowledgeMapShowResponse, KnowledgeMapServiceError> {
111 let mut map = self.load_map().await?;
112 if let Some(topic) = topic {
113 map.sources.retain(|source| source.topic == topic);
114 map.routes.retain(|route| route.topic == topic);
115 map.topics.retain(|entry| entry.id == topic);
116 }
117 Ok(KnowledgeMapShowResponse {
118 metadata: metadata(context),
119 path: KNOWLEDGE_MAP_RELATIVE_PATH.to_owned(),
120 map,
121 })
122 }
123
124 pub async fn route(
125 &self,
126 context: &RequestContext,
127 topic: String,
128 ) -> Result<KnowledgeMapRouteResponse, KnowledgeMapServiceError> {
129 let map = self.load_map().await?;
130 let route = map
131 .routes
132 .iter()
133 .find(|route| route.topic == topic)
134 .cloned();
135 let source_order = route
136 .as_ref()
137 .map(|route| route.source_order.as_slice())
138 .unwrap_or(&[]);
139 let sources = source_order
140 .iter()
141 .filter_map(|id| map.sources.iter().find(|source| &source.id == id).cloned())
142 .collect();
143
144 Ok(KnowledgeMapRouteResponse {
145 metadata: metadata(context),
146 path: KNOWLEDGE_MAP_RELATIVE_PATH.to_owned(),
147 topic,
148 route,
149 sources,
150 })
151 }
152
153 pub async fn add_source(
154 &self,
155 context: &RequestContext,
156 request: KnowledgeMapSourceAddRequest,
157 ) -> Result<KnowledgeMapMutationResponse, KnowledgeMapServiceError> {
158 let _lock = self.acquire_write_lock().await?;
159 let mut map = self.load_or_initial().await?;
160 let id = request.id.clone();
161 let topic = request.topic.clone();
162 let source = KnowledgeMapSource::new(
163 request.id,
164 request.topic,
165 request.kind,
166 request.uri,
167 request.source_scope,
168 request.description,
169 )?;
170 map.add_source(source)?;
171 map.record_change(
172 "source.add",
173 format!("Added source '{id}' to topic '{topic}'."),
174 now_stamp(),
175 );
176 self.write_map(&map).await?;
177 Ok(self.mutation_response(context, map.map_version, format!("added source {id}")))
178 }
179
180 pub async fn update_source(
181 &self,
182 context: &RequestContext,
183 change: KnowledgeMapChange,
184 ) -> Result<KnowledgeMapMutationResponse, KnowledgeMapServiceError> {
185 let _lock = self.acquire_write_lock().await?;
186 let mut map = self.load_map().await?;
187 let id = change.id.clone();
188 map.update_source(change)?;
189 map.record_change(
190 "source.update",
191 format!("Updated source '{id}'."),
192 now_stamp(),
193 );
194 self.write_map(&map).await?;
195 Ok(self.mutation_response(context, map.map_version, format!("updated source {id}")))
196 }
197
198 pub async fn remove_source(
199 &self,
200 context: &RequestContext,
201 id: String,
202 ) -> Result<KnowledgeMapMutationResponse, KnowledgeMapServiceError> {
203 let _lock = self.acquire_write_lock().await?;
204 let mut map = self.load_map().await?;
205 map.remove_source(&id)?;
206 map.record_change(
207 "source.remove",
208 format!("Removed source '{id}'."),
209 now_stamp(),
210 );
211 self.write_map(&map).await?;
212 Ok(self.mutation_response(context, map.map_version, format!("removed source {id}")))
213 }
214
215 pub async fn validate(
216 &self,
217 context: &RequestContext,
218 ) -> Result<KnowledgeMapValidationResponse, KnowledgeMapServiceError> {
219 let mut diagnostics = Vec::new();
220 match self.load_map().await {
221 Ok(map) => {
222 if let Err(error) = map.validate() {
223 diagnostics.push(error.to_string());
224 }
225 }
226 Err(error) => diagnostics.push(error.to_string()),
227 }
228
229 let agents_path = self.repository_root.join("AGENTS.md");
230 match fs::read_to_string(&agents_path).await {
231 Ok(contents) if contents.contains(KNOWLEDGE_MAP_RELATIVE_PATH) => {}
232 Ok(_) => diagnostics.push(format!(
233 "AGENTS.md does not reference {KNOWLEDGE_MAP_RELATIVE_PATH}"
234 )),
235 Err(error) => diagnostics.push(format!("failed to read AGENTS.md: {error}")),
236 }
237
238 Ok(KnowledgeMapValidationResponse {
239 metadata: metadata(context),
240 path: KNOWLEDGE_MAP_RELATIVE_PATH.to_owned(),
241 valid: diagnostics.is_empty(),
242 diagnostics,
243 })
244 }
245
246 pub fn agent_snippet(&self, context: &RequestContext) -> KnowledgeMapAgentSnippetResponse {
247 KnowledgeMapAgentSnippetResponse {
248 metadata: metadata(context),
249 snippet: format!("Knowledge map: {KNOWLEDGE_MAP_RELATIVE_PATH}"),
250 }
251 }
252
253 async fn load_or_initial(&self) -> Result<KnowledgeMap, KnowledgeMapServiceError> {
254 let path = self.map_path();
255 if fs::try_exists(&path).await? {
256 self.load_map().await
257 } else {
258 Ok(KnowledgeMap::initial(now_stamp()))
259 }
260 }
261
262 async fn load_map(&self) -> Result<KnowledgeMap, KnowledgeMapServiceError> {
263 let content = fs::read_to_string(self.map_path()).await?;
264 let map = serde_norway::from_str::<KnowledgeMap>(&content)
265 .map_err(|error| KnowledgeMapServiceError::Yaml(error.to_string()))?;
266 map.validate()?;
267 Ok(map)
268 }
269
270 async fn write_map(&self, map: &KnowledgeMap) -> Result<(), KnowledgeMapServiceError> {
271 map.validate()?;
272 let dir = self.repository_root.join(AGENT_CONTRACT_DIR_NAME);
273 fs::create_dir_all(&dir).await?;
274 let path = self.map_path();
275 let temp_path = dir.join(format!(
276 "{KNOWLEDGE_MAP_FILE_NAME}.{}.{}.tmp",
277 std::process::id(),
278 SystemTime::now()
279 .duration_since(UNIX_EPOCH)
280 .map(|duration| duration.as_nanos())
281 .unwrap_or(0)
282 ));
283 let yaml = serde_norway::to_string(map)
284 .map_err(|error| KnowledgeMapServiceError::Yaml(error.to_string()))?;
285 fs::write(&temp_path, yaml).await?;
286 if cfg!(target_os = "windows") && fs::try_exists(&path).await? {
287 fs::remove_file(&path).await?;
288 }
289 fs::rename(&temp_path, &path).await?;
290 Ok(())
291 }
292
293 async fn acquire_write_lock(&self) -> Result<KnowledgeMapWriteLock, KnowledgeMapServiceError> {
294 let dir = self.repository_root.join(AGENT_CONTRACT_DIR_NAME);
295 fs::create_dir_all(&dir).await?;
296 let path = dir.join(format!("{KNOWLEDGE_MAP_FILE_NAME}.lock"));
297 let deadline = Instant::now() + Duration::from_secs(10);
298 loop {
299 match fs::OpenOptions::new()
300 .write(true)
301 .create_new(true)
302 .open(&path)
303 .await
304 {
305 Ok(_) => return Ok(KnowledgeMapWriteLock { path }),
306 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
307 if Instant::now() >= deadline {
308 return Err(KnowledgeMapServiceError::LockTimeout(path));
309 }
310 sleep(Duration::from_millis(25)).await;
311 }
312 Err(error) => return Err(KnowledgeMapServiceError::Io(error)),
313 }
314 }
315 }
316
317 fn mutation_response(
318 &self,
319 context: &RequestContext,
320 map_version: u64,
321 summary: String,
322 ) -> KnowledgeMapMutationResponse {
323 KnowledgeMapMutationResponse {
324 metadata: metadata(context),
325 path: KNOWLEDGE_MAP_RELATIVE_PATH.to_owned(),
326 map_version,
327 summary,
328 }
329 }
330
331 fn map_path(&self) -> PathBuf {
332 self.repository_root
333 .join(Path::new(AGENT_CONTRACT_DIR_NAME))
334 .join(KNOWLEDGE_MAP_FILE_NAME)
335 }
336}
337
338#[derive(Debug)]
340pub enum KnowledgeMapServiceError {
341 Io(std::io::Error),
342 Yaml(String),
343 Domain(crate::domain::DomainError),
344 LockTimeout(PathBuf),
345}
346
347impl fmt::Display for KnowledgeMapServiceError {
348 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
349 match self {
350 Self::Io(error) => write!(formatter, "{error}"),
351 Self::Yaml(error) => write!(formatter, "invalid knowledge map YAML: {error}"),
352 Self::Domain(error) => write!(formatter, "{error}"),
353 Self::LockTimeout(path) => write!(
354 formatter,
355 "timed out waiting for knowledge map write lock '{}'",
356 path.display()
357 ),
358 }
359 }
360}
361
362impl Error for KnowledgeMapServiceError {}
363
364impl From<std::io::Error> for KnowledgeMapServiceError {
365 fn from(error: std::io::Error) -> Self {
366 Self::Io(error)
367 }
368}
369
370impl From<crate::domain::DomainError> for KnowledgeMapServiceError {
371 fn from(error: crate::domain::DomainError) -> Self {
372 Self::Domain(error)
373 }
374}
375
376fn metadata(context: &RequestContext) -> ApiMetadata {
377 ApiMetadata::graph_only(context, crate::domain::GraphVersion::ZERO)
378}
379
380struct KnowledgeMapWriteLock {
381 path: PathBuf,
382}
383
384impl Drop for KnowledgeMapWriteLock {
385 fn drop(&mut self) {
386 let _ = std::fs::remove_file(&self.path);
387 }
388}
389
390fn now_stamp() -> String {
391 let seconds = SystemTime::now()
392 .duration_since(UNIX_EPOCH)
393 .map(|duration| duration.as_secs())
394 .unwrap_or(0);
395 format!("unix:{seconds}")
396}
397
398#[cfg(test)]
399#[path = "mod_tests.rs"]
400mod tests;