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