relay_knowledge/watcher/task_seed/
mod.rs1use std::path::{Path, PathBuf};
2
3use crate::identity::StableHasher64;
4
5use super::WatchedRepository;
6
7pub 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
85pub 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
155pub 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;