use std::collections::BTreeMap;
use crate::{
Database, ExportedFact, FactId, HostError, IfMissing, RecallQuery, RememberInput, Workspace,
};
use super::{DbName, WorkspaceError};
pub const SELF_ENTITY: &str = "plugmem workspace self";
pub const ENTRY_TAG: &str = "plugmem-db";
pub const ARCHIVED_TAG: &str = "archived";
const NAME_KEY: &str = "name";
const OWNER_KEY: &str = "owner";
const OWNED_BY_REL: &str = "owned-by";
const ANCHOR_K: usize = 4;
#[derive(Clone, Copy, Debug, Default)]
pub struct Description<'a> {
pub text: &'a str,
pub tags: &'a [&'a str],
pub owner: Option<&'a str>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct DbEntry {
pub name: DbName,
pub description: String,
pub tags: Vec<String>,
pub owner: Option<String>,
}
impl DbEntry {
pub fn is_archived(&self) -> bool {
self.tags.iter().any(|t| t == ARCHIVED_TAG)
}
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct ReindexReport {
pub indexed: Vec<DbName>,
pub undescribed: Vec<DbName>,
pub busy: Vec<DbName>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
#[non_exhaustive]
pub enum WorkspaceIssue {
Missing {
name: DbName,
},
Undescribed {
name: DbName,
},
Stale {
name: DbName,
},
Unreadable {
name: DbName,
why: String,
},
AmbiguousSelf {
name: DbName,
facts: usize,
},
}
impl Workspace {
pub fn describe(
&self,
name: &DbName,
now_ms: u64,
desc: Description<'_>,
) -> Result<(), WorkspaceError> {
let db = self.get(name, now_ms, IfMissing::Create)?;
write_self(&db, now_ms, desc).map_err(|e| self.blame(name, e))?;
drop(db);
let registry = self.registry()?;
write_entry(®istry, now_ms, name, desc)?;
Ok(())
}
pub fn archive(&self, name: &DbName, now_ms: u64) -> Result<bool, WorkspaceError> {
let Some(entry) = self.entry(name)? else {
return Err(WorkspaceError::NoSuchDatabase {
name: name.clone(),
path: self.layout().path_of(name),
});
};
if entry.is_archived() {
return Ok(false);
}
let mut tags: Vec<&str> = entry.tags.iter().map(String::as_str).collect();
tags.push(ARCHIVED_TAG);
self.describe(
name,
now_ms,
Description {
text: &entry.description,
tags: &tags,
owner: entry.owner.as_deref(),
},
)?;
Ok(true)
}
pub fn entry(&self, name: &DbName) -> Result<Option<DbEntry>, WorkspaceError> {
Ok(self.entries()?.into_iter().find(|e| &e.name == name))
}
pub fn entries(&self) -> Result<Vec<DbEntry>, WorkspaceError> {
let registry = self.registry()?;
let mut out: Vec<DbEntry> = registry.export().iter().filter_map(entry_of).collect();
out.sort_by(|a, b| a.name.cmp(&b.name));
Ok(out)
}
pub fn find(&self, query: &str, k: usize, now_ms: u64) -> Result<Vec<DbEntry>, WorkspaceError> {
let registry = self.registry()?;
let hits = registry.recall(RecallQuery {
entities: &[query],
k,
..RecallQuery::text(now_ms, query)
})?;
let mut by_name: BTreeMap<String, DbEntry> = self
.entries()?
.into_iter()
.map(|e| (e.name.to_string(), e))
.collect();
let mut out = Vec::with_capacity(hits.facts.len());
for hit in &hits.facts {
if let Some(snap) = registry.get(hit.id)
&& let Some(name) = snap.metadata.get(NAME_KEY)
&& let Some(entry) = by_name.remove(name)
{
out.push(entry);
}
}
Ok(out)
}
pub fn reindex(&self, now_ms: u64) -> Result<ReindexReport, WorkspaceError> {
let mut report = ReindexReport::default();
for name in self.layout().list()? {
let db = match self.get(&name, now_ms, IfMissing::Fail) {
Ok(db) => db,
Err(WorkspaceError::Busy { .. }) => {
report.busy.push(name);
continue;
}
Err(e) => return Err(e),
};
let found = self_description(&db, now_ms).map_err(|e| self.blame(&name, e))?;
drop(db);
match found {
Some(desc) => {
let registry = self.registry()?;
let tags: Vec<&str> = desc.tags.iter().map(String::as_str).collect();
write_entry(
®istry,
now_ms,
&name,
Description {
text: &desc.text,
tags: &tags,
owner: desc.owner.as_deref(),
},
)?;
report.indexed.push(name);
}
None => report.undescribed.push(name),
}
}
Ok(report)
}
pub fn verify(&self, now_ms: u64) -> Result<Vec<WorkspaceIssue>, WorkspaceError> {
let on_disk = self.layout().list()?;
let recorded = self.entries()?;
let mut issues = Vec::new();
for entry in &recorded {
if !self.layout().exists(&entry.name) {
issues.push(WorkspaceIssue::Missing {
name: entry.name.clone(),
});
}
}
for name in on_disk {
let db = match self.get(&name, now_ms, IfMissing::Fail) {
Ok(db) => db,
Err(e) => {
issues.push(WorkspaceIssue::Unreadable {
name,
why: e.to_string(),
});
continue;
}
};
let anchored = anchored_facts(&db, now_ms).map_err(|e| self.blame(&name, e))?;
let own = anchored.first().map(|(id, snap)| SelfDescription {
text: snap.text.clone(),
tags: db.tags_of(*id),
owner: snap.metadata.get(OWNER_KEY).cloned(),
});
let anchored = anchored.len();
drop(db);
if anchored > 1 {
issues.push(WorkspaceIssue::AmbiguousSelf {
name: name.clone(),
facts: anchored,
});
}
match (own, recorded.iter().find(|e| e.name == name)) {
(None, None) => issues.push(WorkspaceIssue::Undescribed { name }),
(Some(own), Some(record)) if own.agrees_with(record) => {}
_ => issues.push(WorkspaceIssue::Stale { name }),
}
}
Ok(issues)
}
fn blame(&self, name: &DbName, e: HostError) -> WorkspaceError {
match e {
HostError::Locked { .. } => WorkspaceError::Busy { name: name.clone() },
other => WorkspaceError::Host(other),
}
}
}
fn write_self(db: &Database, now_ms: u64, desc: Description<'_>) -> Result<(), HostError> {
let owner = desc.owner.map(|o| [(OWNER_KEY, o)]);
let input = RememberInput {
entity: Some(SELF_ENTITY),
tags: desc.tags,
metadata: owner.as_ref().map(|m| m.as_slice()),
..RememberInput::text(now_ms, desc.text)
};
match anchored_facts(db, now_ms)?.first() {
Some((id, _)) => db.revise(*id, input)?,
None => db.remember(input)?,
};
Ok(())
}
fn write_entry(
registry: &Database,
now_ms: u64,
name: &DbName,
desc: Description<'_>,
) -> Result<(), WorkspaceError> {
let mut tags: Vec<&str> = Vec::with_capacity(desc.tags.len() + 1);
tags.push(ENTRY_TAG);
tags.extend(desc.tags.iter().copied().filter(|t| *t != ENTRY_TAG));
let mut metadata: Vec<(&str, &str)> = vec![(NAME_KEY, name.as_str())];
if let Some(owner) = desc.owner {
metadata.push((OWNER_KEY, owner));
}
let links: Vec<(&str, &str)> = desc
.owner
.map(|owner| vec![(OWNED_BY_REL, owner)])
.unwrap_or_default();
let input = RememberInput {
entity: Some(name.as_str()),
tags: &tags,
links: &links,
metadata: Some(&metadata),
..RememberInput::text(now_ms, desc.text)
};
match existing_record(registry, now_ms, name)? {
Some(id) => registry.revise(id, input)?,
None => registry.remember(input)?,
};
Ok(())
}
fn existing_record(
registry: &Database,
now_ms: u64,
name: &DbName,
) -> Result<Option<FactId>, HostError> {
let hits = registry.recall(RecallQuery {
entities: &[name.as_str()],
k: ANCHOR_K,
..blank(now_ms)
})?;
for hit in &hits.facts {
if let Some(snap) = registry.get(hit.id)
&& snap.metadata.get(NAME_KEY).map(String::as_str) == Some(name.as_str())
{
return Ok(Some(hit.id));
}
}
Ok(None)
}
fn anchored_facts(
db: &Database,
now_ms: u64,
) -> Result<Vec<(FactId, crate::FactSnapshot)>, HostError> {
let hits = db.recall(RecallQuery {
entities: &[SELF_ENTITY],
k: ANCHOR_K,
..blank(now_ms)
})?;
Ok(hits
.facts
.iter()
.filter_map(|hit| db.get(hit.id).map(|snap| (hit.id, snap)))
.collect())
}
struct SelfDescription {
text: String,
tags: Vec<String>,
owner: Option<String>,
}
impl SelfDescription {
fn agrees_with(&self, record: &DbEntry) -> bool {
self.text == record.description && self.owner == record.owner && self.tags == record.tags
}
}
fn self_description(db: &Database, now_ms: u64) -> Result<Option<SelfDescription>, HostError> {
let Some((id, snap)) = anchored_facts(db, now_ms)?.into_iter().next() else {
return Ok(None);
};
Ok(Some(SelfDescription {
text: snap.text,
tags: db.tags_of(id),
owner: snap.metadata.get(OWNER_KEY).cloned(),
}))
}
fn blank(now_ms: u64) -> RecallQuery<'static> {
RecallQuery {
text: None,
..RecallQuery::text(now_ms, "")
}
}
fn entry_of(fact: &ExportedFact) -> Option<DbEntry> {
if !fact.tags.iter().any(|t| t == ENTRY_TAG) {
return None;
}
entry_from(&fact.text, &fact.tags, &fact.metadata)
}
fn entry_from(text: &str, tags: &[String], metadata: &BTreeMap<String, String>) -> Option<DbEntry> {
let name = DbName::parse(metadata.get(NAME_KEY)?).ok()?;
Some(DbEntry {
name,
description: text.to_string(),
tags: tags.iter().filter(|t| *t != ENTRY_TAG).cloned().collect(),
owner: metadata.get(OWNER_KEY).cloned(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::workspace::testkit::{TempDir, name, workspace};
use crate::{RememberInput, WorkspaceLimits};
fn about(text: &str) -> Description<'_> {
Description {
text,
..Description::default()
}
}
fn names(entries: &[DbEntry]) -> Vec<&str> {
entries.iter().map(|e| e.name.as_str()).collect()
}
#[test]
fn the_registry_is_not_opened_until_something_needs_it() {
let tmp = TempDir::new("registry-lazy");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
ws.get(&name("chat-42"), 1_000, IfMissing::Create).unwrap();
assert!(!ws.layout().registry_path().exists());
ws.entries().unwrap();
assert!(crate::storage::database_exists(
&ws.layout().registry_path()
));
assert!(ws.close_registry());
assert!(!ws.close_registry());
}
#[test]
fn describing_twice_revises_one_record_rather_than_adding_a_second() {
let tmp = TempDir::new("registry-revise");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
let chat = name("chat-42");
ws.describe(&chat, 1_000, about("work chat about plugmem"))
.unwrap();
ws.describe(&chat, 2_000, about("work chat about releases"))
.unwrap();
let entries = ws.entries().unwrap();
assert_eq!(names(&entries), ["chat-42"]);
assert_eq!(entries[0].description, "work chat about releases");
let db = ws.get(&chat, 3_000, IfMissing::Fail).unwrap();
let own = self_description(&db, 3_000).unwrap().unwrap();
assert_eq!(own.text, "work chat about releases");
}
#[test]
fn describing_a_name_that_has_no_database_creates_it() {
let tmp = TempDir::new("registry-create");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
let chat = name("chat-42");
assert!(!ws.layout().exists(&chat));
ws.describe(&chat, 1_000, about("a brand new chat"))
.unwrap();
assert!(ws.layout().exists(&chat));
}
#[test]
fn a_database_is_found_by_what_it_is_for_and_by_who_owns_it() {
let tmp = TempDir::new("registry-find");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
ws.describe(
&name("chat-42"),
1_000,
Description {
text: "release planning and performance work on the engine",
tags: &["kind:chat"],
owner: Some("ann"),
},
)
.unwrap();
ws.describe(
&name("recipes"),
1_000,
Description {
text: "dinner ideas and shopping lists",
tags: &["kind:notes"],
owner: Some("bob"),
},
)
.unwrap();
let hits = ws.find("release planning", 4, 2_000).unwrap();
assert_eq!(hits.first().map(|e| e.name.as_str()), Some("chat-42"));
let hits = ws.find("shopping lists", 4, 2_000).unwrap();
assert_eq!(hits.first().map(|e| e.name.as_str()), Some("recipes"));
let ann = ws.find("ann", 4, 2_000).unwrap();
assert_eq!(names(&ann), ["chat-42"]);
let entry = ws.entry(&name("chat-42")).unwrap().unwrap();
assert_eq!(entry.tags, ["kind:chat"]);
assert_eq!(entry.owner.as_deref(), Some("ann"));
assert!(!entry.is_archived());
assert!(ws.entry(&name("nope")).unwrap().is_none());
}
#[test]
fn archiving_is_a_label_and_is_idempotent() {
let tmp = TempDir::new("registry-archive");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
let chat = name("chat-42");
assert!(matches!(
ws.archive(&chat, 1_000),
Err(WorkspaceError::NoSuchDatabase { .. })
));
ws.describe(
&chat,
1_000,
Description {
text: "an old project",
tags: &["kind:chat"],
owner: Some("ann"),
},
)
.unwrap();
assert!(ws.archive(&chat, 2_000).unwrap());
let entry = ws.entry(&chat).unwrap().unwrap();
assert!(entry.is_archived());
assert_eq!(entry.description, "an old project");
assert_eq!(entry.owner.as_deref(), Some("ann"));
assert!(entry.tags.contains(&"kind:chat".to_string()));
assert!(!ws.archive(&chat, 3_000).unwrap());
assert!(ws.layout().exists(&chat));
}
#[test]
fn a_deleted_registry_is_rebuilt_from_the_databases_themselves() {
let tmp = TempDir::new("registry-reindex");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
for (db, text) in [("chat-42", "release planning"), ("recipes", "dinner ideas")] {
ws.describe(
&name(db),
1_000,
Description {
text,
tags: &["kind:chat"],
owner: Some("ann"),
},
)
.unwrap();
}
ws.get(&name("scratch"), 1_000, IfMissing::Create).unwrap();
ws.close_registry();
for entry in std::fs::read_dir(ws.layout().root()).unwrap() {
let path = entry.unwrap().path();
if path.is_file() {
std::fs::remove_file(path).unwrap();
}
}
assert!(ws.entries().unwrap().is_empty());
let report = ws.reindex(2_000).unwrap();
assert_eq!(
report
.indexed
.iter()
.map(DbName::to_string)
.collect::<Vec<_>>(),
["chat-42", "recipes"]
);
assert_eq!(
report
.undescribed
.iter()
.map(DbName::to_string)
.collect::<Vec<_>>(),
["scratch"]
);
assert!(report.busy.is_empty());
let entry = ws.entry(&name("chat-42")).unwrap().unwrap();
assert_eq!(entry.description, "release planning");
assert_eq!(entry.tags, ["kind:chat"]);
assert_eq!(entry.owner.as_deref(), Some("ann"));
}
#[test]
fn reindex_names_the_databases_it_could_not_read() {
let tmp = TempDir::new("registry-reindex-busy");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
let held = name("chat-42");
ws.describe(&held, 1_000, about("release planning"))
.unwrap();
ws.describe(&name("recipes"), 1_000, about("dinner ideas"))
.unwrap();
ws.close_all();
let outsider = Database::open(ws.layout().path_of(&held), crate::Config::default())
.unwrap()
.0;
let report = ws.reindex(2_000).unwrap();
assert_eq!(report.busy, std::slice::from_ref(&held));
assert_eq!(report.indexed, [name("recipes")]);
drop(outsider);
assert_eq!(ws.reindex(3_000).unwrap().busy, []);
}
#[test]
fn verify_reports_every_way_the_registry_can_disagree() {
let tmp = TempDir::new("registry-verify");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
let agreed = name("chat-42");
ws.describe(&agreed, 1_000, about("release planning"))
.unwrap();
assert_eq!(ws.verify(2_000).unwrap(), []);
let plain = name("scratch");
ws.get(&plain, 1_000, IfMissing::Create).unwrap();
let unlisted = name("orphan");
let db = ws.get(&unlisted, 1_000, IfMissing::Create).unwrap();
write_self(&db, 1_000, about("known only to itself")).unwrap();
drop(db);
let registry = ws.registry().unwrap();
write_entry(®istry, 1_000, &name("ghost"), about("gone")).unwrap();
drop(registry);
let two = name("twins");
let db = ws.get(&two, 1_000, IfMissing::Create).unwrap();
for text in ["first claim", "second claim"] {
db.remember(RememberInput {
entity: Some(SELF_ENTITY),
..RememberInput::text(1_000, text)
})
.unwrap();
}
drop(db);
let issues = ws.verify(2_000).unwrap();
assert!(issues.contains(&WorkspaceIssue::Missing {
name: name("ghost")
}));
assert!(issues.contains(&WorkspaceIssue::Undescribed { name: plain }));
assert!(issues.contains(&WorkspaceIssue::Stale { name: unlisted }));
assert!(issues.iter().any(|i| matches!(
i,
WorkspaceIssue::AmbiguousSelf { name, facts: 2 } if name == &two
)));
assert!(!issues.iter().any(|i| matches!(
i,
WorkspaceIssue::Stale { name } | WorkspaceIssue::Undescribed { name } if name == &agreed
)));
}
#[test]
fn verify_reports_a_database_it_could_not_open() {
let tmp = TempDir::new("registry-verify-busy");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
let held = name("chat-42");
ws.describe(&held, 1_000, about("release planning"))
.unwrap();
ws.close_all();
let _outsider = Database::open(ws.layout().path_of(&held), crate::Config::default())
.unwrap()
.0;
let issues = ws.verify(2_000).unwrap();
assert!(issues.iter().any(|i| matches!(
i,
WorkspaceIssue::Unreadable { name, why } if name == &held && why.contains("chat-42")
)));
}
#[test]
fn a_stray_fact_in_the_registry_file_is_not_read_as_a_record() {
let tmp = TempDir::new("registry-stray");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
ws.describe(&name("chat-42"), 1_000, about("planning"))
.unwrap();
let registry = ws.registry().unwrap();
registry
.remember(RememberInput::text(1_000, "somebody's shopping list"))
.unwrap();
drop(registry);
assert_eq!(names(&ws.entries().unwrap()), ["chat-42"]);
}
#[test]
fn a_failure_that_is_not_a_lock_stops_a_reindex_instead_of_being_counted() {
let tmp = TempDir::new("registry-reindex-err");
let (seed, _) = workspace(&tmp, WorkspaceLimits::default());
seed.get(&name("chat-42"), 1_000, IfMissing::Create)
.unwrap();
drop(seed);
let broken: crate::Opener = Box::new(|_| Err(HostError::Embed("no provider".into())));
let ws = Workspace::new(
crate::WorkspaceLayout::new(&tmp.0),
broken,
WorkspaceLimits::default(),
);
assert!(matches!(
ws.reindex(2_000),
Err(WorkspaceError::Host(HostError::Embed(_)))
));
}
#[test]
fn a_host_failure_is_attributed_to_the_database_it_came_from() {
let tmp = TempDir::new("registry-blame");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
let chat = name("chat-42");
let locked = HostError::Locked {
path: ws.layout().path_of(&chat),
};
assert!(matches!(
ws.blame(&chat, locked),
WorkspaceError::Busy { name } if name == chat
));
assert!(matches!(
ws.blame(&chat, HostError::Embed("no".into())),
WorkspaceError::Host(HostError::Embed(_))
));
}
#[test]
fn a_stale_record_is_one_a_reindex_would_change() {
let tmp = TempDir::new("registry-stale");
let (ws, _) = workspace(&tmp, WorkspaceLimits::default());
let chat = name("chat-42");
ws.describe(&chat, 1_000, about("release planning"))
.unwrap();
let registry = ws.registry().unwrap();
write_entry(®istry, 2_000, &chat, about("something else")).unwrap();
drop(registry);
assert_eq!(
ws.verify(3_000).unwrap(),
[WorkspaceIssue::Stale { name: chat.clone() }]
);
ws.reindex(4_000).unwrap();
assert_eq!(ws.verify(5_000).unwrap(), []);
assert_eq!(
ws.entry(&chat).unwrap().unwrap().description,
"release planning"
);
}
}