feat(identity): route captures by the identity the client resolved (#708)

/hook, /hook/batch items and /handoff accept identity / identity_src;
only the explicit and git_remote rungs route by identity, anything else
routes by name as before. The path-keyed project cache includes the
identity, so one path holding two repositories on two machines does
not answer for both.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_014aqFKAuGVkuBoewmpA3cx9
This commit is contained in:
luisfnicolau
2026-09-28 09:43:22 -04:00
co-authored by Claude Opus 5.5
parent b2621692b1
commit 2185d74116
4 changed files with 353 additions and 15 deletions
@@ -734,6 +734,27 @@ pub fn discover_repo_root(start: &Path) -> Result<PathBuf, BootstrapError> {
.ok_or_else(|| BootstrapError::NotARepo(start.to_path_buf()))
}
/// URLs of the `upstream` and `origin` remotes of the repository containing
/// `start`, for resolving its repository identity (#708).
///
/// Only those two names are read. Remote names are personal, so picking some
/// other one would give the same repository a different identity per person;
/// see `ai_memory_core::repository_identity`. A linked worktree reads the
/// main repository's remotes, which it shares. `(None, None)` when `start` is
/// not inside a repository.
#[must_use]
pub fn read_identity_remotes(start: &Path) -> (Option<String>, Option<String>) {
let Ok(repo) = git2::Repository::discover(start) else {
return (None, None);
};
let url = |name: &str| {
repo.find_remote(name)
.ok()
.and_then(|remote| remote.url().ok().map(str::to_owned))
};
(url("upstream"), url("origin"))
}
/// Like [`discover_repo_root`], but when `start` is inside a git
/// worktree the function follows the `.git` commondir pointer back to
/// the **main** repository root rather than returning the worktree's
+1 -1
View File
@@ -56,7 +56,7 @@ pub use bootstrap::{
Bootstrap, BootstrapConfig, BootstrapError, BootstrapOutcome, BootstrapSource,
DEFAULT_CHUNK_INPUT_TOKENS, ProjectNameStrategy, SourceCounts, SourceKind, collect_sources,
derive_project_name, discover_main_repo_root, discover_repo_root, effective_chunk_budget,
plan_bootstrap_chunks, prune_sources_to_budget,
plan_bootstrap_chunks, prune_sources_to_budget, read_identity_remotes,
};
pub use cold_cluster::{adaptive_eps, cosine_distance, dbscan};
pub use compaction::build_compacted_markdown;
+20
View File
@@ -89,6 +89,14 @@ pub struct HookQuery {
/// `.ai-memory.toml` named it, `repo-root` when the host hook derived it
/// from the enclosing checkout. Absent on older clients (#394).
pub project_src: Option<String>,
/// Repository identity the client resolved for this checkout (#708):
/// an explicit marker `identity`, or a normalised git remote. Paired
/// with `identity_src`. Absent on older clients and for checkouts that
/// declare a `project`, which keep routing by name.
pub identity: Option<String>,
/// Which rung produced `identity`: `explicit` or `git_remote`. Anything
/// else, or a malformed identity, is ignored and the event routes by name.
pub identity_src: Option<String>,
}
/// Coalesced view of an incoming hook event after light parsing of the
@@ -117,6 +125,10 @@ pub struct HookEnvelope {
/// Where `project_override` came from. Always [`ProjectSource::Unspecified`]
/// when there is no override, so the two can never disagree (#394).
pub project_source: ProjectSource,
/// Repository identity from the client, already validated by
/// [`ai_memory_core::repository_identity::accept_wire_identity`]. `None`
/// routes by project name, as every event did before identities existed.
pub identity: Option<ai_memory_core::repository_identity::RepositoryIdentity>,
/// Whether this project opted into `drop_subagent_captures` via its
/// `.ai-memory.toml` (forwarded as the `drop_subagent` query flag). The
/// ingest router consults this per-event so the drop is scoped to the
@@ -164,6 +176,7 @@ impl std::fmt::Debug for HookEnvelope {
.field("workspace_override", &self.workspace_override)
.field("project_override", &self.project_override)
.field("project_strategy", &self.project_strategy)
.field("identity", &self.identity)
.field("drop_subagent_requested", &self.drop_subagent_requested)
.field(
"recall_default_global_requested",
@@ -502,6 +515,12 @@ impl HookEnvelope {
} else {
ProjectSource::Unspecified
};
let identity = match (query.identity.as_deref(), query.identity_src.as_deref()) {
(Some(identity), Some(source)) => {
ai_memory_core::repository_identity::accept_wire_identity(identity, source)
}
_ => None,
};
let drop_subagent_requested = query_flag_truthy(query.drop_subagent.as_deref());
let recall_default_global_requested = query_flag_truthy(query.default_global.as_deref());
let all_owners_requested = query_flag_truthy(query.all_owners.as_deref());
@@ -563,6 +582,7 @@ impl HookEnvelope {
project_override,
project_strategy,
project_source,
identity,
drop_subagent_requested,
recall_default_global_requested,
all_owners_requested,
+311 -14
View File
@@ -88,8 +88,8 @@ pub const DEFAULT_PROJECT_CACHE_MAX_ENTRIES: usize = 4096;
const SUBAGENT_SESSIONS_MAX: usize = 4096;
/// Resolved-project cache key:
/// `(cwd, workspace_override, project_override, project_strategy)`.
pub type ProjectCacheKey = (String, String, String, String);
/// `(cwd, workspace_override, project_override, project_strategy, identity)`.
pub type ProjectCacheKey = (String, String, String, String, String);
/// Shared bounded resolved-project cache.
pub type ProjectCache = Arc<tokio::sync::Mutex<ProjectCacheStore>>;
@@ -1191,6 +1191,7 @@ async fn should_drop_subagent(
env.workspace_override.as_deref(),
env.project_override.as_deref(),
env.project_strategy,
env.identity.as_ref(),
viewer,
)
.await
@@ -1267,6 +1268,24 @@ pub struct HandoffQuery {
pub managed_run: Option<String>,
/// Native session identifier observed in the SessionStart payload.
pub session_id: Option<String>,
/// Repository identity the client resolved (#708). Same contract as
/// [`crate::payload::HookQuery::identity`].
pub identity: Option<String>,
/// Rung that produced `identity`; see
/// [`crate::payload::HookQuery::identity_src`].
pub identity_src: Option<String>,
}
impl HandoffQuery {
/// The validated repository identity, or `None` to route by name.
fn repository_identity(
&self,
) -> Option<ai_memory_core::repository_identity::RepositoryIdentity> {
ai_memory_core::repository_identity::accept_wire_identity(
self.identity.as_deref()?,
self.identity_src.as_deref()?,
)
}
}
/// Synchronous endpoint used by `session-start.sh` to discover any
@@ -1353,6 +1372,7 @@ async fn fetch_and_accept_handoff_at(
query.workspace.as_deref(),
query.project.as_deref(),
ProjectStrategy::parse(query.project_strategy.as_deref()),
query.repository_identity().as_ref(),
viewer,
)
.await?;
@@ -2091,12 +2111,19 @@ fn cache_key_for(
workspace_override: Option<&str>,
project_override: Option<&str>,
project_strategy: ProjectStrategy,
) -> (String, String, String, String) {
identity: Option<&ai_memory_core::repository_identity::RepositoryIdentity>,
) -> ProjectCacheKey {
(
cwd_norm.unwrap_or_default().to_string(),
workspace_override.unwrap_or_default().to_string(),
project_override.unwrap_or_default().to_string(),
project_strategy.as_str().to_string(),
// Two repositories can sit at the same path on two machines. Without
// the identity in the key, whichever resolved first would answer for
// both.
identity
.map(|i| format!("{}:{}", i.source.as_str(), i.identity))
.unwrap_or_default(),
)
}
@@ -2145,6 +2172,7 @@ async fn resolve_project_ids_inner(
workspace_override: Option<&str>,
project_override: Option<&str>,
project_strategy: ProjectStrategy,
identity: Option<&ai_memory_core::repository_identity::RepositoryIdentity>,
creator: Option<ai_memory_core::UserId>,
) -> anyhow::Result<(WorkspaceId, ProjectId)> {
let cwd_raw = cwd.filter(|s| !s.is_empty());
@@ -2156,11 +2184,15 @@ async fn resolve_project_ids_inner(
return Ok((state.workspace_id, state.project_id));
}
// Only the rungs that carry something a name does not route by identity:
// a declared `project` and a folder name keep routing by name.
let identity = identity.filter(|i| i.source.routes_by_identity());
let cache_key = cache_key_for(
cwd_norm.as_deref(),
workspace_override,
project_override,
project_strategy,
identity,
);
{
@@ -2343,17 +2375,44 @@ async fn resolve_project_ids_inner(
// The match is keyed on the actual cwd (`cwd_norm`), not the stored
// `repo_path`: `repo_path` is now the git root or None (issue #103),
// whereas cwd->parent matching needs the full deep path.
let proj = if project_override.is_none()
&& let Some(rp) = cwd_norm.as_deref().filter(|s| !s.is_empty())
&& let Some((parent_id, parent_name)) = state
let parent = match cwd_norm.as_deref().filter(|s| !s.is_empty()) {
Some(rp) if project_override.is_none() => state
.reader
.find_project_by_cwd_prefix(ws, rp.to_string(), state.home_dir.as_deref())
.await
.map_err(|e| anyhow::anyhow!("find_project_by_cwd_prefix: {e}"))?
.map_err(|e| anyhow::anyhow!("find_project_by_cwd_prefix: {e}"))?,
_ => None,
};
let proj = if let Some(identity) = identity {
// A repository identity decides the project before any name does: the
// project already carrying it wins, whatever it is called, and the
// name (or the cwd-prefix parent) is only the candidate for a project
// that has not been claimed yet. One writer transaction, so two
// captures racing on a new repository cannot both create it.
let (proj, resolution) = state
.writer
.resolve_project_by_identity(
ws,
identity.clone(),
project_name,
repo_path,
parent.map(|(parent_id, _)| parent_id),
creator,
)
.await
.map_err(|e| anyhow::anyhow!("resolve_project_by_identity: {e}"))?;
debug!(
identity = %identity.identity,
source = identity.source.as_str(),
?resolution,
"hook router: resolved project by repository identity"
);
proj
} else if let Some((parent_id, parent_name)) = parent
&& parent_name != project_name
{
debug!(
cwd = rp,
cwd = ?cwd_norm,
derived = %project_name,
parent = %parent_name,
"hook router: cwd inside existing project — using parent instead of \
@@ -2397,6 +2456,7 @@ async fn resolve_project_ids(
project_override,
project_strategy,
None,
None,
)
.await?;
if has_publishable_scope_hint(cwd, project_override) {
@@ -2763,6 +2823,7 @@ async fn process_authorized(
env.workspace_override.as_deref(),
env.project_override.as_deref(),
env.project_strategy,
env.identity.as_ref(),
viewer,
)
.await?
@@ -2849,6 +2910,7 @@ async fn process_authorized(
env.workspace_override.as_deref(),
env.project_override.as_deref(),
env.project_strategy,
env.identity.as_ref(),
);
let mut attempts = 0;
// Keep the successful keyed-ingest gate until every downstream effect has
@@ -2912,6 +2974,7 @@ async fn process_authorized(
env.workspace_override.as_deref(),
env.project_override.as_deref(),
env.project_strategy,
env.identity.as_ref(),
viewer,
)
.await?;
@@ -3846,6 +3909,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
}
}
@@ -4300,6 +4365,178 @@ mod tests {
);
}
/// Session start routes by the same identity a capture does, so its
/// handoff and briefing come from the project the captures land in. Same
/// validation: a rung that routes by name, or a malformed value, is `None`.
#[test]
fn handoff_query_accepts_identity_on_the_same_terms_as_a_capture() {
let query = |identity: &str, source: &str| HandoffQuery {
identity: Some(identity.to_owned()),
identity_src: Some(source.to_owned()),
..Default::default()
};
let got = query("github.com/orga/api", "git_remote")
.repository_identity()
.unwrap();
assert_eq!(got.identity, "github.com/orga/api");
assert!(
query("github.com/orga/api", "manifest")
.repository_identity()
.is_none()
);
assert!(
query("GitHub.com/Orga/API", "git_remote")
.repository_identity()
.is_none()
);
assert!(HandoffQuery::default().repository_identity().is_none());
}
/// A capture that carries the repository identity its client resolved.
async fn capture_with_identity(
state: &HookState,
cwd: &std::path::Path,
session: &str,
identity: Option<(&str, &str)>,
) -> ProjectId {
let env = HookEnvelope::from_query_and_body(
HookQuery {
event: "user-prompt-submit".into(),
agent: Some("claude-code".into()),
identity: identity.map(|(identity, _)| identity.to_owned()),
identity_src: identity.map(|(_, source)| source.to_owned()),
..Default::default()
},
serde_json::json!({
"session_id": session,
"cwd": cwd.to_string_lossy(),
"prompt": "hello",
}),
);
process_authorized(
state,
env,
None,
ai_memory_core::AuthLevel::User,
Vec::new(),
None,
)
.await
.unwrap();
state
.reader
.observations_for_session(session.parse().unwrap())
.await
.unwrap()[0]
.project_id
}
/// #708's collision, end to end: two unrelated repositories both checked
/// out as `api/` used to land in one project. With the client's identity
/// they land in two, and the same repository under a different folder name
/// lands back in its own. The two `api/` checkouts share a folder name and
/// nothing else, so the path-keyed project cache must not answer for both.
#[tokio::test]
async fn two_api_checkouts_with_different_remotes_get_two_projects() {
let tmp = TempDir::new().unwrap();
let state = make_state(&tmp).await;
let org_a = tmp.path().join("org-a").join("api");
let org_b = tmp.path().join("org-b").join("api");
let clone = tmp.path().join("elsewhere").join("acme-api");
for dir in [&org_a, &org_b, &clone] {
std::fs::create_dir_all(dir).unwrap();
}
let a = capture_with_identity(
&state,
&org_a,
&SessionId::new().to_string(),
Some(("github.com/orga/api", "git_remote")),
)
.await;
let b = capture_with_identity(
&state,
&org_b,
&SessionId::new().to_string(),
Some(("github.com/orgb/api", "git_remote")),
)
.await;
assert_ne!(a, b, "unrelated repositories must not share a project");
assert_eq!(
state
.reader
.project_name_by_id(state.workspace_id, b)
.await
.unwrap(),
Some("orgb-api".to_owned())
);
let again = capture_with_identity(
&state,
&clone,
&SessionId::new().to_string(),
Some(("github.com/orga/api", "git_remote")),
)
.await;
assert_eq!(again, a, "one repository, one project, whatever the folder");
// A client that sends nothing, or a rung that routes by name, keeps
// today's routing: the folder name.
let plain = tmp.path().join("plain").join("api");
std::fs::create_dir_all(&plain).unwrap();
let by_name =
capture_with_identity(&state, &plain, &SessionId::new().to_string(), None).await;
assert_eq!(
by_name, a,
"no identity routes by the folder name, as before"
);
let declared = capture_with_identity(
&state,
&plain,
&SessionId::new().to_string(),
Some(("github.com/orgc/api", "manifest")),
)
.await;
assert_eq!(declared, a, "a manifest identity is ignored on the wire");
}
/// Two machines sharing one server can hold unrelated repositories at the
/// same path (`~/src/api` on each). The project cache is keyed by path, so
/// without the identity in its key whichever resolved first would answer
/// for both, and the second machine's captures would land in the first's
/// project.
#[tokio::test]
async fn one_path_with_two_remotes_is_two_projects() {
let tmp = TempDir::new().unwrap();
let state = make_state(&tmp).await;
let shared = tmp.path().join("src").join("api");
std::fs::create_dir_all(&shared).unwrap();
let first = capture_with_identity(
&state,
&shared,
&SessionId::new().to_string(),
Some(("github.com/orga/api", "git_remote")),
)
.await;
let second = capture_with_identity(
&state,
&shared,
&SessionId::new().to_string(),
Some(("github.com/orgb/api", "git_remote")),
)
.await;
assert_ne!(first, second, "the cache answered for the wrong repository");
// Control: the same repository at that path keeps resolving to its own.
let again = capture_with_identity(
&state,
&shared,
&SessionId::new().to_string(),
Some(("github.com/orga/api", "git_remote")),
)
.await;
assert_eq!(again, first);
}
/// A capture, with the viewer it is authenticated as.
async fn capture_as(
state: &HookState,
@@ -7781,10 +8018,11 @@ mod tests {
String::new(),
String::new(),
ProjectStrategy::Basename.as_str().to_string(),
String::new(),
);
assert!(
cache.contains_key(&key),
"cache keyed by (cwd, ws_override, proj_override, project_strategy)"
"cache keyed by (cwd, ws_override, proj_override, project_strategy, identity)"
);
}
@@ -7811,9 +8049,27 @@ mod tests {
#[test]
fn project_cache_store_evicts_oldest_untouched_entry() {
let mut cache = ProjectCacheStore::new(2);
let key_a = ("/a".into(), String::new(), String::new(), "basename".into());
let key_b = ("/b".into(), String::new(), String::new(), "basename".into());
let key_c = ("/c".into(), String::new(), String::new(), "basename".into());
let key_a = (
"/a".into(),
String::new(),
String::new(),
"basename".into(),
String::new(),
);
let key_b = (
"/b".into(),
String::new(),
String::new(),
"basename".into(),
String::new(),
);
let key_c = (
"/c".into(),
String::new(),
String::new(),
"basename".into(),
String::new(),
);
cache.insert(key_a.clone(), (WorkspaceId::new(), ProjectId::new()));
cache.insert(key_b.clone(), (WorkspaceId::new(), ProjectId::new()));
@@ -7834,8 +8090,20 @@ mod tests {
let mut cache = ProjectCacheStore::new(4);
let doomed_ws = WorkspaceId::new();
let kept_ws = WorkspaceId::new();
let key_a = ("/a".into(), String::new(), String::new(), "basename".into());
let key_b = ("/b".into(), String::new(), String::new(), "basename".into());
let key_a = (
"/a".into(),
String::new(),
String::new(),
"basename".into(),
String::new(),
);
let key_b = (
"/b".into(),
String::new(),
String::new(),
"basename".into(),
String::new(),
);
cache.insert(key_a.clone(), (doomed_ws, ProjectId::new()));
cache.insert(key_b.clone(), (kept_ws, ProjectId::new()));
@@ -8096,6 +8364,7 @@ mod tests {
Some("default"),
Some("scratch"),
ProjectStrategy::Basename,
None,
);
let mut cache = state.project_cache.lock().await;
assert_eq!(cache.get(&cache_key), Some((cached_ws, cached_proj)));
@@ -10028,6 +10297,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
}),
Some(axum::Extension(ai_memory_core::ActorContext {
issuer: Some("https://idp.example".into()),
@@ -10638,6 +10909,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
};
let state = Arc::new(state);
let session_start = |viewer: Option<ai_memory_core::UserId>| {
@@ -10719,6 +10992,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
};
let state = Arc::new(state);
@@ -10893,6 +11168,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
}),
None,
None,
@@ -11709,6 +11986,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
},
None,
Vec::new(),
@@ -11783,6 +12062,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
},
None,
Vec::new(),
@@ -11850,6 +12131,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: Some(session_id.into()),
identity: None,
identity_src: None,
};
let empty_sid = "empty-native-session";
let rendered = fetch_and_accept_handoff(&state, query(empty_sid), None, Vec::new(), None)
@@ -11987,6 +12270,8 @@ mod tests {
briefing_budget: None,
managed_run: Some(run.run_id.to_string()),
session_id: Some("native-2".into()),
identity: None,
identity_src: None,
};
let rendered = fetch_and_accept_handoff(&state, query.clone(), None, Vec::new(), None)
.await
@@ -12092,6 +12377,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
};
let named = ai_memory_core::ActorContext {
@@ -12168,6 +12455,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
};
let rendered =
@@ -12228,6 +12517,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
};
// Non-truthy opt-in: no handoff pending, nothing to inject.
@@ -12374,6 +12665,8 @@ mod tests {
briefing_budget: None,
managed_run: Some(kimi.run_id.to_string()),
session_id: Some("kimi-session".into()),
identity: None,
identity_src: None,
};
let rendered =
@@ -12586,6 +12879,8 @@ mod tests {
briefing_budget: None,
managed_run: None,
session_id: None,
identity: None,
identity_src: None,
},
None,
Vec::new(),
@@ -12710,6 +13005,7 @@ mod tests {
String::new(),
String::new(),
strat.clone(),
String::new(),
)),
"cache key must remain case-folded for #806 stickiness"
);
@@ -12719,6 +13015,7 @@ mod tests {
String::new(),
String::new(),
strat,
String::new(),
)),
"cache key must not carry the original-case basename"
);