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