Skip to main content

ito_core/audit/
store.rs

1//! Audit event storage abstractions.
2
3use 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/// Storage location descriptor for audit events.
30#[derive(Debug, Clone, PartialEq, Eq)]
31pub enum AuditStorageLocation {
32    /// Filesystem-backed storage at the given path.
33    Filesystem(PathBuf),
34    /// Non-filesystem or abstract storage identified by a short label.
35    Other(String),
36}
37
38/// Combined read/write abstraction for audit event storage.
39pub trait AuditEventStore: AuditWriter + Send + Sync {
40    /// Read all available events from the underlying storage.
41    fn read_all(&self) -> Vec<AuditEvent>;
42
43    /// Describe the underlying storage location for diagnostics and routing.
44    fn location(&self) -> AuditStorageLocation;
45}
46
47/// Build a stable deduplication key for an audit storage location.
48pub 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
285/// Resolve the default audit store for the current project.
286///
287/// Today this is filesystem-backed. Later tasks can route this to internal-branch
288/// or backend-managed storage without changing reader call sites again.
289pub 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;