Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
212 changes: 85 additions & 127 deletions crates/tracedecay-code-index-runtime/src/git_transactions/owner.rs
Original file line number Diff line number Diff line change
Expand Up @@ -204,46 +204,34 @@ impl DaemonGitAuthoritySource for ProductionDaemonGitAuthoritySource {
}
}

enum DaemonGitAuthority {
Installed(Arc<dyn DaemonGitAuthoritySource>),
/// Daemon shutdown revoked the authority of an owner a request may
/// already hold.
Revoked,
}

/// A slot is only ever constructed around an installed source, so a mounted
/// owner always carries authority until shutdown revokes it.
struct DaemonGitAuthoritySlot {
source: ProfiledStdRwLock<Option<Arc<dyn DaemonGitAuthoritySource>>>,
source: ProfiledStdRwLock<DaemonGitAuthority>,
}

impl Default for DaemonGitAuthoritySlot {
fn default() -> Self {
impl DaemonGitAuthoritySlot {
fn new(source: Arc<dyn DaemonGitAuthoritySource>) -> Self {
Self {
source: hotpath::rw_lock!(
std::sync::RwLock::new(None),
std::sync::RwLock::new(DaemonGitAuthority::Installed(source)),
label = "daemon.git.tx.authority_source"
),
}
}
}

impl DaemonGitAuthoritySlot {
fn install(
&self,
source: Arc<dyn DaemonGitAuthoritySource>,
) -> Result<(), GitIndexTransactionPortError> {
fn set(&self, authority: DaemonGitAuthority) -> Result<(), GitIndexTransactionPortError> {
*self
.source
.write()
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)? = Some(source);
Ok(())
}

fn installed(&self) -> Result<bool, GitIndexTransactionPortError> {
Ok(self
.source
.read()
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?
.is_some())
}

fn clear(&self) -> Result<(), GitIndexTransactionPortError> {
self.source
.write()
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?
.take();
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)? = authority;
Ok(())
}
}
Expand All @@ -253,13 +241,16 @@ impl DaemonGitAuthoritySource for DaemonGitAuthoritySlot {
&self,
capability_id: &CapabilityId,
) -> Result<DaemonGitAuthorityStateV1, GitIndexTransactionPortError> {
let source = self
let source = match &*self
.source
.read()
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?
.as_ref()
.ok_or(GitIndexTransactionPortError::PolicyDenied)?
.clone();
{
DaemonGitAuthority::Installed(source) => Arc::clone(source),
DaemonGitAuthority::Revoked => {
return Err(GitIndexTransactionPortError::DaemonUnavailable);
}
};
source.current_capability(capability_id)
}
}
Expand Down Expand Up @@ -415,6 +406,17 @@ struct ServiceEntry {
authority: Arc<DaemonGitAuthoritySlot>,
}

impl ServiceEntry {
fn owner(&self) -> DaemonGitInvocationOwner {
DaemonGitInvocationOwner {
project_id: self.project_id.clone(),
repository_root: self.repository_root.clone(),
service: Arc::clone(&self.service),
authority: Arc::clone(&self.authority),
}
}
}

#[derive(Clone, Debug, PartialEq, Eq, Hash)]
struct ServiceKey {
database_path: PathBuf,
Expand Down Expand Up @@ -481,40 +483,57 @@ pub struct DaemonGitIndexShutdownReceiptV1 {
}

impl DaemonGitIndexTransactionServiceRegistry {
#[hotpath::measure(label = "daemon.git.tx.ensure", future = true)]
pub async fn ensure(
/// Mounts the owner for one exact project/worktree with its authority in
/// a single publication: the service starts around an already installed
/// authority and becomes resolvable only once both exist. Remounting the
/// same identity reuses the started service and replaces its authority.
#[hotpath::measure(label = "daemon.git.tx.mount", future = true)]
pub async fn mount(
&self,
database: RegisteredGlobalDbLeaseV1,
repository_root: PathBuf,
project_id: ProjectId,
access: ProjectSourceAccessSnapshot,
observed_at: UtcMicros,
) -> Result<Arc<DaemonProjectGitIndexTransactionService>, GitIndexTransactionPortError> {
if self.shutdown_fenced.load(Ordering::SeqCst) {
return Err(GitIndexTransactionPortError::DaemonUnavailable);
runtime: tokio::runtime::Handle,
) -> Result<DaemonGitInvocationOwner, GitIndexTransactionPortError> {
let project_id = access.scope.project_id.clone();
if access.binding.authority != AuthorityRef::Project(project_id.clone()) {
return Err(GitIndexTransactionPortError::PolicyDenied);
}
// `RegisteredGlobalDb` already carries the canonical path admitted by
// its retained runtime authority. Do not rediscover identity through
// the filesystem: a newly opened SQLite shard may not have
// materialized its directory entry yet.
let database_path = database.db_path().to_path_buf();
self.ensure_with(
let source = Arc::new(ProductionDaemonGitAuthoritySource {
access,
configuration:
OwnedGlobalDbConfigurationControlStore::from_registered_project_runtime_db(
database.clone(),
),
catalog: self.catalog.clone(),
runtime,
});
self.mount_with(
database_path,
repository_root,
project_id,
source,
observed_at,
|| self.stores.ensure(database),
)
.await
}

async fn ensure_with<F>(
pub(super) async fn mount_with<F>(
&self,
database_path: PathBuf,
repository_root: PathBuf,
project_id: ProjectId,
source: Arc<dyn DaemonGitAuthoritySource>,
observed_at: UtcMicros,
open_store: F,
) -> Result<Arc<DaemonProjectGitIndexTransactionService>, GitIndexTransactionPortError>
) -> Result<DaemonGitInvocationOwner, GitIndexTransactionPortError>
where
F: FnOnce() -> tracedecay_store::GitIndexTransactionStoreResult<
super::SharedDaemonGitIndexTransactionStore,
Expand All @@ -525,29 +544,32 @@ impl DaemonGitIndexTransactionServiceRegistry {
}
let repository_root = canonicalize_repository_root(&repository_root)
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?;
if let Some(service) = self
.existing(&database_path, &repository_root, &project_id)
.await?
{
return Ok(service);
}
let service_key = ServiceKey::new(&database_path, &project_id, &repository_root);

let _creation = self.creation_gate.lock().await;
if self.shutdown_fenced.load(Ordering::SeqCst) {
return Err(GitIndexTransactionPortError::DaemonUnavailable);
}
if let Some(service) = self
.existing(&database_path, &repository_root, &project_id)
.await?
{
return Ok(service);
let services = self.services.lock().await;
// Another identity mounted at this root makes the root ambiguous.
if services
.iter()
.any(|(key, entry)| *key != service_key && entry.repository_root == repository_root)
{
return Err(GitIndexTransactionPortError::PolicyDenied);
}
if let Some(entry) = services.get(&service_key) {
entry.authority.set(DaemonGitAuthority::Installed(source))?;
return Ok(entry.owner());
}
}

// Open/retain the store actor under the creation gate before native
// recovery runs on a blocking thread. Later ensures reuse this actor.
// recovery runs on a blocking thread. Later mounts reuse this actor.
let store = open_store().map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?;
let native_root = repository_root.clone();
let authority = Arc::new(DaemonGitAuthoritySlot::default());
let authority = Arc::new(DaemonGitAuthoritySlot::new(source));
let service_authority = Arc::clone(&authority);
let mutation_queue = Arc::clone(&self.mutation_queue);
let (project_id, service) = tokio::task::spawn_blocking(move || {
Expand All @@ -566,71 +588,15 @@ impl DaemonGitIndexTransactionServiceRegistry {
})
.await
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)??;
let service = Arc::new(service);
let service_key = ServiceKey::new(&database_path, &project_id, &repository_root);
self.services.lock().await.insert(
service_key,
ServiceEntry {
project_id,
repository_root,
service: Arc::clone(&service),
authority,
},
);
Ok(service)
}

async fn existing(
&self,
database_path: &Path,
repository_root: &Path,
project_id: &ProjectId,
) -> Result<Option<Arc<DaemonProjectGitIndexTransactionService>>, GitIndexTransactionPortError>
{
let services = self.services.lock().await;
Ok(services
.get(&ServiceKey::new(database_path, project_id, repository_root))
.map(|entry| Arc::clone(&entry.service)))
}

#[hotpath::measure(label = "daemon.git.tx.install_authority", future = true)]
pub async fn install_authority(
&self,
repository_root: &std::path::Path,
access: ProjectSourceAccessSnapshot,
configuration_database: RegisteredGlobalDbLeaseV1,
runtime: tokio::runtime::Handle,
) -> Result<(), GitIndexTransactionPortError> {
if self.shutdown_fenced.load(Ordering::SeqCst) {
return Err(GitIndexTransactionPortError::DaemonUnavailable);
}
let repository_root = canonicalize_repository_root(repository_root)
.map_err(|_| GitIndexTransactionPortError::DaemonUnavailable)?;
let services = self.services.lock().await;
let mut matches = services
.values()
.filter(|entry| entry.repository_root == repository_root);
let Some(entry) = matches.next() else {
return Err(GitIndexTransactionPortError::PolicyDenied);
let entry = ServiceEntry {
project_id,
repository_root,
service: Arc::new(service),
authority,
};
if matches.next().is_some()
|| access.scope.project_id != entry.project_id
|| access.binding.authority != AuthorityRef::Project(entry.project_id.clone())
{
return Err(GitIndexTransactionPortError::PolicyDenied);
}
entry
.authority
.install(Arc::new(ProductionDaemonGitAuthoritySource {
access,
configuration:
OwnedGlobalDbConfigurationControlStore::from_registered_project_runtime_db(
configuration_database,
),
catalog: self.catalog.clone(),
runtime,
}))?;
Ok(())
let owner = entry.owner();
self.services.lock().await.insert(service_key, entry);
Ok(owner)
}

/// Retires every invocation owner attached to one exact project-session
Expand All @@ -654,9 +620,6 @@ impl DaemonGitIndexTransactionServiceRegistry {

/// Resolve only an owner already mounted by project-open admission.
/// Missing and ambiguous roots deliberately share the same outcome.
/// Project open publishes the service before it installs the service's
/// authority; until then the owner is still mounting, and resolving it
/// would answer every read as a policy denial.
pub async fn for_repository_root(
&self,
repository_root: &std::path::Path,
Expand All @@ -673,15 +636,10 @@ impl DaemonGitIndexTransactionServiceRegistry {
let Some(entry) = matches.next() else {
return Ok(None);
};
if matches.next().is_some() || !entry.authority.installed()? {
if matches.next().is_some() {
return Ok(None);
}
Ok(Some(DaemonGitInvocationOwner {
project_id: entry.project_id.clone(),
repository_root: entry.repository_root.clone(),
service: Arc::clone(&entry.service),
authority: Arc::clone(&entry.authority),
}))
Ok(Some(entry.owner()))
}

#[hotpath::measure(label = "daemon.git.tx.shutdown", future = true)]
Expand All @@ -696,7 +654,7 @@ impl DaemonGitIndexTransactionServiceRegistry {
let services = {
let mut retained = self.services.lock().await;
for entry in retained.values() {
entry.authority.clear()?;
entry.authority.set(DaemonGitAuthority::Revoked)?;
}
retained.drain().map(|(_, entry)| entry).collect::<Vec<_>>()
};
Expand Down
Loading
Loading