Skip to main content

relay_knowledge/application/knowledge/
map.rs

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/// Request to register a source in the repository knowledge map.
22#[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/// Response shared by map mutation commands.
33#[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/// Response returned by read-only map commands.
42#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
43pub struct KnowledgeMapShowResponse {
44    pub metadata: ApiMetadata,
45    pub path: String,
46    pub map: KnowledgeMap,
47}
48
49/// Response returned by topic routing commands.
50#[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/// Response returned by validation commands.
60#[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/// Response that contains the AGENTS.md reference snippet.
69#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
70pub struct KnowledgeMapAgentSnippetResponse {
71    pub metadata: ApiMetadata,
72    pub snippet: String,
73}
74
75/// File-backed service for the shared YAML knowledge navigation contract.
76pub 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/// Error surfaced by the file-backed knowledge map service.
339#[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)]
399mod tests {
400    use super::*;
401
402    #[tokio::test]
403    async fn writes_and_reads_yaml_contract() {
404        let root = std::env::temp_dir().join(format!(
405            "relay-knowledge-map-{}",
406            SystemTime::now()
407                .duration_since(UNIX_EPOCH)
408                .expect("time should work")
409                .as_nanos()
410        ));
411        fs::create_dir_all(&root).await.expect("root should create");
412        fs::write(
413            root.join("AGENTS.md"),
414            format!("Knowledge map: {KNOWLEDGE_MAP_RELATIVE_PATH}"),
415        )
416        .await
417        .expect("agents should write");
418        let service = KnowledgeMapService::new(root.clone());
419        let context = RequestContext::for_interface(crate::api::InterfaceKind::Cli);
420
421        service.init(&context).await.expect("init should work");
422        service
423            .add_source(
424                &context,
425                KnowledgeMapSourceAddRequest {
426                    id: "build-cargo".to_owned(),
427                    topic: "build".to_owned(),
428                    kind: KnowledgeMapSourceKind::Config,
429                    uri: "Cargo.toml".to_owned(),
430                    source_scope: Some("repo".to_owned()),
431                    description: None,
432                },
433            )
434            .await
435            .expect("source should add");
436        service
437            .update_source(
438                &context,
439                crate::domain::KnowledgeMapChange {
440                    id: "build-cargo".to_owned(),
441                    topic: None,
442                    kind: None,
443                    uri: None,
444                    source_scope: None,
445                    description: Some("Cargo package manifest".to_owned()),
446                },
447            )
448            .await
449            .expect("existing map should be replaceable");
450        let route = service
451            .route(&context, "build".to_owned())
452            .await
453            .expect("route should load");
454        let validation = service
455            .validate(&context)
456            .await
457            .expect("validate should run");
458
459        assert_eq!(route.sources[0].id, "build-cargo");
460        assert!(validation.valid);
461        let _ = fs::remove_dir_all(root).await;
462    }
463
464    #[tokio::test]
465    async fn concurrent_source_adds_preserve_both_changes() {
466        let root = std::env::temp_dir().join(format!(
467            "relay-knowledge-map-concurrent-{}",
468            SystemTime::now()
469                .duration_since(UNIX_EPOCH)
470                .expect("time should work")
471                .as_nanos()
472        ));
473        fs::create_dir_all(&root).await.expect("root should create");
474        let service = KnowledgeMapService::new(root.clone());
475        let context = RequestContext::for_interface(crate::api::InterfaceKind::Cli);
476        service.init(&context).await.expect("init should work");
477
478        let first = service.add_source(
479            &context,
480            KnowledgeMapSourceAddRequest {
481                id: "build-cargo".to_owned(),
482                topic: "build".to_owned(),
483                kind: KnowledgeMapSourceKind::Config,
484                uri: "Cargo.toml".to_owned(),
485                source_scope: Some("repo".to_owned()),
486                description: None,
487            },
488        );
489        let second = service.add_source(
490            &context,
491            KnowledgeMapSourceAddRequest {
492                id: "build-readme".to_owned(),
493                topic: "build".to_owned(),
494                kind: KnowledgeMapSourceKind::Doc,
495                uri: "README.md".to_owned(),
496                source_scope: Some("repo".to_owned()),
497                description: None,
498            },
499        );
500
501        let (first, second) = tokio::join!(first, second);
502        first.expect("first add should succeed");
503        second.expect("second add should succeed");
504        let route = service
505            .route(&context, "build".to_owned())
506            .await
507            .expect("route should load");
508
509        assert_eq!(route.sources.len(), 2);
510        let _ = fs::remove_dir_all(root).await;
511    }
512}