Files
leadcast/backend/src/api.rs
T
2026-07-26 21:17:57 -03:00

1963 lines
64 KiB
Rust

use std::time::Duration;
use actix_governor::{Governor, GovernorConfig, PeerIpKeyExtractor, governor::middleware::NoOpMiddleware};
use actix_web::{HttpMessage, HttpRequest, HttpResponse, http::header::CONTENT_DISPOSITION, web};
use chrono::{DateTime, Utc};
use futures::StreamExt;
use serde::Deserialize;
use serde_json::{Value, json};
use sha2::{Digest, Sha256};
use url::Url;
use uuid::Uuid;
use crate::{
api_queries,
api_response::DataResponse,
auth::{AUTH_FAILURE_MESSAGE, LOGIN_JSON_LIMIT_BYTES, SESSION_COOKIE_NAME, SessionIdentity},
best_contacts, contact_normalizer,
db::{
models::{
Contact, Interviewee, IntervieweePatch, NewAppearance, NewAuditEvent, NewCategory,
NewContact, NewContactEvidenceDraft, NewInterviewee, NewJob, NewMediaAsset, NewOrigin,
NewPipelineRun, NewPodcastChannel,
},
querys::{
appearance, audit_event, category, contact, contact_candidate, interviewee,
interviewee_candidate, job, maintenance, media_asset, origin, pipeline_run,
podcast_channel,
},
},
error::{AppError, AppResult},
export,
identity_resolution::normalize_name,
logs,
media::MediaKind,
request_id,
state::AppState,
views::PodcastDiscoveryView,
};
const MODULE: &str = "api";
/// Quantidade máxima de ids aceita em uma ação em massa. Protege o backend de
/// payloads desproporcionais vindos do frontend (ex.: seleção "marcar tudo").
const MAX_BULK_IDS: usize = 500;
pub fn configure(
config: &mut web::ServiceConfig,
login_rate_limit: &GovernorConfig<PeerIpKeyExtractor, NoOpMiddleware>,
) {
config.service(
web::scope("/api")
.service(
web::scope("/auth")
.service(
web::resource("/login")
.wrap(Governor::new(login_rate_limit))
.app_data(web::JsonConfig::default().limit(LOGIN_JSON_LIMIT_BYTES))
.route(web::post().to(login)),
)
.route("/session", web::get().to(session))
.route("/logout", web::post().to(logout)),
)
.route("/health", web::get().to(health))
.route("/dashboard", web::get().to(dashboard))
.route("/usage", web::get().to(get_ai_usage))
.route("/maintenance/reset", web::post().to(maintenance_reset))
.service(
web::scope("/podcasts")
.route("", web::get().to(list_podcasts))
.route("", web::post().to(add_podcast))
.route("", web::delete().to(bulk_delete_podcasts))
.route("/discover", web::post().to(discover_podcasts))
.route("/{id}", web::delete().to(delete_podcast)),
)
.service(
web::scope("/interviewees")
.route("", web::get().to(list_interviewees))
.route("", web::post().to(create_interviewee))
.route("", web::delete().to(bulk_delete_interviewees))
.route("/extract", web::post().to(start_interviewee_extraction))
.route("/{id}", web::patch().to(update_interviewee))
.route("/{id}", web::delete().to(delete_interviewee)),
)
.service(
web::scope("/contacts")
.route("", web::get().to(list_contacts))
.route("", web::post().to(create_contact))
.route("", web::delete().to(bulk_delete_contacts))
.route("/extract", web::post().to(start_contact_extraction))
.route("/export", web::get().to(export_contacts))
.route("/{id}", web::patch().to(update_contact))
.route("/{id}", web::delete().to(delete_contact)),
)
.service(
web::scope("/runs")
.route("", web::get().to(list_runs))
.route("/{id}/cancel", web::post().to(cancel_run))
.route("/{id}/retry", web::post().to(retry_run)),
)
.service(
web::scope("/reviews")
.route("", web::get().to(list_reviews))
.route("", web::delete().to(reject_all_reviews))
.route("/{id}/resolve", web::post().to(resolve_review)),
),
);
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct BulkIdsBody {
ids: Vec<Uuid>,
}
fn validate_bulk_ids(ids: &[Uuid]) -> AppResult<()> {
if ids.is_empty() {
return Err(AppError::Validation("ids não pode ficar vazio".into()));
}
if ids.len() > MAX_BULK_IDS {
return Err(AppError::Validation(format!(
"ids excede o máximo de {MAX_BULK_IDS} itens por requisição"
)));
}
Ok(())
}
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct LoginRequest {
username: String,
password: String,
}
async fn login(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<LoginRequest>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let peer = request.peer_addr().map(|address| address.ip());
let login = match state.auth.login(peer, &body.username, &body.password).await {
Ok(login) => login,
Err(error) => {
logs::warn(MODULE, &request_id, "login rejeitado");
return Err(error);
}
};
let cookie = state.auth.session_cookie(login.token);
logs::info(MODULE, &request_id, "sessão autenticada criada");
Ok(HttpResponse::Ok()
.cookie(cookie)
.json(DataResponse::new(json!({ "username": login.username }))))
}
async fn session(request: HttpRequest) -> AppResult<HttpResponse> {
let identity = request
.extensions()
.get::<SessionIdentity>()
.cloned()
.ok_or_else(|| AppError::Unauthorized(AUTH_FAILURE_MESSAGE.into()))?;
Ok(HttpResponse::Ok().json(DataResponse::new(json!({
"username": identity.username,
}))))
}
async fn logout(request: HttpRequest, state: web::Data<AppState>) -> AppResult<HttpResponse> {
let token = request
.cookie(SESSION_COOKIE_NAME)
.map(|cookie| cookie.value().to_owned())
.ok_or_else(|| AppError::Unauthorized(AUTH_FAILURE_MESSAGE.into()))?;
state.auth.logout(&token).await;
let request_id = request_id::from_request(&request);
logs::info(MODULE, &request_id, "sessão autenticada encerrada");
Ok(HttpResponse::NoContent()
.cookie(state.auth.removal_cookie())
.finish())
}
async fn health(state: web::Data<AppState>) -> AppResult<HttpResponse> {
state.db.health().await?;
Ok(HttpResponse::Ok().json(DataResponse::new(json!({
"status": "ok",
"database": "ok",
"openAiConfigured": state.ai.configured(),
"model": state.ai.model(),
}))))
}
async fn dashboard(state: web::Data<AppState>) -> AppResult<HttpResponse> {
Ok(HttpResponse::Ok().json(DataResponse::new(api_queries::dashboard(&state.db).await?)))
}
async fn get_ai_usage(state: web::Data<AppState>) -> AppResult<HttpResponse> {
Ok(HttpResponse::Ok().json(DataResponse::new(
api_queries::get_ai_usage_view(&state.db, &state.config).await?,
)))
}
async fn maintenance_reset(
request: HttpRequest,
state: web::Data<AppState>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let summary = maintenance::reset_interrupted_work(&state.db, &state.runs, &request_id).await?;
logs::info(
MODULE,
&request_id,
"manutenção solicitada manualmente via API",
);
Ok(HttpResponse::Ok().json(DataResponse::new(summary)))
}
#[derive(Debug, Deserialize, Default)]
#[serde(default, rename_all = "camelCase")]
struct ListQuery {
query: Option<String>,
status: Option<String>,
category: Option<String>,
#[serde(rename = "type")]
item_type: Option<String>,
relationship: Option<String>,
kind: Option<String>,
priority: Option<String>,
sort: Option<String>,
page: i64,
page_size: i64,
}
fn paging(query: &ListQuery) -> (i64, i64) {
let page_size = if query.page_size <= 0 {
10
} else {
query.page_size.clamp(1, 250)
};
(query.page.max(1), page_size)
}
fn filter(value: Option<&str>) -> Option<&str> {
value
.map(str::trim)
.filter(|value| !value.is_empty() && !value.eq_ignore_ascii_case("all"))
}
async fn list_podcasts(
state: web::Data<AppState>,
query: web::Query<ListQuery>,
) -> AppResult<HttpResponse> {
let (page, page_size) = paging(&query);
let data = api_queries::list_podcasts(
&state.db,
filter(query.query.as_deref()),
filter(query.status.as_deref()),
page,
page_size,
)
.await?;
Ok(HttpResponse::Ok().json(DataResponse::new(data)))
}
async fn list_interviewees(
state: web::Data<AppState>,
query: web::Query<ListQuery>,
) -> AppResult<HttpResponse> {
let (page, page_size) = paging(&query);
let data = api_queries::list_interviewees(
&state.db,
filter(query.query.as_deref()),
filter(query.status.as_deref()),
filter(query.category.as_deref()),
page,
page_size,
)
.await?;
Ok(HttpResponse::Ok().json(DataResponse::new(data)))
}
async fn list_contacts(
state: web::Data<AppState>,
query: web::Query<ListQuery>,
) -> AppResult<HttpResponse> {
let (page, page_size) = paging(&query);
let data = api_queries::list_contacts(
&state.db,
filter(query.query.as_deref()),
filter(query.item_type.as_deref()),
filter(query.relationship.as_deref()),
filter(query.status.as_deref()),
page,
page_size,
)
.await?;
Ok(HttpResponse::Ok().json(DataResponse::new(data)))
}
#[derive(Debug, Deserialize, Default)]
#[serde(default, rename_all = "camelCase")]
struct ExportContactsQuery {
query: Option<String>,
category: Option<String>,
format: Option<String>,
columns: Option<String>,
only_with_best: bool,
}
/// Exporta contatos processados pelo backend em CSV ou XLSX: uma linha por
/// entrevistado (somente os que têm ao menos um contato ativo), com uma
/// coluna por tipo de contato e duas colunas inteligentes
/// (`melhor_email`/`melhor_numero`) escolhidas pela IA entre os contatos já
/// verificados, priorizando pessoal e caindo para comercial só quando não há
/// pessoal daquele tipo. O cálculo é armazenado em cache em `interviewees` e
/// só é refeito quando algum contato ativo muda depois do último cálculo.
async fn export_contacts(
request: HttpRequest,
state: web::Data<AppState>,
query: web::Query<ExportContactsQuery>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let format = export::ExportFormat::parse(query.format.as_deref());
let columns = export::parse_columns(query.columns.as_deref());
let candidates = api_queries::export_contacts_candidates(
&state.db,
filter(query.query.as_deref()),
filter(query.category.as_deref()),
)
.await?;
let concurrency = state.config.worker_concurrency.clamp(1, 8);
let rows: Vec<export::ExportRow> = futures::stream::iter(candidates)
.map(|candidate| {
let state = state.clone();
let request_id = request_id.clone();
async move {
let stale = best_contacts::is_stale(
candidate.interviewee.best_contacts_computed_at,
candidate.contacts_last_activity,
);
let (best_email, best_phone) = if stale {
let resolved = best_contacts::resolve_and_persist(
&state.db,
&state.ai,
&request_id,
&candidate.interviewee,
)
.await;
(resolved.best_email, resolved.best_phone)
} else {
(
candidate.interviewee.best_email.clone(),
candidate.interviewee.best_phone.clone(),
)
};
export::ExportRow {
interviewee_id: candidate.interviewee.id,
display_name: candidate.interviewee.display_name,
category: candidate.category,
profession: candidate.interviewee.profession,
description: candidate.interviewee.professional_summary,
created_at: candidate.interviewee.created_at,
best_email,
best_phone,
by_type: candidate.by_type,
}
}
})
.buffer_unordered(concurrency)
.collect()
.await;
// Um entrevistado sem nenhum contato ativo não tem o que exportar: a
// exportação é sempre restrita a quem já tem pelo menos um contato
// encontrado, para não encher o arquivo com linhas em branco.
let rows: Vec<export::ExportRow> = rows
.into_iter()
.filter(|row| !row.by_type.is_empty())
.collect();
let rows: Vec<export::ExportRow> = if query.only_with_best {
rows.into_iter()
.filter(|row| row.best_email.is_some() || row.best_phone.is_some())
.collect()
} else {
rows
};
logs::info(
MODULE,
&request_id,
format!(
"exportação de contatos gerada com {} entrevistados",
rows.len()
),
);
let bytes = export::build(&rows, &columns, format)?;
Ok(HttpResponse::Ok()
.content_type(format.content_type())
.insert_header((
CONTENT_DISPOSITION,
format!("attachment; filename=\"{}\"", format.filename()),
))
.body(bytes))
}
async fn list_runs(
state: web::Data<AppState>,
query: web::Query<ListQuery>,
) -> AppResult<HttpResponse> {
let (page, page_size) = paging(&query);
let data = api_queries::list_runs(
&state.db,
filter(query.query.as_deref()),
filter(query.status.as_deref()),
filter(query.item_type.as_deref()),
page,
page_size,
)
.await?;
Ok(HttpResponse::Ok().json(DataResponse::new(data)))
}
async fn list_reviews(
state: web::Data<AppState>,
query: web::Query<ListQuery>,
) -> AppResult<HttpResponse> {
let (page, page_size) = paging(&query);
let data = api_queries::list_reviews(
&state.db,
filter(query.query.as_deref()),
filter(query.kind.as_deref()),
filter(query.priority.as_deref()),
filter(query.sort.as_deref()),
page,
page_size,
)
.await?;
Ok(HttpResponse::Ok().json(DataResponse::new(data)))
}
async fn reject_all_reviews(
request: HttpRequest,
state: web::Data<AppState>,
query: web::Query<ListQuery>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let refs = api_queries::list_review_refs(
&state.db,
filter(query.query.as_deref()),
filter(query.kind.as_deref()),
filter(query.priority.as_deref()),
)
.await?;
for (id, kind) in &refs {
let run_id = if kind == "identity" {
let candidate = interviewee_candidate::decide(&state.db, *id, "rejected", None).await?;
candidate.map(|candidate| candidate.run_id)
} else {
let candidate =
contact_candidate::decide(&state.db, *id, "rejected", None, None).await?;
candidate.map(|candidate| candidate.run_id)
};
audit_event::append(
&state.db,
&NewAuditEvent {
run_id,
actor_type: "user".into(),
actor_id: None,
action: "reject".into(),
entity_type: if kind == "identity" {
"interviewee_candidate"
} else {
"contact_candidate"
}
.into(),
entity_id: Some(*id),
before_data: None,
after_data: Some(json!({ "decision": "rejected" })),
metadata: json!({ "note": "bulk rejection" }),
},
)
.await?;
}
logs::info(
MODULE,
&request_id,
format!("revisões rejeitadas em lote total={}", refs.len()),
);
Ok(HttpResponse::NoContent().finish())
}
#[derive(Debug, Default)]
enum PatchField<T> {
#[default]
Missing,
Null,
Value(T),
}
impl<'de, T> Deserialize<'de> for PatchField<T>
where
T: Deserialize<'de>,
{
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
Option::<T>::deserialize(deserializer)
.map(|value| value.map(Self::Value).unwrap_or(Self::Null))
}
}
impl<T> PatchField<T> {
fn is_missing(&self) -> bool {
matches!(self, Self::Missing)
}
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct CreateIntervieweeBody {
display_name: String,
real_name: Option<String>,
brand_name: Option<String>,
category: Option<String>,
professional_summary: Option<String>,
public_bio: Option<String>,
profession: Option<String>,
content_type: Option<String>,
audience: Option<String>,
}
#[derive(Debug, Default, Deserialize)]
#[serde(default, rename_all = "camelCase", deny_unknown_fields)]
struct UpdateIntervieweeBody {
display_name: PatchField<String>,
real_name: PatchField<String>,
brand_name: PatchField<String>,
category: PatchField<String>,
professional_summary: PatchField<String>,
public_bio: PatchField<String>,
profession: PatchField<String>,
content_type: PatchField<String>,
audience: PatchField<String>,
}
impl UpdateIntervieweeBody {
fn is_empty(&self) -> bool {
self.display_name.is_missing()
&& self.real_name.is_missing()
&& self.brand_name.is_missing()
&& self.category.is_missing()
&& self.professional_summary.is_missing()
&& self.public_bio.is_missing()
&& self.profession.is_missing()
&& self.content_type.is_missing()
&& self.audience.is_missing()
}
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct CreateContactBody {
interviewee_id: Uuid,
#[serde(rename = "type")]
contact_type: String,
value: String,
relationship: String,
label: Option<String>,
confidence: Option<f32>,
source_name: String,
source_url: String,
status: Option<String>,
}
#[derive(Debug, Default, Deserialize)]
#[serde(default, rename_all = "camelCase", deny_unknown_fields)]
struct UpdateContactBody {
interviewee_id: PatchField<Uuid>,
#[serde(rename = "type")]
contact_type: PatchField<String>,
value: PatchField<String>,
relationship: PatchField<String>,
label: PatchField<String>,
confidence: PatchField<f32>,
source_name: PatchField<String>,
source_url: PatchField<String>,
status: PatchField<String>,
}
impl UpdateContactBody {
fn is_empty(&self) -> bool {
self.interviewee_id.is_missing()
&& self.contact_type.is_missing()
&& self.value.is_missing()
&& self.relationship.is_missing()
&& self.label.is_missing()
&& self.confidence.is_missing()
&& self.source_name.is_missing()
&& self.source_url.is_missing()
&& self.status.is_missing()
}
}
async fn create_interviewee(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<CreateIntervieweeBody>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let input = create_interviewee_input(&state, &body).await?;
let created = interviewee::create(&state.db, &input).await?;
let saved = interviewee::update(
&state.db,
created.id,
&IntervieweePatch {
dedup_review_status: Some("confirmed".into()),
..IntervieweePatch::default()
},
)
.await?
.ok_or_else(|| AppError::NotFound("entrevistado recém-criado não encontrado".into()))?;
append_manual_audit(
&state,
&request_id,
"create",
"interviewee",
saved.id,
None,
Some(serde_json::to_value(&saved)?),
)
.await?;
logs::info(
MODULE,
&request_id,
format!("entrevistado criado manualmente id={}", saved.id),
);
Ok(HttpResponse::Created().json(DataResponse::new(
api_queries::get_interviewee_view(&state.db, saved.id).await?,
)))
}
async fn update_interviewee(
request: HttpRequest,
state: web::Data<AppState>,
id: web::Path<Uuid>,
body: web::Json<UpdateIntervieweeBody>,
) -> AppResult<HttpResponse> {
if body.is_empty() {
return Err(AppError::Validation(
"informe ao menos um campo para atualizar o entrevistado".into(),
));
}
let request_id = request_id::from_request(&request);
let interviewee_id = id.into_inner();
let before = interviewee::get(&state.db, interviewee_id)
.await?
.ok_or_else(|| AppError::NotFound("entrevistado não encontrado".into()))?;
let profile = apply_interviewee_patch(&state, &before, &body).await?;
let saved = interviewee::replace_manual(&state.db, interviewee_id, &profile)
.await?
.ok_or_else(|| AppError::NotFound("entrevistado não encontrado".into()))?;
append_manual_audit(
&state,
&request_id,
"update",
"interviewee",
saved.id,
Some(serde_json::to_value(&before)?),
Some(serde_json::to_value(&saved)?),
)
.await?;
logs::info(
MODULE,
&request_id,
format!("entrevistado atualizado manualmente id={}", saved.id),
);
Ok(HttpResponse::Ok().json(DataResponse::new(
api_queries::get_interviewee_view(&state.db, saved.id).await?,
)))
}
async fn delete_interviewee(
request: HttpRequest,
state: web::Data<AppState>,
id: web::Path<Uuid>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
if !soft_delete_interviewee(&state, &request_id, id.into_inner()).await? {
return Err(AppError::NotFound("entrevistado não encontrado".into()));
}
Ok(HttpResponse::NoContent().finish())
}
/// Remove em lote para o frontend evitar disparar uma requisição HTTP por
/// item selecionado (o que facilmente estoura o rate limit por IP em seleções
/// grandes). Ids inexistentes ou já removidos são ignorados silenciosamente.
async fn bulk_delete_interviewees(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<BulkIdsBody>,
) -> AppResult<HttpResponse> {
validate_bulk_ids(&body.ids)?;
let request_id = request_id::from_request(&request);
let mut removed = 0usize;
for id in &body.ids {
if soft_delete_interviewee(&state, &request_id, *id).await? {
removed += 1;
}
}
logs::info(
MODULE,
&request_id,
format!(
"entrevistados arquivados em lote total={removed} solicitados={}",
body.ids.len()
),
);
Ok(HttpResponse::NoContent().finish())
}
async fn soft_delete_interviewee(
state: &AppState,
request_id: &str,
interviewee_id: Uuid,
) -> AppResult<bool> {
let Some(before) = interviewee::get(&state.db, interviewee_id).await? else {
return Ok(false);
};
if !interviewee::soft_delete(&state.db, interviewee_id).await? {
return Ok(false);
}
let after = interviewee::get_including_inactive(&state.db, interviewee_id)
.await?
.map(serde_json::to_value)
.transpose()?;
append_manual_audit(
state,
request_id,
"delete",
"interviewee",
interviewee_id,
Some(serde_json::to_value(before)?),
after,
)
.await?;
logs::info(
MODULE,
request_id,
format!("entrevistado arquivado manualmente id={interviewee_id}"),
);
Ok(true)
}
async fn create_contact(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<CreateContactBody>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let (input, storage_status, last_verified_at) = create_contact_input(&state, &body).await?;
let saved =
contact::create_manual(&state.db, &input, &storage_status, last_verified_at).await?;
append_manual_audit(
&state,
&request_id,
"create",
"contact",
saved.id,
None,
Some(serde_json::to_value(&saved)?),
)
.await?;
logs::info(
MODULE,
&request_id,
format!("contato criado manualmente id={}", saved.id),
);
Ok(HttpResponse::Created().json(DataResponse::new(
api_queries::get_contact_view(&state.db, saved.id).await?,
)))
}
async fn update_contact(
request: HttpRequest,
state: web::Data<AppState>,
id: web::Path<Uuid>,
body: web::Json<UpdateContactBody>,
) -> AppResult<HttpResponse> {
if body.is_empty() {
return Err(AppError::Validation(
"informe ao menos um campo para atualizar o contato".into(),
));
}
let request_id = request_id::from_request(&request);
let contact_id = id.into_inner();
let before = contact::get(&state.db, contact_id)
.await?
.ok_or_else(|| AppError::NotFound("contato não encontrado".into()))?;
let (input, storage_status, last_verified_at) =
apply_contact_patch(&state, &before, &body).await?;
let saved = contact::replace_manual(
&state.db,
contact_id,
&input,
&storage_status,
last_verified_at,
)
.await?
.ok_or_else(|| AppError::NotFound("contato não encontrado".into()))?;
append_manual_audit(
&state,
&request_id,
"update",
"contact",
saved.id,
Some(serde_json::to_value(&before)?),
Some(serde_json::to_value(&saved)?),
)
.await?;
logs::info(
MODULE,
&request_id,
format!("contato atualizado manualmente id={}", saved.id),
);
Ok(HttpResponse::Ok().json(DataResponse::new(
api_queries::get_contact_view(&state.db, saved.id).await?,
)))
}
async fn delete_contact(
request: HttpRequest,
state: web::Data<AppState>,
id: web::Path<Uuid>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
if !soft_delete_contact(&state, &request_id, id.into_inner()).await? {
return Err(AppError::NotFound("contato não encontrado".into()));
}
Ok(HttpResponse::NoContent().finish())
}
/// Remove em lote para o frontend evitar disparar uma requisição HTTP por
/// item selecionado. Ids inexistentes ou já removidos são ignorados
/// silenciosamente.
async fn bulk_delete_contacts(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<BulkIdsBody>,
) -> AppResult<HttpResponse> {
validate_bulk_ids(&body.ids)?;
let request_id = request_id::from_request(&request);
let mut removed = 0usize;
for id in &body.ids {
if soft_delete_contact(&state, &request_id, *id).await? {
removed += 1;
}
}
logs::info(
MODULE,
&request_id,
format!(
"contatos removidos em lote total={removed} solicitados={}",
body.ids.len()
),
);
Ok(HttpResponse::NoContent().finish())
}
async fn soft_delete_contact(state: &AppState, request_id: &str, contact_id: Uuid) -> AppResult<bool> {
let Some(before) = contact::get(&state.db, contact_id).await? else {
return Ok(false);
};
if !contact::soft_delete(&state.db, contact_id).await? {
return Ok(false);
}
let after = contact::get_including_deleted(&state.db, contact_id)
.await?
.map(serde_json::to_value)
.transpose()?;
append_manual_audit(
state,
request_id,
"delete",
"contact",
contact_id,
Some(serde_json::to_value(before)?),
after,
)
.await?;
logs::info(
MODULE,
request_id,
format!("contato removido manualmente id={contact_id}"),
);
Ok(true)
}
async fn create_interviewee_input(
state: &AppState,
body: &CreateIntervieweeBody,
) -> AppResult<NewInterviewee> {
let display_name = required_text(&body.display_name, "displayName", 200)?;
let real_name = optional_non_empty_text(body.real_name.as_deref(), "realName", 200)?;
let brand_name = optional_non_empty_text(body.brand_name.as_deref(), "brandName", 200)?;
let professional_summary = body
.professional_summary
.as_deref()
.map(|value| clearable_text(value, "professionalSummary", 10_000))
.transpose()?
.unwrap_or_default();
let public_bio = body
.public_bio
.as_deref()
.map(|value| clearable_text(value, "publicBio", 20_000))
.transpose()?
.unwrap_or_default();
let profession = optional_non_empty_text(body.profession.as_deref(), "profession", 200)?;
let creator_content_type =
optional_non_empty_text(body.content_type.as_deref(), "contentType", 500)?;
let creator_audience = optional_non_empty_text(body.audience.as_deref(), "audience", 500)?;
let normalized_display_name = normalized_person_name(&display_name, "displayName")?;
let normalized_real_name = real_name
.as_deref()
.map(|value| normalized_person_name(value, "realName"))
.transpose()?;
let normalized_brand_name = brand_name
.as_deref()
.map(|value| normalized_person_name(value, "brandName"))
.transpose()?;
let primary_category_id = match body.category.as_deref() {
Some(name) => Some(upsert_manual_category(state, name).await?),
None => None,
};
Ok(NewInterviewee {
primary_category_id,
normalized_display_name,
normalized_real_name,
normalized_brand_name,
display_name,
real_name,
brand_name,
professional_summary,
public_bio,
profession,
creator_content_type,
creator_audience,
professional_image_asset_id: None,
personal_image_asset_id: None,
created_in_run_id: None,
metadata: json!({ "manual": true }),
})
}
async fn apply_interviewee_patch(
state: &AppState,
before: &Interviewee,
patch: &UpdateIntervieweeBody,
) -> AppResult<NewInterviewee> {
let display_name = patched_required_text(
&patch.display_name,
&before.display_name,
"displayName",
200,
)?;
let real_name = patched_nullable_text(&patch.real_name, &before.real_name, "realName", 200)?;
let brand_name =
patched_nullable_text(&patch.brand_name, &before.brand_name, "brandName", 200)?;
let professional_summary = patched_clearable_text(
&patch.professional_summary,
&before.professional_summary,
"professionalSummary",
10_000,
)?;
let public_bio =
patched_clearable_text(&patch.public_bio, &before.public_bio, "publicBio", 20_000)?;
let profession =
patched_nullable_text(&patch.profession, &before.profession, "profession", 200)?;
let creator_content_type = patched_nullable_text(
&patch.content_type,
&before.creator_content_type,
"contentType",
500,
)?;
let creator_audience =
patched_nullable_text(&patch.audience, &before.creator_audience, "audience", 500)?;
let normalized_display_name = normalized_person_name(&display_name, "displayName")?;
let normalized_real_name = real_name
.as_deref()
.map(|value| normalized_person_name(value, "realName"))
.transpose()?;
let normalized_brand_name = brand_name
.as_deref()
.map(|value| normalized_person_name(value, "brandName"))
.transpose()?;
let primary_category_id = match &patch.category {
PatchField::Missing => before.primary_category_id,
PatchField::Null => None,
PatchField::Value(name) => Some(upsert_manual_category(state, name).await?),
};
Ok(NewInterviewee {
primary_category_id,
normalized_display_name,
normalized_real_name,
normalized_brand_name,
display_name,
real_name,
brand_name,
professional_summary,
public_bio,
profession,
creator_content_type,
creator_audience,
professional_image_asset_id: before.professional_image_asset_id,
personal_image_asset_id: before.personal_image_asset_id,
created_in_run_id: before.created_in_run_id,
metadata: before.metadata.clone(),
})
}
async fn upsert_manual_category(state: &AppState, value: &str) -> AppResult<Uuid> {
let display_name = required_text(value, "category", 200)?;
let normalized_name = normalize_name(&display_name);
if normalized_name.is_empty() {
return Err(AppError::Validation(
"category deve conter letras ou números".into(),
));
}
Ok(category::upsert(
&state.db,
&NewCategory {
display_name,
normalized_name,
description: String::new(),
created_by: "manual".into(),
},
)
.await?
.id)
}
async fn create_contact_input(
state: &AppState,
body: &CreateContactBody,
) -> AppResult<(NewContact, String, Option<DateTime<Utc>>)> {
ensure_active_interviewee(state, body.interviewee_id).await?;
let normalized = normalize_manual_contact(&body.contact_type, &body.value)?;
let relationship = validate_relationship(&body.relationship)?;
let label = optional_non_empty_text(body.label.as_deref(), "label", 200)?;
let confidence = validate_confidence(body.confidence.unwrap_or(100.0))?;
let (storage_status, last_verified_at) =
contact_storage_status(body.status.as_deref().unwrap_or("pending"), false)?;
let source = upsert_manual_origin(
state,
&body.source_name,
&body.source_url,
&normalized.contact_type,
)
.await?;
Ok((
NewContact {
interviewee_id: body.interviewee_id,
primary_origin_id: source.id,
contact_type: normalized.contact_type,
raw_value: normalized.raw_value,
normalized_value: normalized.normalized_value,
relationship_kind: relationship,
label,
confidence,
discovered_in_run_id: None,
},
storage_status,
last_verified_at,
))
}
async fn apply_contact_patch(
state: &AppState,
before: &Contact,
patch: &UpdateContactBody,
) -> AppResult<(NewContact, String, Option<DateTime<Utc>>)> {
let interviewee_id = match patch.interviewee_id {
PatchField::Missing => before.interviewee_id,
PatchField::Null => {
return Err(AppError::Validation(
"intervieweeId não pode ser null".into(),
));
}
PatchField::Value(value) => value,
};
ensure_active_interviewee(state, interviewee_id).await?;
let contact_type =
patched_required_text(&patch.contact_type, &before.contact_type, "type", 30)?;
let raw_value = patched_required_text(&patch.value, &before.raw_value, "value", 2_048)?;
let normalized = normalize_manual_contact(&contact_type, &raw_value)?;
let relationship = validate_relationship(&patched_required_text(
&patch.relationship,
&before.relationship_kind,
"relationship",
30,
)?)?;
let label = patched_nullable_text(&patch.label, &before.label, "label", 200)?;
let confidence = match patch.confidence {
PatchField::Missing => before.confidence,
PatchField::Null => {
return Err(AppError::Validation("confidence não pode ser null".into()));
}
PatchField::Value(value) => validate_confidence(value)?,
};
let (storage_status, last_verified_at) = match &patch.status {
PatchField::Missing => (before.status.clone(), before.last_verified_at),
PatchField::Null => {
return Err(AppError::Validation("status não pode ser null".into()));
}
PatchField::Value(value) => contact_storage_status(value, true)?,
};
let primary_origin_id = if patch.source_name.is_missing() && patch.source_url.is_missing() {
before.primary_origin_id
} else {
let current = origin::get(&state.db, before.primary_origin_id)
.await?
.ok_or_else(|| AppError::NotFound("origem atual do contato não encontrada".into()))?;
let source_name =
patched_required_text(&patch.source_name, &current.display_name, "sourceName", 200)?;
let source_url = patched_required_text(
&patch.source_url,
&current.canonical_url,
"sourceUrl",
4_096,
)?;
upsert_manual_origin(state, &source_name, &source_url, &normalized.contact_type)
.await?
.id
};
Ok((
NewContact {
interviewee_id,
primary_origin_id,
contact_type: normalized.contact_type,
raw_value: normalized.raw_value,
normalized_value: normalized.normalized_value,
relationship_kind: relationship,
label,
confidence,
discovered_in_run_id: before.discovered_in_run_id,
},
storage_status,
last_verified_at,
))
}
async fn ensure_active_interviewee(state: &AppState, id: Uuid) -> AppResult<()> {
if interviewee::get(&state.db, id).await?.is_none() {
return Err(AppError::Validation(
"intervieweeId deve identificar um entrevistado ativo".into(),
));
}
Ok(())
}
async fn upsert_manual_origin(
state: &AppState,
name: &str,
url: &str,
contact_type: &str,
) -> AppResult<crate::db::models::Origin> {
let display_name = required_text(name, "sourceName", 200)?;
let raw_url = required_text(url, "sourceUrl", 4_096)?;
let canonical_url = contact_normalizer::canonical_url(&raw_url)?;
let parsed = Url::parse(&canonical_url)?;
if !parsed.username().is_empty() || parsed.password().is_some() {
return Err(AppError::Validation(
"sourceUrl não pode conter credenciais".into(),
));
}
let domain = parsed
.host_str()
.map(|value| value.trim_start_matches("www.").to_ascii_lowercase())
.filter(|value| !value.is_empty())
.ok_or_else(|| AppError::Validation("sourceUrl deve conter um domínio".into()))?;
if domain.chars().count() > 253 {
return Err(AppError::Validation(
"o domínio de sourceUrl excede 253 caracteres".into(),
));
}
let source_type = infer_origin_type(&domain, contact_type);
origin::upsert(
&state.db,
&NewOrigin {
display_name,
canonical_url,
domain,
source_type,
icon_asset_id: None,
metadata: json!({ "manual": true }),
},
)
.await
}
fn infer_origin_type(domain: &str, contact_type: &str) -> String {
if domain.contains("youtube.com") || domain == "youtu.be" || contact_type == "youtube" {
"youtube"
} else if domain.contains("instagram.com") || contact_type == "instagram" {
"instagram"
} else if matches!(
domain,
"linktr.ee" | "beacons.ai" | "campsite.bio" | "solo.to"
) {
"linktree"
} else {
"website"
}
.into()
}
fn normalize_manual_contact(
contact_type: &str,
value: &str,
) -> AppResult<contact_normalizer::NormalizedContact> {
let contact_type = contact_type.trim().to_ascii_lowercase();
if !matches!(
contact_type.as_str(),
"email"
| "phone"
| "whatsapp"
| "instagram"
| "linkedin"
| "facebook"
| "tiktok"
| "x"
| "telegram"
| "youtube"
| "website"
| "other"
) {
return Err(AppError::Validation(format!(
"type inválido: {contact_type}"
)));
}
if value.chars().count() > 2_048 {
return Err(AppError::Validation("value excede 2048 caracteres".into()));
}
contact_normalizer::normalize(&contact_type, value)
}
fn validate_relationship(value: &str) -> AppResult<String> {
let relationship = value.trim().to_ascii_lowercase();
if !matches!(relationship.as_str(), "personal" | "commercial") {
return Err(AppError::Validation(
"relationship deve ser personal ou commercial".into(),
));
}
Ok(relationship)
}
fn validate_confidence(value: f32) -> AppResult<f32> {
if !value.is_finite() || !(0.0..=100.0).contains(&value) {
return Err(AppError::Validation(
"confidence deve estar entre 0 e 100".into(),
));
}
Ok(value / 100.0)
}
fn contact_storage_status(
value: &str,
allow_rejected: bool,
) -> AppResult<(String, Option<DateTime<Utc>>)> {
match value.trim().to_ascii_lowercase().as_str() {
"pending" => Ok(("active".into(), None)),
"verified" => Ok(("active".into(), Some(Utc::now()))),
"rejected" if allow_rejected => Ok(("suppressed".into(), None)),
_ if allow_rejected => Err(AppError::Validation(
"status deve ser pending, verified ou rejected".into(),
)),
_ => Err(AppError::Validation(
"status deve ser pending ou verified na criação".into(),
)),
}
}
fn required_text(value: &str, field: &str, max: usize) -> AppResult<String> {
let value = value.trim();
if value.is_empty() {
return Err(AppError::Validation(format!(
"{field} não pode ficar vazio"
)));
}
if value.chars().count() > max {
return Err(AppError::Validation(format!(
"{field} excede {max} caracteres"
)));
}
Ok(value.to_owned())
}
fn normalized_person_name(value: &str, field: &str) -> AppResult<String> {
let normalized = normalize_name(value);
if normalized.is_empty() {
return Err(AppError::Validation(format!(
"{field} deve conter letras ou números"
)));
}
Ok(normalized)
}
fn clearable_text(value: &str, field: &str, max: usize) -> AppResult<String> {
let value = value.trim();
if value.chars().count() > max {
return Err(AppError::Validation(format!(
"{field} excede {max} caracteres"
)));
}
Ok(value.to_owned())
}
fn optional_non_empty_text(
value: Option<&str>,
field: &str,
max: usize,
) -> AppResult<Option<String>> {
value
.map(|value| required_text(value, field, max))
.transpose()
}
fn patched_required_text(
patch: &PatchField<String>,
current: &str,
field: &str,
max: usize,
) -> AppResult<String> {
match patch {
PatchField::Missing => Ok(current.to_owned()),
PatchField::Null => Err(AppError::Validation(format!("{field} não pode ser null"))),
PatchField::Value(value) => required_text(value, field, max),
}
}
fn patched_nullable_text(
patch: &PatchField<String>,
current: &Option<String>,
field: &str,
max: usize,
) -> AppResult<Option<String>> {
match patch {
PatchField::Missing => Ok(current.clone()),
PatchField::Null => Ok(None),
PatchField::Value(value) => required_text(value, field, max).map(Some),
}
}
fn patched_clearable_text(
patch: &PatchField<String>,
current: &str,
field: &str,
max: usize,
) -> AppResult<String> {
match patch {
PatchField::Missing => Ok(current.to_owned()),
PatchField::Null => Ok(String::new()),
PatchField::Value(value) => clearable_text(value, field, max),
}
}
async fn append_manual_audit(
state: &AppState,
request_id: &str,
action: &str,
entity_type: &str,
entity_id: Uuid,
before_data: Option<Value>,
after_data: Option<Value>,
) -> AppResult<()> {
audit_event::append(
&state.db,
&NewAuditEvent {
run_id: None,
actor_type: "user".into(),
actor_id: None,
action: action.into(),
entity_type: entity_type.into(),
entity_id: Some(entity_id),
before_data,
after_data,
metadata: json!({
"requestId": request_id,
"source": "manual_api",
}),
},
)
.await?;
Ok(())
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct DiscoverPodcastsBody {
query: String,
limit: Option<usize>,
}
async fn discover_podcasts(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<DiscoverPodcastsBody>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let query = body.query.trim();
if query.is_empty() {
return Err(AppError::Validation("informe o termo da busca".into()));
}
let limit = body.limit.unwrap_or(10).clamp(1, 30);
logs::info(
MODULE,
&request_id,
format!("descoberta de podcasts query_len={}", query.len()),
);
let channels = tokio::time::timeout(
state.config.browser_timeout + Duration::from_secs(15),
state
.youtube
.discover_podcast_channels(&request_id, query, limit),
)
.await
.map_err(|_| AppError::Timeout("descoberta de podcasts".into()))??;
let mut output = Vec::with_capacity(channels.len());
for channel in channels {
let already_added = podcast_channel::get_by_youtube_id(&state.db, &channel.channel_id)
.await?
.is_some();
output.push(PodcastDiscoveryView {
id: channel.channel_id,
name: channel.name,
url: channel.url,
logo_url: channel.logo_url,
description: channel.description.unwrap_or_default(),
subscribers_text: channel.subscriber_count_text,
already_added,
});
}
Ok(HttpResponse::Ok().json(DataResponse::new(output)))
}
#[derive(Debug, Deserialize)]
struct AddPodcastBody {
url: String,
name: Option<String>,
}
async fn add_podcast(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<AddPodcastBody>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let channel = tokio::time::timeout(
state.config.browser_timeout + Duration::from_secs(15),
state.youtube.channel_from_url(&request_id, body.url.trim()),
)
.await
.map_err(|_| AppError::Timeout("leitura do canal do YouTube".into()))??;
let display_name = body
.name
.as_deref()
.map(str::trim)
.filter(|value| !value.is_empty())
.unwrap_or(&channel.name)
.to_owned();
let mut logo_asset_id = None;
if let Some(logo_url) = channel.logo_url.as_deref() {
match state
.media
.download_named(&request_id, logo_url, MediaKind::PodcastLogo, &display_name)
.await
{
Ok(stored) => {
let asset = media_asset::upsert(
&state.db,
&NewMediaAsset {
kind: "channel_logo".into(),
source_url: stored.source_url,
storage_path: stored.relative_path.to_string_lossy().into_owned(),
sha256: stored.sha256,
mime_type: stored.content_type,
size_bytes: stored.byte_size.min(i64::MAX as u64) as i64,
width: None,
height: None,
},
)
.await?;
logo_asset_id = Some(asset.id);
}
Err(error) => logs::warn(
MODULE,
&request_id,
format!("logo do canal não pôde ser armazenado: {error}"),
),
}
}
let mut saved = podcast_channel::upsert(
&state.db,
&NewPodcastChannel {
youtube_channel_id: channel.channel_id,
name: display_name,
canonical_url: channel.url,
logo_asset_id,
status: "active".into(),
metadata: json!({
"description": channel.description,
"handle": channel.handle,
"subscriberCountText": channel.subscriber_count_text,
"remoteLogoUrl": channel.logo_url,
}),
},
)
.await?;
if saved.status == "removed" {
podcast_channel::restore(&state.db, saved.id).await?;
saved = podcast_channel::update(
&state.db,
saved.id,
&crate::db::models::PodcastChannelPatch {
status: Some("active".into()),
..Default::default()
},
)
.await?
.ok_or_else(|| AppError::NotFound("canal restaurado não encontrado".into()))?;
}
logs::info(
MODULE,
&request_id,
format!("podcast salvo id={}", saved.id),
);
Ok(HttpResponse::Ok().json(DataResponse::new(
api_queries::get_podcast_view(&state.db, saved.id).await?,
)))
}
async fn delete_podcast(
request: HttpRequest,
state: web::Data<AppState>,
id: web::Path<Uuid>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
if !podcast_channel::soft_delete(&state.db, id.into_inner()).await? {
return Err(AppError::NotFound("podcast não encontrado".into()));
}
logs::info(MODULE, &request_id, "podcast removido");
Ok(HttpResponse::NoContent().finish())
}
/// Remove em lote para o frontend evitar disparar uma requisição HTTP por
/// item selecionado. Ids inexistentes ou já removidos são ignorados
/// silenciosamente.
async fn bulk_delete_podcasts(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<BulkIdsBody>,
) -> AppResult<HttpResponse> {
validate_bulk_ids(&body.ids)?;
let request_id = request_id::from_request(&request);
let mut removed = 0usize;
for id in &body.ids {
if podcast_channel::soft_delete(&state.db, *id).await? {
removed += 1;
}
}
logs::info(
MODULE,
&request_id,
format!(
"podcasts removidos em lote total={removed} solicitados={}",
body.ids.len()
),
);
Ok(HttpResponse::NoContent().finish())
}
#[derive(Debug, Deserialize)]
#[serde(rename_all = "camelCase")]
struct ExtractionBody {
#[serde(default)]
podcast_ids: Vec<Uuid>,
#[serde(default)]
interviewee_ids: Vec<Uuid>,
max_videos_per_channel: Option<usize>,
}
async fn start_interviewee_extraction(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<ExtractionBody>,
) -> AppResult<HttpResponse> {
let input = json!({
"podcastIds": body.podcast_ids,
"maxVideosPerChannel": body.max_videos_per_channel,
});
enqueue_pipeline(&request, &state, "interviewee_extraction", "manual", input).await
}
async fn start_contact_extraction(
request: HttpRequest,
state: web::Data<AppState>,
body: web::Json<ExtractionBody>,
) -> AppResult<HttpResponse> {
let input = json!({ "intervieweeIds": body.interviewee_ids });
enqueue_pipeline(&request, &state, "contact_extraction", "manual", input).await
}
async fn enqueue_pipeline(
request: &HttpRequest,
state: &web::Data<AppState>,
kind: &str,
mode: &str,
input: Value,
) -> AppResult<HttpResponse> {
if !state.ai.configured() {
return Err(AppError::Config(
"configure OPENAI_API_KEY em backend/.env antes de iniciar a extração".into(),
));
}
let usage = api_queries::get_ai_usage_view(&state.db, &state.config).await?;
if usage.limit_exceeded {
return Err(AppError::BudgetExceeded(format!(
"gasto mensal com IA (${:.2}) atingiu o limite configurado (${:.2}); aguarde o próximo mês ou aumente OPENAI_MONTHLY_BUDGET_USD",
usage.spent_usd, usage.limit_usd
)));
}
let request_id = request_id::from_request(request);
let run = pipeline_run::create(
&state.db,
&NewPipelineRun {
kind: kind.into(),
mode: mode.into(),
idempotency_key: None,
requested_by: None,
input: input.clone(),
progress_total: None,
},
)
.await?;
if let Err(error) = job::enqueue(
&state.db,
&NewJob {
run_id: Some(run.id),
parent_job_id: None,
kind: kind.into(),
payload: input,
priority: 5,
idempotency_key: None,
max_attempts: 6,
available_at: None,
},
)
.await
{
let _ = pipeline_run::fail(
&state.db,
run.id,
Some(error.code()),
Some(&error.to_string()),
)
.await;
return Err(error);
}
logs::info(
MODULE,
&request_id,
format!("pipeline enfileirado kind={kind} run_id={}", run.id),
);
Ok(HttpResponse::Accepted().json(DataResponse::new(
api_queries::get_run_view(&state.db, run.id).await?,
)))
}
async fn cancel_run(
request: HttpRequest,
state: web::Data<AppState>,
id: web::Path<Uuid>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let run_id = id.into_inner();
pipeline_run::request_cancel(&state.db, run_id)
.await?
.ok_or_else(|| AppError::Conflict("execução já terminou ou não existe".into()))?;
state.runs.cancel(run_id).await;
logs::info(
MODULE,
&request_id,
format!("cancelamento solicitado run_id={run_id}"),
);
Ok(HttpResponse::Ok().json(DataResponse::new(
api_queries::get_run_view(&state.db, run_id).await?,
)))
}
async fn retry_run(
request: HttpRequest,
state: web::Data<AppState>,
id: web::Path<Uuid>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let run_id = id.into_inner();
let source = pipeline_run::get(&state.db, run_id)
.await?
.ok_or_else(|| AppError::NotFound("execução não encontrada".into()))?;
if !matches!(source.status.as_str(), "failed" | "cancelled" | "paused") {
return Err(AppError::Conflict(
"somente execuções falhas, canceladas ou pausadas podem ser repetidas".into(),
));
}
if source.status == "paused" {
job::make_retry_available(&state.db, run_id)
.await?
.ok_or_else(|| {
AppError::Conflict(
"a execução pausada não possui job aguardando nova tentativa".into(),
)
})?;
pipeline_run::resume(&state.db, run_id)
.await?
.ok_or_else(|| AppError::NotFound("execução pausada não encontrada".into()))?;
logs::info(
MODULE,
&request_id,
format!("execução pausada retomada no mesmo job run_id={run_id}"),
);
return Ok(HttpResponse::Ok().json(DataResponse::new(
api_queries::get_run_view(&state.db, run_id).await?,
)));
}
enqueue_pipeline(&request, &state, &source.kind, &source.mode, source.input).await
}
#[derive(Debug, Deserialize)]
struct ResolveReviewBody {
decision: String,
note: Option<String>,
}
async fn resolve_review(
request: HttpRequest,
state: web::Data<AppState>,
id: web::Path<Uuid>,
body: web::Json<ResolveReviewBody>,
) -> AppResult<HttpResponse> {
let request_id = request_id::from_request(&request);
let review_id = id.into_inner();
let approved = match body.decision.as_str() {
"approved" => true,
"rejected" => false,
_ => {
return Err(AppError::Validation(
"decision deve ser approved ou rejected".into(),
));
}
};
let (entity_type, run_id) = if let Some(candidate) =
interviewee_candidate::get(&state.db, review_id).await?
{
if candidate.status != "pending" {
return Err(AppError::Conflict(
"revisão de identidade já resolvida".into(),
));
}
if approved {
let matched_id = if let Some(id) = candidate.matched_interviewee_id {
id
} else {
let fallback = category::upsert(
&state.db,
&NewCategory {
display_name: "Sem categoria".into(),
normalized_name: "sem categoria".into(),
description: "Classificação manual pendente".into(),
created_by: "system".into(),
},
)
.await?;
let person = interviewee::create(
&state.db,
&NewInterviewee {
primary_category_id: Some(fallback.id),
display_name: candidate.proposed_name.clone(),
real_name: candidate.proposed_real_name.clone(),
brand_name: candidate.proposed_brand_name.clone(),
normalized_display_name: candidate.normalized_name.clone(),
normalized_real_name: candidate
.proposed_real_name
.as_deref()
.map(normalize_name),
normalized_brand_name: candidate
.proposed_brand_name
.as_deref()
.map(normalize_name),
professional_summary: candidate.professional_summary.clone(),
public_bio: candidate.personal_summary.clone().unwrap_or_default(),
profession: candidate.profession.clone(),
creator_content_type: candidate.creator_content_type.clone(),
creator_audience: candidate.creator_audience.clone(),
professional_image_asset_id: None,
personal_image_asset_id: None,
created_in_run_id: Some(candidate.run_id),
metadata: json!({ "manualReview": true }),
},
)
.await?;
appearance::upsert(
&state.db,
&NewAppearance {
interviewee_id: person.id,
video_id: candidate.video_id,
confidence: candidate.confidence,
evidence: candidate.evidence.clone(),
evidence_hash: candidate.evidence_hash.clone(),
extraction_source: "manual".into(),
},
)
.await?;
person.id
};
interviewee_candidate::decide(&state.db, review_id, "created", Some(matched_id))
.await?;
} else {
interviewee_candidate::decide(&state.db, review_id, "rejected", None).await?;
}
("interviewee_candidate", Some(candidate.run_id))
} else if let Some(candidate) = contact_candidate::get(&state.db, review_id).await? {
if !matches!(candidate.status.as_str(), "pending" | "needs_review") {
return Err(AppError::Conflict("revisão de contato já resolvida".into()));
}
if approved {
let relationship = candidate
.proposed_relationship_kind
.clone()
.unwrap_or_else(|| "commercial".into());
let saved = contact::upsert_with_evidence(
&state.db,
&NewContact {
interviewee_id: candidate.interviewee_id,
primary_origin_id: candidate.origin_id,
contact_type: candidate.contact_type.clone(),
raw_value: candidate.raw_value.clone(),
normalized_value: candidate.normalized_value.clone(),
relationship_kind: relationship,
label: candidate.proposed_label.clone(),
confidence: candidate.confidence,
discovered_in_run_id: Some(candidate.run_id),
},
&NewContactEvidenceDraft {
origin_id: candidate.origin_id,
page_url: origin_url(&state, candidate.origin_id).await?,
evidence_text: candidate.evidence.clone(),
evidence_hash: Some(hash_text(&candidate.evidence)),
confidence: candidate.confidence,
},
)
.await?;
contact_candidate::decide(&state.db, review_id, "accepted", Some(saved.id), None)
.await?;
} else {
contact_candidate::decide(&state.db, review_id, "rejected", None, body.note.as_deref())
.await?;
}
("contact_candidate", Some(candidate.run_id))
} else {
return Err(AppError::NotFound("revisão não encontrada".into()));
};
audit_event::append(
&state.db,
&NewAuditEvent {
run_id,
actor_type: "user".into(),
actor_id: None,
action: if approved { "approve" } else { "reject" }.into(),
entity_type: entity_type.into(),
entity_id: Some(review_id),
before_data: None,
after_data: Some(json!({ "decision": body.decision })),
metadata: json!({ "note": body.note }),
},
)
.await?;
logs::info(
MODULE,
&request_id,
format!("revisão resolvida id={review_id}"),
);
Ok(HttpResponse::NoContent().finish())
}
async fn origin_url(state: &AppState, origin_id: Uuid) -> AppResult<String> {
let client = state.db.client().await?;
client
.query_opt(
"SELECT canonical_url FROM origins WHERE id = $1",
&[&origin_id],
)
.await?
.map(|row| row.get(0))
.ok_or_else(|| AppError::NotFound("origem do contato não encontrada".into()))
}
fn hash_text(value: &str) -> String {
format!("{:x}", Sha256::digest(value.as_bytes()))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn patch_distinguishes_missing_null_and_value() {
let empty: UpdateIntervieweeBody = serde_json::from_value(json!({})).unwrap();
assert!(empty.is_empty());
let clear: UpdateIntervieweeBody =
serde_json::from_value(json!({ "realName": null })).unwrap();
assert!(!clear.is_empty());
assert!(matches!(clear.real_name, PatchField::Null));
let update: UpdateIntervieweeBody =
serde_json::from_value(json!({ "realName": "Ada Lovelace" })).unwrap();
assert!(matches!(update.real_name, PatchField::Value(value) if value == "Ada Lovelace"));
}
#[test]
fn manual_contact_validation_uses_public_percentage() {
assert_eq!(validate_confidence(87.5).unwrap(), 0.875);
assert!(validate_confidence(100.1).is_err());
let contact = normalize_manual_contact("EMAIL", " USER@Example.com ").unwrap();
assert_eq!(contact.contact_type, "email");
assert_eq!(contact.normalized_value, "user@example.com");
}
#[test]
fn manual_payload_rejects_unknown_fields() {
let parsed = serde_json::from_value::<CreateContactBody>(json!({
"intervieweeId": Uuid::nil(),
"type": "email",
"value": "user@example.com",
"relationship": "commercial",
"sourceName": "Site",
"sourceUrl": "https://example.com",
"unexpected": true
}));
assert!(parsed.is_err());
}
}