use std::collections::HashMap;
use std::path::PathBuf;
use std::sync::Mutex;
use rayon::prelude::*;
use crate::{
ActivityOptions, Phase, Progress, RpoError, SkippedFile, SnapshotSelector, StreamStats,
backend::{Commit, CommitId, GitBackend, Signature, WalkOptions},
filters::FilterSet,
frames::{BlameRecord, extension_of},
sinks::FrameSink,
};
pub(crate) struct Snapshot {
id: CommitId,
label: Option<String>,
time_ms: i64,
}
pub(crate) fn resolve_snapshots<B: GitBackend>(
backend: &B,
selector: &SnapshotSelector,
activity: &ActivityOptions,
) -> Result<Vec<Snapshot>, RpoError> {
match selector {
SnapshotSelector::Head => {
let id = backend.head_commit()?;
let meta = single_commit_meta(backend, &id)?;
Ok(vec![Snapshot {
id,
label: Some("HEAD".to_string()),
time_ms: meta.committer.time_ms,
}])
}
SnapshotSelector::AtRevs(revs) => {
let mut out = Vec::with_capacity(revs.len());
for rev in revs {
let id = backend.resolve_rev(rev)?;
let meta = single_commit_meta(backend, &id)?;
out.push(Snapshot {
id,
label: Some(rev.clone()),
time_ms: meta.committer.time_ms,
});
}
Ok(out)
}
SnapshotSelector::Tags => {
let mut out = Vec::new();
for (name, id) in backend.tags()? {
let meta = single_commit_meta(backend, &id)?;
out.push(Snapshot {
id,
label: Some(name),
time_ms: meta.committer.time_ms,
});
}
ensure_head_tagged(backend, &mut out)?;
Ok(out)
}
SnapshotSelector::EveryNCommits(_)
| SnapshotSelector::Daily
| SnapshotSelector::Weekly
| SnapshotSelector::Monthly
| SnapshotSelector::AllCommits => {
let opts = WalkOptions {
first_parent_only: true,
include_merges: !activity.ignore_merges || activity.first_parent_only,
};
let commits: Vec<Commit> = backend.iter_commits(opts).collect::<Result<Vec<_>, _>>()?;
let mut out = select_commits(commits, selector);
ensure_head_tagged(backend, &mut out)?;
Ok(out)
}
}
}
fn ensure_head_tagged<B: GitBackend>(
backend: &B,
snapshots: &mut Vec<Snapshot>,
) -> Result<(), RpoError> {
let head_id = backend.head_commit()?;
if let Some(existing) = snapshots.iter_mut().find(|s| s.id == head_id) {
existing.label = Some(match existing.label.take() {
Some(l) if !l.is_empty() => format!("{l}, HEAD"),
_ => "HEAD".to_string(),
});
return Ok(());
}
let meta = single_commit_meta(backend, &head_id)?;
snapshots.insert(
0,
Snapshot {
id: head_id,
label: Some("HEAD".to_string()),
time_ms: meta.committer.time_ms,
},
);
Ok(())
}
struct CommitMeta {
committer: Signature,
}
fn single_commit_meta<B: GitBackend>(backend: &B, id: &CommitId) -> Result<CommitMeta, RpoError> {
let c = backend.commit_meta(id)?;
Ok(CommitMeta {
committer: c.committer,
})
}
fn select_commits(commits: Vec<Commit>, selector: &SnapshotSelector) -> Vec<Snapshot> {
match selector {
SnapshotSelector::AllCommits => commits
.into_iter()
.map(|c| Snapshot {
time_ms: c.committer.time_ms,
id: c.id,
label: None,
})
.collect(),
SnapshotSelector::EveryNCommits(n) => {
let step = (*n).max(1);
commits
.into_iter()
.enumerate()
.filter(|(i, _)| i % step == 0)
.map(|(_, c)| Snapshot {
time_ms: c.committer.time_ms,
id: c.id,
label: None,
})
.collect()
}
SnapshotSelector::Daily => bucketed_snapshots(commits, day_bucket_utc),
SnapshotSelector::Weekly => bucketed_snapshots(commits, iso_week_bucket_utc),
SnapshotSelector::Monthly => bucketed_snapshots(commits, month_bucket_utc),
_ => Vec::new(), }
}
fn bucketed_snapshots<K: Eq + std::hash::Hash>(
commits: Vec<Commit>,
bucket: impl Fn(i64) -> (K, String),
) -> Vec<Snapshot> {
let mut seen: HashMap<K, ()> = HashMap::new();
let mut out = Vec::new();
for c in commits {
let (key, label) = bucket(c.committer.time_ms);
if seen.insert(key, ()).is_none() {
out.push(Snapshot {
time_ms: c.committer.time_ms,
id: c.id,
label: Some(label),
});
}
}
out
}
fn zoned_utc(unix_ms: i64) -> jiff::civil::Date {
let ts = jiff::Timestamp::from_millisecond(unix_ms).unwrap_or(jiff::Timestamp::UNIX_EPOCH);
ts.to_zoned(jiff::tz::TimeZone::UTC).date()
}
fn day_bucket_utc(unix_ms: i64) -> ((i16, u16), String) {
let d = zoned_utc(unix_ms);
let y = d.year();
let doy = d.day_of_year();
let label = format!("{y:04}-{:02}-{:02}", d.month(), d.day());
((y, doy as u16), label)
}
fn iso_week_bucket_utc(unix_ms: i64) -> ((i16, i8), String) {
let d = zoned_utc(unix_ms);
let iso = d.iso_week_date();
let label = format!("{:04}-W{:02}", iso.year(), iso.week());
((iso.year(), iso.week()), label)
}
fn month_bucket_utc(unix_ms: i64) -> ((i16, i8), String) {
let d = zoned_utc(unix_ms);
let y = d.year();
let m = d.month();
let label = format!("{y:04}-{m:02}");
((y, m), label)
}
pub struct BlameRun {
pub records: Vec<BlameRecord>,
pub skipped: Vec<SkippedFile>,
}
pub fn run_in_memory<B: GitBackend>(
backend: &B,
filters: &FilterSet,
selector: &SnapshotSelector,
activity: &ActivityOptions,
thread_count: Option<usize>,
progress: Option<&(dyn Fn(Progress) + Send + Sync)>,
commit_meta_by_sha: &HashMap<String, CommitMetaBlame>,
) -> Result<BlameRun, RpoError> {
let snapshots = resolve_snapshots(backend, selector, activity)?;
let total = snapshots.len();
tracing::info!(snapshots = total, "blame: start");
let started_all = std::time::Instant::now();
let mut records: Vec<BlameRecord> = Vec::new();
let mut skipped: Vec<SkippedFile> = Vec::new();
for (i, snap) in snapshots.into_iter().enumerate() {
let snap_label = snap
.label
.clone()
.unwrap_or_else(|| snap.id.short_hex().to_string());
tracing::debug!(
snapshot = %snap_label,
index = i + 1,
total,
"blame: snapshot start"
);
let started_snap = std::time::Instant::now();
let (snap_records, snap_skipped) =
blame_single_snapshot(backend, &snap, filters, thread_count, progress)?;
tracing::debug!(
snapshot = %snap_label,
rows = snap_records.len(),
skipped = snap_skipped.len(),
elapsed_ms = started_snap.elapsed().as_millis() as u64,
"blame: snapshot done"
);
for unit in snap_records {
let sha_hex = unit.path_hunk.commit_id.to_hex();
let meta = commit_meta_by_sha
.get(&sha_hex)
.cloned()
.unwrap_or_default();
records.push(BlameRecord {
snapshot_sha: snap.id.to_hex(),
snapshot_time_ms: snap.time_ms,
snapshot_label: snap.label.clone(),
path: unit.path.to_string_lossy().to_string(),
start_line: unit.path_hunk.start_line,
line_count: unit.path_hunk.line_count,
commit_sha: sha_hex,
canonical_author_name: meta.canonical_author_name.clone(),
canonical_author_email: meta.canonical_author_email.clone(),
canonical_committer_name: meta.canonical_committer_name.clone(),
canonical_committer_email: meta.canonical_committer_email.clone(),
commit_time_ms: meta.commit_time_ms,
extension: extension_of(&unit.path),
});
}
skipped.extend(snap_skipped);
}
tracing::info!(
snapshots = total,
rows = records.len(),
skipped = skipped.len(),
elapsed_ms = started_all.elapsed().as_millis() as u64,
"blame: done"
);
Ok(BlameRun { records, skipped })
}
#[allow(clippy::too_many_arguments)]
pub fn run_streaming<B: GitBackend, S: FrameSink + ?Sized>(
backend: &B,
filters: &FilterSet,
selector: &SnapshotSelector,
activity: &ActivityOptions,
thread_count: Option<usize>,
progress: Option<&(dyn Fn(Progress) + Send + Sync)>,
commit_meta_by_sha: &HashMap<String, CommitMetaBlame>,
sink: &mut S,
) -> Result<StreamStats, RpoError> {
let snapshots = resolve_snapshots(backend, selector, activity)?;
let mut stats = StreamStats {
snapshots_written: 0,
rows_written: 0,
bytes_written: 0,
skipped_files: Vec::new(),
};
for snap in snapshots {
let (snap_units, snap_skipped) =
blame_single_snapshot(backend, &snap, filters, thread_count, progress)?;
let mut recs = Vec::with_capacity(snap_units.len());
for unit in snap_units {
let sha_hex = unit.path_hunk.commit_id.to_hex();
let meta = commit_meta_by_sha
.get(&sha_hex)
.cloned()
.unwrap_or_default();
recs.push(BlameRecord {
snapshot_sha: snap.id.to_hex(),
snapshot_time_ms: snap.time_ms,
snapshot_label: snap.label.clone(),
path: unit.path.to_string_lossy().to_string(),
start_line: unit.path_hunk.start_line,
line_count: unit.path_hunk.line_count,
commit_sha: sha_hex,
canonical_author_name: meta.canonical_author_name,
canonical_author_email: meta.canonical_author_email,
canonical_committer_name: meta.canonical_committer_name,
canonical_committer_email: meta.canonical_committer_email,
commit_time_ms: meta.commit_time_ms,
extension: extension_of(&unit.path),
});
}
let row_count = recs.len() as u64;
let frame = crate::frames::blame::build(recs)?;
sink.write_snapshot(frame)?;
stats.snapshots_written += 1;
stats.rows_written += row_count;
stats.skipped_files.extend(snap_skipped);
}
sink.finish()?;
Ok(stats)
}
#[derive(Clone, Debug, Default)]
pub struct CommitMetaBlame {
pub canonical_author_name: String,
pub canonical_author_email: String,
pub canonical_committer_name: String,
pub canonical_committer_email: String,
pub commit_time_ms: i64,
}
struct SnapshotUnit {
path: PathBuf,
path_hunk: crate::backend::BlameHunk,
}
fn blame_single_snapshot<B: GitBackend>(
backend: &B,
snap: &Snapshot,
filters: &FilterSet,
thread_count: Option<usize>,
progress: Option<&(dyn Fn(Progress) + Send + Sync)>,
) -> Result<(Vec<SnapshotUnit>, Vec<SkippedFile>), RpoError> {
let all_paths = backend.list_tree_paths(&snap.id)?;
let paths: Vec<PathBuf> = all_paths
.into_iter()
.filter(|p| filters.blame_includes(p))
.collect();
let total = paths.len() as u64;
if let Some(cb) = progress {
cb(Progress {
phase: Phase::Blaming {
snapshot: snap.id.to_hex(),
},
completed: 0,
total,
});
}
let pool = if let Some(n) = thread_count {
Some(
rayon::ThreadPoolBuilder::new()
.num_threads(n)
.build()
.map_err(|e| RpoError::Backend(e.to_string()))?,
)
} else {
None
};
let completed = Mutex::new(0u64);
let skipped = Mutex::new(Vec::<SkippedFile>::new());
let units_mx = Mutex::new(Vec::<SnapshotUnit>::new());
let work = |path: PathBuf| -> Result<(), RpoError> {
let h = backend.thread_handle()?;
let result = h.blame_file(&snap.id, &path);
let bump = {
let mut n = completed.lock().unwrap();
*n += 1;
*n
};
match result {
Ok(hunks) => {
let mut units = units_mx.lock().unwrap();
for h in hunks {
units.push(SnapshotUnit {
path: path.clone(),
path_hunk: h,
});
}
}
Err(e) => {
let mut sk = skipped.lock().unwrap();
sk.push(SkippedFile {
snapshot_sha: snap.id.to_hex(),
path: path.clone(),
reason: e.to_string(),
});
}
}
if let Some(cb) = progress {
cb(Progress {
phase: Phase::Blaming {
snapshot: snap.id.to_hex(),
},
completed: bump,
total,
});
}
Ok(())
};
match pool {
Some(p) => p.install(|| paths.into_par_iter().try_for_each(work))?,
None => paths.into_par_iter().try_for_each(work)?,
};
Ok((
units_mx.into_inner().unwrap(),
skipped.into_inner().unwrap(),
))
}
#[cfg(test)]
mod tests {
use super::*;
fn ts(iso: &str) -> i64 {
let ts: jiff::Timestamp = iso.parse().unwrap();
ts.as_millisecond()
}
#[test]
fn day_bucket_labels_and_keys() {
let (key, label) = day_bucket_utc(ts("2025-03-05T12:34:56Z"));
assert_eq!(label, "2025-03-05");
assert_eq!(key.0, 2025);
assert_eq!(key.1, 64);
}
#[test]
fn day_bucket_spans_midnight() {
let a = day_bucket_utc(ts("2025-03-05T23:59:59Z"));
let b = day_bucket_utc(ts("2025-03-06T00:00:00Z"));
assert_ne!(a.0, b.0);
assert_eq!(a.1, "2025-03-05");
assert_eq!(b.1, "2025-03-06");
}
#[test]
fn iso_week_crosses_year_boundary() {
let (key, label) = iso_week_bucket_utc(ts("2024-12-30T12:00:00Z"));
assert_eq!(label, "2025-W01");
assert_eq!(key, (2025, 1));
let (key2, label2) = iso_week_bucket_utc(ts("2023-01-01T12:00:00Z"));
assert_eq!(label2, "2022-W52");
assert_eq!(key2, (2022, 52));
}
#[test]
fn month_bucket_matches_prior_behavior() {
let (key, label) = month_bucket_utc(ts("2020-02-29T08:00:00Z"));
assert_eq!(key, (2020, 2));
assert_eq!(label, "2020-02");
}
#[test]
fn bucketed_keeps_first_seen_per_bucket() {
let commits = vec![
mk_commit(0x01, ts("2025-03-05T18:00:00Z")),
mk_commit(0x02, ts("2025-03-05T08:00:00Z")),
mk_commit(0x03, ts("2025-03-06T12:00:00Z")),
];
let out = bucketed_snapshots(commits, day_bucket_utc);
assert_eq!(out.len(), 2);
assert_eq!(out[0].id.0[0], 0x01);
assert_eq!(out[0].label.as_deref(), Some("2025-03-05"));
assert_eq!(out[1].id.0[0], 0x03);
assert_eq!(out[1].label.as_deref(), Some("2025-03-06"));
}
fn mk_commit(first_byte: u8, time_ms: i64) -> Commit {
let mut bytes = [0u8; 20];
bytes[0] = first_byte;
Commit {
id: CommitId(bytes),
author: Signature {
name: String::new(),
email: String::new(),
time_ms,
},
committer: Signature {
name: String::new(),
email: String::new(),
time_ms,
},
parent_ids: vec![],
message_subject: String::new(),
}
}
}