1use std::collections::HashSet;
4use std::path::Path;
5use std::path::PathBuf;
6use std::sync::OnceLock;
7
8use ito_config::{ConfigContext, load_cascading_project_config, resolve_audit_mirror_settings};
9use ito_domain::audit::event::AuditEvent;
10use ito_domain::audit::writer::AuditWriter;
11#[cfg(feature = "backend")]
12use ito_domain::backend::{BackendEventIngestClient, EventBatch};
13
14#[cfg(feature = "backend")]
15use crate::backend_client::idempotency_key;
16#[cfg(feature = "backend")]
17use crate::backend_http::BackendHttpClient;
18use crate::process::{ProcessRequest, ProcessRunner, SystemProcessRunner};
19#[cfg(feature = "backend")]
20use crate::repository_runtime::{PersistenceMode, resolve_repository_runtime};
21
22use super::mirror::{
23 InternalBranchLogRead, append_jsonl_to_internal_branch, read_internal_branch_log,
24};
25use super::writer::{
26 append_event_to_file, audit_log_path, parse_events_from_jsonl, read_events_from_path,
27};
28
29#[derive(Debug, Clone, PartialEq, Eq)]
31pub enum AuditStorageLocation {
32 Filesystem(PathBuf),
34 Other(String),
36}
37
38pub trait AuditEventStore: AuditWriter + Send + Sync {
40 fn read_all(&self) -> Vec<AuditEvent>;
42
43 fn location(&self) -> AuditStorageLocation;
45}
46
47pub fn audit_storage_location_key(location: &AuditStorageLocation) -> String {
49 match location {
50 AuditStorageLocation::Filesystem(path) => format!("fs:{}", path.display()),
51 AuditStorageLocation::Other(label) => format!("other:{label}"),
52 }
53}
54
55#[cfg(feature = "backend")]
56struct BackendAuditStore {
57 client: BackendHttpClient,
58}
59
60#[cfg(feature = "backend")]
61impl BackendAuditStore {
62 fn new(client: BackendHttpClient) -> Self {
63 Self { client }
64 }
65}
66
67#[cfg(feature = "backend")]
68impl AuditWriter for BackendAuditStore {
69 fn append(&self, event: &AuditEvent) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
70 let batch = EventBatch {
71 events: vec![event.clone()],
72 idempotency_key: idempotency_key("audit-write"),
73 };
74
75 if let Err(err) = self.client.ingest(&batch) {
76 tracing::warn!("backend audit write failed: {err}");
77 }
78
79 Ok(())
80 }
81}
82
83#[cfg(feature = "backend")]
84impl AuditEventStore for BackendAuditStore {
85 fn read_all(&self) -> Vec<AuditEvent> {
86 match self.client.list_audit_events() {
87 Ok(events) => events,
88 Err(err) => {
89 tracing::warn!("backend audit read failed: {err}");
90 Vec::new()
91 }
92 }
93 }
94
95 fn location(&self) -> AuditStorageLocation {
96 AuditStorageLocation::Other("backend".to_string())
97 }
98}
99
100struct LocalAuditStore {
101 ito_path: PathBuf,
102 branch: String,
103 fallback_path: PathBuf,
104 legacy_migration_done: OnceLock<()>,
105}
106
107impl LocalAuditStore {
108 fn new(ito_path: &Path, branch: String, fallback_path: PathBuf) -> Self {
109 Self {
110 ito_path: ito_path.to_path_buf(),
111 branch,
112 fallback_path,
113 legacy_migration_done: OnceLock::new(),
114 }
115 }
116
117 fn repo_root(&self) -> Option<&Path> {
118 self.ito_path.parent()
119 }
120
121 fn append_to_branch(&self, event: &AuditEvent) -> Result<(), String> {
122 let Some(repo_root) = self.repo_root() else {
123 return Err("unable to resolve project root for internal audit branch".to_string());
124 };
125 let json = serde_json::to_string(event)
126 .map_err(|err| format!("failed to serialize audit event: {err}"))?;
127 append_jsonl_to_internal_branch(repo_root, &self.branch, &format!("{json}\n"))
128 .map_err(|err| err.to_string())
129 }
130
131 fn read_from_branch(&self) -> Result<InternalBranchRead, String> {
132 let Some(repo_root) = self.repo_root() else {
133 return Err("unable to resolve project root for internal audit branch".to_string());
134 };
135 let branch_read =
136 read_internal_branch_log(repo_root, &self.branch).map_err(|err| err.to_string())?;
137 Ok(match branch_read {
138 InternalBranchLogRead::BranchMissing => InternalBranchRead::BranchMissing,
139 InternalBranchLogRead::LogMissing => InternalBranchRead::LogMissing,
140 InternalBranchLogRead::Contents(contents) => {
141 InternalBranchRead::Events(parse_events_from_jsonl(&contents))
142 }
143 })
144 }
145
146 fn append_to_fallback(&self, event: &AuditEvent) {
147 if let Err(err) = append_event_to_file(&self.fallback_path, event) {
148 tracing::warn!("fallback audit write failed: {err}");
149 }
150 }
151
152 fn read_fallback_events(&self) -> Vec<AuditEvent> {
153 read_events_from_path(&self.fallback_path)
154 }
155
156 fn replay_fallback_into_branch(&self) -> Result<(), String> {
157 let Ok(contents) = std::fs::read_to_string(&self.fallback_path) else {
158 return Ok(());
159 };
160 if contents.trim().is_empty() {
161 return Ok(());
162 }
163
164 let Some(repo_root) = self.repo_root() else {
165 return Err("unable to resolve project root for fallback audit replay".to_string());
166 };
167 append_jsonl_to_internal_branch(repo_root, &self.branch, &contents)
168 .map_err(|err| err.to_string())?;
169 remove_file_if_present(&self.fallback_path)
170 .map_err(|err| format!("failed to remove fallback audit log: {err}"))?;
171 Ok(())
172 }
173
174 fn merged_events_with_fallback(&self, branch_events: Vec<AuditEvent>) -> Vec<AuditEvent> {
175 let fallback_events = self.read_fallback_events();
176 if fallback_events.is_empty() {
177 return branch_events;
178 }
179
180 if let Err(err) = self.replay_fallback_into_branch() {
181 tracing::warn!("fallback audit replay failed: {err}");
182 }
183
184 merge_events(branch_events, fallback_events)
185 }
186
187 fn migrate_legacy_worktree_log(&self) {
188 let legacy_path = audit_log_path(&self.ito_path);
189 let Ok(contents) = std::fs::read_to_string(&legacy_path) else {
190 return;
191 };
192 if contents.trim().is_empty() {
193 return;
194 }
195
196 if let Some(repo_root) = self.repo_root() {
197 match append_jsonl_to_internal_branch(repo_root, &self.branch, &contents) {
198 Ok(()) => {
199 if let Err(err) = remove_file_if_present(&legacy_path) {
200 tracing::warn!("failed to remove migrated legacy audit log: {err}");
201 }
202 return;
203 }
204 Err(err) => {
205 tracing::warn!("legacy tracked audit log import failed: {err}");
206 eprintln!(
207 "Warning: durable internal audit storage unavailable; migrating legacy tracked audit log into local fallback store '{}': {err}",
208 self.fallback_path.display()
209 );
210 }
211 }
212 }
213
214 if let Err(err) = merge_jsonl_file(&self.fallback_path, &contents) {
215 tracing::warn!("legacy tracked audit fallback import failed: {err}");
216 return;
217 }
218
219 if let Err(err) = remove_file_if_present(&legacy_path) {
220 tracing::warn!("failed to remove migrated legacy audit log: {err}");
221 }
222 }
223
224 fn ensure_legacy_worktree_log_migrated(&self) {
225 self.legacy_migration_done.get_or_init(|| {
226 self.migrate_legacy_worktree_log();
227 });
228 }
229
230 fn warn_and_fallback(&self, err: &str, event: &AuditEvent) {
231 tracing::warn!("internal audit branch unavailable: {err}");
232 eprintln!(
233 "Warning: durable internal audit storage unavailable; using local fallback store '{}': {err}",
234 self.fallback_path.display()
235 );
236 self.append_to_fallback(event);
237 }
238}
239
240impl AuditWriter for LocalAuditStore {
241 fn append(&self, event: &AuditEvent) -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
242 self.ensure_legacy_worktree_log_migrated();
243 if let Err(err) = self.append_to_branch(event) {
244 self.warn_and_fallback(&err, event);
245 }
246 Ok(())
247 }
248}
249
250impl AuditEventStore for LocalAuditStore {
251 fn read_all(&self) -> Vec<AuditEvent> {
252 self.ensure_legacy_worktree_log_migrated();
253 match self.read_from_branch() {
254 Ok(InternalBranchRead::Events(events)) => self.merged_events_with_fallback(events),
255 Ok(InternalBranchRead::BranchMissing) => self.read_fallback_events(),
256 Ok(InternalBranchRead::LogMissing) => {
257 tracing::warn!(
258 branch = %self.branch,
259 "internal audit branch exists but has no audit log yet; treating as empty history"
260 );
261 self.merged_events_with_fallback(Vec::new())
262 }
263 Err(err) => {
264 tracing::warn!("internal audit branch read failed: {err}");
265 self.read_fallback_events()
266 }
267 }
268 }
269
270 fn location(&self) -> AuditStorageLocation {
271 if self.repo_root().is_some() {
272 AuditStorageLocation::Other(format!("internal-branch:{}", self.branch))
273 } else {
274 AuditStorageLocation::Filesystem(self.fallback_path.clone())
275 }
276 }
277}
278
279enum InternalBranchRead {
280 BranchMissing,
281 LogMissing,
282 Events(Vec<AuditEvent>),
283}
284
285pub fn default_audit_store(ito_path: &Path) -> Box<dyn AuditEventStore> {
290 let ctx = ConfigContext::from_process_env();
291 #[cfg(feature = "backend")]
292 if let Ok(runtime) = resolve_repository_runtime(ito_path, &ctx)
293 && runtime.mode() == PersistenceMode::Remote
294 && let Some(backend_runtime) = runtime.backend_runtime().cloned()
295 {
296 return Box::new(BackendAuditStore::new(BackendHttpClient::new(
297 backend_runtime,
298 )));
299 }
300
301 let branch = resolve_internal_audit_branch(ito_path, &ctx);
302 let fallback_path = fallback_audit_log_path(ito_path);
303 Box::new(LocalAuditStore::new(ito_path, branch, fallback_path))
304}
305
306fn resolve_internal_audit_branch(ito_path: &Path, ctx: &ConfigContext) -> String {
307 let Some(project_root) = ito_path.parent() else {
308 return "ito/internal/audit".to_string();
309 };
310 let resolved = load_cascading_project_config(project_root, ito_path, ctx);
311 let (_, branch) = resolve_audit_mirror_settings(&resolved.merged);
312 branch
313}
314
315fn fallback_audit_log_path(ito_path: &Path) -> PathBuf {
316 let runner = SystemProcessRunner;
317 if let Some(project_root) = ito_path.parent()
318 && let Some(git_dir) = git_dir_path(&runner, project_root)
319 {
320 return git_dir.join("ito").join("audit").join("events.jsonl");
321 }
322
323 ito_path
324 .join(".state-local")
325 .join("audit")
326 .join("events.jsonl")
327}
328
329fn git_dir_path(runner: &dyn ProcessRunner, project_root: &Path) -> Option<PathBuf> {
330 let out = runner
331 .run(
332 &ProcessRequest::new("git")
333 .args(["rev-parse", "--absolute-git-dir"])
334 .current_dir(project_root),
335 )
336 .ok()?;
337 if !out.success {
338 return None;
339 }
340
341 let path = out.stdout.trim();
342 if path.is_empty() {
343 return None;
344 }
345 Some(PathBuf::from(path))
346}
347
348fn merge_jsonl_file(path: &Path, incoming: &str) -> std::io::Result<()> {
349 let existing = std::fs::read_to_string(path).unwrap_or_default();
350 let merged = merge_jsonl_contents(&existing, incoming);
351 if merged == existing {
352 return Ok(());
353 }
354
355 if let Some(parent) = path.parent() {
356 std::fs::create_dir_all(parent)?;
357 }
358 std::fs::write(path, merged)
359}
360
361fn merge_jsonl_contents(existing: &str, incoming: &str) -> String {
362 let mut merged = Vec::new();
363 let mut seen = HashSet::new();
364
365 for line in existing.lines().chain(incoming.lines()) {
366 let line = line.trim();
367 if line.is_empty() || !seen.insert(line.to_string()) {
368 continue;
369 }
370 merged.push(line.to_string());
371 }
372
373 if merged.is_empty() {
374 String::new()
375 } else {
376 format!("{}\n", merged.join("\n"))
377 }
378}
379
380fn remove_file_if_present(path: &Path) -> std::io::Result<()> {
381 match std::fs::remove_file(path) {
382 Ok(()) => Ok(()),
383 Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(()),
384 Err(err) => Err(err),
385 }
386}
387
388fn merge_events(primary: Vec<AuditEvent>, secondary: Vec<AuditEvent>) -> Vec<AuditEvent> {
389 let mut merged = Vec::new();
390 let mut seen = HashSet::new();
391
392 for event in primary.into_iter().chain(secondary) {
393 let Ok(key) = serde_json::to_string(&event) else {
394 merged.push(event);
395 continue;
396 };
397 if seen.insert(key) {
398 merged.push(event);
399 }
400 }
401
402 merged
403}
404
405#[cfg(test)]
406#[path = "store_tests.rs"]
407mod store_tests;