use std::sync::Arc;
use std::time::Instant;
use serde::Deserialize;
use tracing::{debug, warn};
use trusty_common::credentials::{default_store, KeyStore};
use trusty_common::inference::{
register_default_factories, ChatMessage, ChatRequest, Configurator, InferenceAdapter,
InferenceError,
};
use super::types::{Effort, Finding, LongitudinalFinding, PeriodBatch, TokenCostSummary};
pub const PERIOD_REVIEWER_TEMPERATURE: f32 = 0.2;
pub const PERIOD_REVIEWER_MAX_TOKENS: u32 = 2048;
const MAX_DIFFS_IN_PROMPT: usize = 10;
pub fn period_findings_schema() -> serde_json::Value {
serde_json::json!({
"type": "object",
"properties": {
"findings": {
"type": "array",
"items": {
"type": "object",
"properties": {
"kind": {"type": "string"},
"description": {"type": "string"},
"suggestion": {"type": "string"},
"confidence": {"type": "number", "minimum": 0.0, "maximum": 1.0},
"file": {"type": "string"},
"severity": {
"type": "string",
"enum": ["low", "medium", "high", "critical"]
}
},
"required": ["kind", "description"]
}
}
},
"required": ["findings"]
})
}
pub fn period_reviewer_system_prompt() -> &'static str {
r#"You are a senior software engineer reviewing a sample of one engineer's commits
over a specific time window as part of a longitudinal quality analysis.
## Task
Identify code-quality findings present in the sampled diffs. Focus on:
- Correctness bugs, error-handling gaps, resource leaks
- Security weaknesses (injection, auth, secrets in code)
- Logic errors, off-by-one issues, data-loss risks
- Missing tests or test-quality issues
- Recurring anti-patterns visible across multiple commits in this window
## Output (REQUIRED)
Populate the structured response with a `findings` array.
Each finding must include:
- `kind`: short category label (e.g. error_handling, security, logic)
- `description`: concise description of the issue observed
- `suggestion`: concrete improvement suggestion
- `confidence`: float in [0.0, 1.0]
- `file`: most relevant file path; use "multiple" if the issue spans files
- `severity`: one of low, medium, high, critical
`findings` may be an empty array if the sample looks clean."#
}
pub fn build_period_user_message(batch: &PeriodBatch) -> String {
let s = &batch.stats;
let mut msg = String::with_capacity(4096);
msg.push_str(&format!(
"## Period: {}\nFrom {} to {}\n\n",
s.period_label, s.since, s.until
));
msg.push_str("### Statistics\n");
msg.push_str(&format!("- Commits: {}\n", s.commit_count));
msg.push_str(&format!("- Quality score: {:.2}\n", s.quality_score));
msg.push_str(&format!("- Ticketed %: {:.0}%\n", s.ticketed_pct * 100.0));
if !s.categories.is_empty() {
let mut cats: Vec<(&String, &u64)> = s.categories.iter().collect();
cats.sort_by_key(|(k, _)| k.as_str());
let cat_str: Vec<String> = cats.iter().map(|(k, v)| format!("{k}={v}")).collect();
msg.push_str(&format!("- Categories: {}\n", cat_str.join(", ")));
}
if !s.repositories.is_empty() {
msg.push_str(&format!("- Repositories: {}\n", s.repositories.join(", ")));
}
msg.push('\n');
if batch.sampled_diffs.is_empty() {
msg.push_str("### Sampled diffs\n*(no diffs available for this period)*\n\n");
} else {
msg.push_str("### Sampled diffs\n\n");
for (i, diff) in batch
.sampled_diffs
.iter()
.enumerate()
.take(MAX_DIFFS_IN_PROMPT)
{
let cat = diff.category.as_deref().unwrap_or("unknown");
let effort = diff.effort.as_deref().unwrap_or("?");
msg.push_str(&format!(
"#### Diff {} — {} ({repo}) [category={cat}, effort={effort}]\n",
i + 1,
&diff.sha[..8.min(diff.sha.len())],
repo = diff.repository,
));
msg.push_str(&format!("Commit: {}\n\n", diff.message));
msg.push_str("```diff\n");
msg.push_str(&diff.diff_text);
if !diff.diff_text.ends_with('\n') {
msg.push('\n');
}
msg.push_str("```\n\n");
}
}
msg.push_str(
"Please review the diffs above and populate the structured `findings` \
array as specified in the system prompt.\n",
);
msg
}
pub fn build_period_request(batch: &PeriodBatch, model: &str) -> ChatRequest {
let schema = serde_json::to_string_pretty(&period_findings_schema()).unwrap_or_default();
let system = format!(
"{}\n\n## Response schema\nReturn ONLY a JSON object conforming to this schema:\n\
```json\n{schema}\n```",
period_reviewer_system_prompt()
);
let mut req = ChatRequest::new(
model,
vec![
ChatMessage::system(system),
ChatMessage::user(build_period_user_message(batch)),
],
);
req.temperature = Some(PERIOD_REVIEWER_TEMPERATURE);
req.max_tokens = Some(PERIOD_REVIEWER_MAX_TOKENS);
req
}
#[derive(Debug, Default)]
#[non_exhaustive]
pub struct PeriodReview {
pub findings: Vec<LongitudinalFinding>,
pub skipped: Option<InferenceError>,
}
impl PeriodReview {
pub fn was_skipped(&self) -> bool {
self.skipped.is_some()
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SkippedPeriod {
pub period_label: String,
pub reason: String,
}
#[derive(Debug, Default, Clone)]
#[non_exhaustive]
pub struct PeriodRunSummary {
pub reviewed: usize,
pub skipped: Vec<SkippedPeriod>,
}
impl PeriodRunSummary {
pub fn record(&mut self, period_label: &str, review: PeriodReview) -> Vec<LongitudinalFinding> {
match review.skipped {
Some(err) => {
self.skipped.push(SkippedPeriod {
period_label: period_label.to_string(),
reason: err.to_string(),
});
Vec::new()
}
None => {
self.reviewed += 1;
review.findings
}
}
}
pub fn attempted(&self) -> usize {
self.reviewed + self.skipped.len()
}
pub fn is_complete(&self) -> bool {
self.skipped.is_empty()
}
pub fn coverage_line(&self) -> String {
let attempted = self.attempted();
if self.is_complete() {
return format!("{}/{attempted} period(s) reviewed", self.reviewed);
}
let labels: Vec<&str> = self
.skipped
.iter()
.map(|s| s.period_label.as_str())
.collect();
format!(
"{}/{attempted} period(s) reviewed — {} SKIPPED (provider failure): {}",
self.reviewed,
self.skipped.len(),
labels.join(", ")
)
}
pub fn coverage_note(&self) -> Option<String> {
if self.is_complete() {
return None;
}
let mut note = String::from("## Coverage\n\n");
note.push_str(&format!(
"{} of {} period(s) were reviewed. The following period(s) were \
**skipped** because the inference provider call failed — they are \
absent from the findings and the trajectory, and are NOT evidence \
of clean work:\n\n",
self.reviewed,
self.attempted()
));
for s in &self.skipped {
note.push_str(&format!("- `{}` — {}\n", s.period_label, s.reason));
}
note.push('\n');
Some(note)
}
}
pub struct PeriodReviewer {
adapter: Arc<dyn InferenceAdapter>,
model: String,
}
impl PeriodReviewer {
pub fn from_slug(model: &str) -> super::Result<Self> {
Self::from_slug_with_store(model, default_store().as_ref())
}
pub fn from_slug_with_store(model: &str, store: &dyn KeyStore) -> super::Result<Self> {
let mut configurator = Configurator::new();
register_default_factories(&mut configurator);
#[cfg(feature = "bedrock")]
trusty_common::inference::register_bedrock_factory(&mut configurator);
let adapter = configurator.build(model, store)?;
Ok(Self {
adapter: Arc::from(adapter),
model: model.to_string(),
})
}
pub fn with_adapter(adapter: Arc<dyn InferenceAdapter>, model: impl Into<String>) -> Self {
Self {
adapter,
model: model.into(),
}
}
pub async fn review_period(
&self,
batch: &PeriodBatch,
cost_out: &mut TokenCostSummary,
) -> PeriodReview {
let period = &batch.stats.period_label;
let request = build_period_request(batch, &self.model);
let start = Instant::now();
let response = match self.adapter.chat(&request).await {
Ok(r) => r,
Err(e) => {
warn!(
period = %period,
model = %self.model,
error = %e,
"batch_reviewer: inference call failed — period skipped"
);
return PeriodReview {
findings: Vec::new(),
skipped: Some(e),
};
}
};
let latency_ms = start.elapsed().as_millis() as u64;
let usage = response.usage();
cost_out.accumulate(
u64::from(usage.prompt_tokens),
u64::from(usage.completion_tokens),
usage.cost_usd.unwrap_or(0.0),
latency_ms,
);
debug!(
period = %period,
model = %response.resolved_model(&self.model),
input_tokens = usage.prompt_tokens,
output_tokens = usage.completion_tokens,
latency_ms,
"batch_reviewer: inference call complete"
);
PeriodReview {
findings: parse_period_findings(&response.first_text().unwrap_or_default(), period),
skipped: None,
}
}
}
#[derive(Debug, Deserialize)]
struct PeriodFindingsBlock {
#[serde(default)]
findings: Vec<PeriodFindingWire>,
}
#[derive(Debug, Deserialize)]
struct PeriodFindingWire {
#[serde(default)]
kind: String,
#[serde(default)]
description: String,
#[serde(default)]
suggestion: String,
#[serde(default)]
confidence: f32,
#[serde(default)]
file: String,
#[serde(default)]
severity: String,
}
pub fn parse_period_findings(body: &str, period_label: &str) -> Vec<LongitudinalFinding> {
let body = body.trim();
if body.is_empty() {
warn!(period = %period_label, "batch_reviewer: empty response — returning empty findings");
return Vec::new();
}
if body.starts_with('{') {
if let Ok(block) = serde_json::from_str::<PeriodFindingsBlock>(body) {
debug!(
period = %period_label,
findings = block.findings.len(),
"batch_reviewer: parsed via direct JSON"
);
return convert_period_block(block, period_label);
}
}
let Some(fence_start) = body.rfind("```json") else {
warn!(
period = %period_label,
"batch_reviewer: no JSON block in response — returning empty findings"
);
return Vec::new();
};
let after = &body[fence_start + 7..];
let Some(fence_end) = after.find("```") else {
warn!(period = %period_label, "batch_reviewer: unclosed JSON block — returning empty findings");
return Vec::new();
};
match serde_json::from_str::<PeriodFindingsBlock>(after[..fence_end].trim()) {
Ok(block) => convert_period_block(block, period_label),
Err(e) => {
warn!(
period = %period_label,
error = %e,
"batch_reviewer: JSON parse error — returning empty findings"
);
Vec::new()
}
}
}
fn convert_period_block(
block: PeriodFindingsBlock,
period_label: &str,
) -> Vec<LongitudinalFinding> {
block
.findings
.into_iter()
.map(|f| {
let effort = severity_to_effort(&f.severity);
let file = if f.file.is_empty() {
"unknown".to_string()
} else {
f.file
};
let kind = if f.kind.is_empty() {
"general".to_string()
} else {
f.kind
};
LongitudinalFinding {
period_label: period_label.to_string(),
finding: Finding::new(
file,
kind,
f.description,
f.suggestion,
f.confidence,
effort,
),
trend_tag: None,
}
})
.collect()
}
pub fn severity_to_effort(severity: &str) -> Effort {
match severity.to_lowercase().as_str() {
"high" | "critical" => Effort::High,
"medium" => Effort::Medium,
_ => Effort::Low,
}
}
#[cfg(test)]
#[path = "batch_reviewer_tests.rs"]
mod tests;