Skip to main content

relay_knowledge/watcher/task_seed/
mod.rs

1use std::path::{Path, PathBuf};
2
3use crate::identity::StableHasher64;
4
5use super::WatchedRepository;
6
7/// Builds the single durable reconciliation slot for a repository's checked-out ref.
8///
9/// The fingerprint intentionally excludes the moving commit pair: while one
10/// commit update is unfinished, repeated hints coalesce into that task. Once it
11/// publishes, the same slot can be reset with the next immutable pair.
12pub fn build_commit_task_seed(
13    repository: &WatchedRepository,
14    base_commit: &str,
15    head_commit: &str,
16    tree_hash: &str,
17    now_ms: u64,
18) -> Option<crate::storage::CodeIndexTaskSeed> {
19    if base_commit.is_empty()
20        || head_commit.is_empty()
21        || tree_hash.is_empty()
22        || base_commit == head_commit
23    {
24        return None;
25    }
26    let mode = crate::domain::CodeIndexMode::incremental(base_commit, head_commit).ok()?;
27    let request = crate::domain::CodeIndexRequest {
28        repository: crate::domain::CodeRepositorySelector {
29            repository: repository.alias.clone(),
30            ref_selector: head_commit.to_owned(),
31            path_filters: Vec::new(),
32            language_filters: Vec::new(),
33        },
34        mode: mode.clone(),
35        workspace_detection: Default::default(),
36        freshness_policy: crate::domain::FreshnessPolicy::WaitUntilFresh,
37        reuse_historical: false,
38    };
39    let mut payload = serde_json::to_value(&request).ok()?;
40    if let Some(object) = payload.as_object_mut() {
41        object.insert(
42            "git_event".to_owned(),
43            serde_json::json!({
44                "kind": "ref_reconcile",
45                "ref": "HEAD",
46                "old_oid": base_commit,
47                "new_oid": head_commit,
48            }),
49        );
50    }
51    let path_filters_json = serde_json::to_string(&repository.path_filters).ok()?;
52    let language_filters_json = serde_json::to_string(&repository.language_filters).ok()?;
53    let source_scope = crate::domain::code_snapshot_scope_id(
54        &repository.repository_id,
55        tree_hash,
56        &repository.path_filters,
57        &repository.language_filters,
58    );
59
60    Some(crate::storage::CodeIndexTaskSeed {
61        repository_id: repository.repository_id.clone(),
62        alias: repository.alias.clone(),
63        ref_selector: "HEAD".to_owned(),
64        resolved_commit_sha: head_commit.to_owned(),
65        tree_hash: tree_hash.to_owned(),
66        source_scope,
67        path_filters: repository.path_filters.clone(),
68        language_filters: repository.language_filters.clone(),
69        mode,
70        input_fingerprint: format!(
71            "git_ref_reconcile:{}:HEAD:{}:{}",
72            repository.repository_id, path_filters_json, language_filters_json
73        ),
74        resource_budget: crate::domain::CodeIndexResourceBudget::default(),
75        payload_json: serde_json::to_string(&payload).ok()?,
76        now_ms,
77    })
78}
79
80pub(super) struct ChangedPathSnapshot {
81    pub path: PathBuf,
82    pub content_hash: u64,
83}
84
85/// Builds a worktree task pinned to the repository's last clean indexed base.
86pub fn build_incremental_task_seed(
87    repository: &WatchedRepository,
88    changed_paths: &[PathBuf],
89    content_fingerprint: u64,
90    now_ms: u64,
91) -> Option<crate::storage::CodeIndexTaskSeed> {
92    if changed_paths.is_empty() {
93        return None;
94    }
95    let relative_paths = changed_path_labels(repository, changed_paths);
96    if relative_paths.is_empty() {
97        return None;
98    }
99    let base_commit = immutable_worktree_base(&repository.last_indexed_commit)?;
100    let path_hash = stable_path_fingerprint(&relative_paths);
101    let task_tree_hash = format!("worktree:pending:{base_commit}");
102    let source_scope = crate::domain::code_snapshot_scope_id(
103        &repository.repository_id,
104        &task_tree_hash,
105        &repository.path_filters,
106        &repository.language_filters,
107    );
108
109    let input_fingerprint = format!(
110        "worktree_overlay:{}:{}:{}:{path_hash:016x}:{content_fingerprint:016x}",
111        repository.repository_id, task_tree_hash, source_scope,
112    );
113
114    let request = crate::domain::CodeIndexRequest {
115        repository: crate::domain::CodeRepositorySelector {
116            repository: repository.alias.clone(),
117            ref_selector: base_commit.to_owned(),
118            path_filters: Vec::new(),
119            language_filters: Vec::new(),
120        },
121        mode: crate::domain::CodeIndexMode::WorktreeOverlay,
122        workspace_detection: Default::default(),
123        freshness_policy: crate::domain::FreshnessPolicy::WaitUntilFresh,
124        reuse_historical: false,
125    };
126    let mut payload = serde_json::to_value(&request).ok()?;
127    if let Some(object) = payload.as_object_mut() {
128        object.insert(
129            "watcher".to_owned(),
130            serde_json::json!({
131                "repository_id": repository.repository_id.clone(),
132                "changed_paths": relative_paths,
133                "content_fingerprint": format!("{content_fingerprint:016x}"),
134            }),
135        );
136    }
137
138    Some(crate::storage::CodeIndexTaskSeed {
139        repository_id: repository.repository_id.clone(),
140        alias: repository.alias.clone(),
141        ref_selector: base_commit.to_owned(),
142        resolved_commit_sha: task_tree_hash.clone(),
143        tree_hash: task_tree_hash,
144        source_scope,
145        path_filters: repository.path_filters.clone(),
146        language_filters: repository.language_filters.clone(),
147        mode: crate::domain::CodeIndexMode::WorktreeOverlay,
148        input_fingerprint,
149        resource_budget: crate::domain::CodeIndexResourceBudget::default(),
150        payload_json: serde_json::to_string(&payload).ok()?,
151        now_ms,
152    })
153}
154
155/// Builds the stable periodic worktree reconciliation slot after a missed or
156/// rejected file event. The worker still scans bounded Git status; the marker
157/// is task metadata only and is never interpreted as a source path.
158pub fn build_worktree_reconcile_task_seed(
159    repository: &WatchedRepository,
160    observation_fingerprint: u64,
161    now_ms: u64,
162) -> Option<crate::storage::CodeIndexTaskSeed> {
163    let base_commit = immutable_worktree_base(&repository.last_indexed_commit)?;
164    let task_tree_hash = format!("worktree:pending:{base_commit}");
165    let source_scope = crate::domain::code_snapshot_scope_id(
166        &repository.repository_id,
167        &task_tree_hash,
168        &repository.path_filters,
169        &repository.language_filters,
170    );
171    let request = crate::domain::CodeIndexRequest {
172        repository: crate::domain::CodeRepositorySelector {
173            repository: repository.alias.clone(),
174            ref_selector: base_commit.to_owned(),
175            path_filters: Vec::new(),
176            language_filters: Vec::new(),
177        },
178        mode: crate::domain::CodeIndexMode::WorktreeOverlay,
179        workspace_detection: Default::default(),
180        freshness_policy: crate::domain::FreshnessPolicy::WaitUntilFresh,
181        reuse_historical: false,
182    };
183    let mut payload = serde_json::to_value(&request).ok()?;
184    if let Some(object) = payload.as_object_mut() {
185        object.insert(
186            "watcher".to_owned(),
187            serde_json::json!({
188                "kind": "periodic_worktree_reconcile",
189                "observation_fingerprint": format!("{observation_fingerprint:016x}"),
190            }),
191        );
192    }
193    Some(crate::storage::CodeIndexTaskSeed {
194        repository_id: repository.repository_id.clone(),
195        alias: repository.alias.clone(),
196        ref_selector: base_commit.to_owned(),
197        resolved_commit_sha: task_tree_hash.clone(),
198        tree_hash: task_tree_hash,
199        source_scope,
200        path_filters: repository.path_filters.clone(),
201        language_filters: repository.language_filters.clone(),
202        mode: crate::domain::CodeIndexMode::WorktreeOverlay,
203        input_fingerprint: format!(
204            "worktree_reconcile:{}:{}:{observation_fingerprint:016x}",
205            repository.repository_id, base_commit,
206        ),
207        resource_budget: crate::domain::CodeIndexResourceBudget::default(),
208        payload_json: serde_json::to_string(&payload).ok()?,
209        now_ms,
210    })
211}
212
213fn immutable_worktree_base(snapshot_identity: &str) -> Option<&str> {
214    crate::domain::clean_git_commit_from_snapshot_identity(snapshot_identity)
215}
216
217pub(super) fn changed_content_fingerprint(
218    repository: &WatchedRepository,
219    changes: &[&ChangedPathSnapshot],
220) -> u64 {
221    let mut entries = changes
222        .iter()
223        .filter_map(|change| {
224            let relative = change.path.strip_prefix(&repository.root).ok()?;
225            let label = path_label(relative)?;
226            Some((label, change.content_hash))
227        })
228        .collect::<Vec<_>>();
229    entries.sort();
230    entries.dedup();
231    stable_content_fingerprint(&entries)
232}
233
234pub(super) fn unreadable_path_fingerprint(path: &Path) -> u64 {
235    let label = path_label(path).unwrap_or_else(|| "<unreadable>".to_owned());
236    stable_content_fingerprint(&[(label, 0)])
237}
238
239fn changed_path_labels(repository: &WatchedRepository, changed_paths: &[PathBuf]) -> Vec<String> {
240    let mut labels = changed_paths
241        .iter()
242        .filter_map(|path| path.strip_prefix(&repository.root).ok())
243        .filter_map(path_label)
244        .collect::<Vec<_>>();
245    labels.sort();
246    labels.dedup();
247    labels
248}
249
250fn path_label(path: &Path) -> Option<String> {
251    let value = path
252        .to_string_lossy()
253        .replace(std::path::MAIN_SEPARATOR, "/");
254    (!value.is_empty()).then_some(value)
255}
256
257fn stable_path_fingerprint(paths: &[String]) -> u64 {
258    let mut hasher = StableHasher64::new();
259    for path in paths {
260        hasher.update(path.as_bytes());
261        hasher.update(&[0]);
262    }
263    hasher.finish()
264}
265
266fn stable_content_fingerprint(entries: &[(String, u64)]) -> u64 {
267    let mut hasher = StableHasher64::new();
268    for (path, content_hash) in entries {
269        hasher.update(path.as_bytes());
270        hasher.update(&[0]);
271        hasher.update(&content_hash.to_le_bytes());
272        hasher.update(&[0]);
273    }
274    hasher.finish()
275}
276
277#[cfg(test)]
278#[path = "mod_tests.rs"]
279mod tests;