Conector Scardua — recebe da API da FLS e entrega no Oracle do cliente

Servico que roda dentro da rede da Comercial Scardua e faz a ponte entre a
API publica da FLS (na VPS) e o Oracle privado (10.16.x), inalcancavel pela
internet.

Fluxo: Holmes -> API na VPS (trata e cifra) -> ESTE conector (abre e valida)
-> Oracle.

Recebimento:
- app/v1/compras.py: POST /v1/compras/dados
- app/services/cripto.py: abre o envelope AES-256-GCM + RSA-OAEP-SHA256 com a
  chave privada. O GCM autentica: corpo adulterado levanta InvalidTag em vez
  de devolver lixo
- app/schemas.py: modulo folha (so pydantic) com a config e o contrato
  PayloadCompras, que espelha o da API. Fora de sincronia devolve 422 de
  proposito, pra falhar explicito em vez de gravar dado torto

Seguranca:
- app/seguranca.py: header X-Token com compare_digest (nao vaza por tempo de
  resposta) + allowlist de IP da VPS
- .gitignore barra configs.json.*, *.bak-*, *.pem e *.key

Config:
- app/config.py e so o carregamento do configs.json
- o bloco "vps" guarda chave privada, token e ips permitidos

Estado: ainda NAO persiste. Recebe, valida e descarta (persistido: false).
O INSERT no Oracle e a proxima fase, e tem que ser MERGE por id_processo
porque o Holmes reentrega webhook.

Docs em docs/ — arquitetura, a decisao do conector (opcao A, com as
alternativas descartadas), deploy e problemas conhecidos.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
Ricardo 2026-08-18 06:05:03 -03:00
commit f902ccbecb
26 changed files with 1824 additions and 0 deletions

6
app/config.py Normal file
View file

@ -0,0 +1,6 @@
from pathlib import Path
from app.schemas import Settings
_PATH = Path(__file__).resolve().parent.parent / "configs.json"
settings = Settings.model_validate_json(_PATH.read_text(encoding="utf-8"))

45
app/middleware.py Normal file
View file

@ -0,0 +1,45 @@
import re
from fastapi import Request, status
from fastapi.responses import Response
BLOCKED_PATHS = {
"/",
"/metrics",
"/security.txt",
"/.env",
"/wp-admin",
"/wp-login.php",
"/admin",
"/config",
"/actuator",
"/nice%20ports%2C/Trinity.txt.bak",
}
BLOCKED_UA_PATTERNS = re.compile(
r"(nmap|nikto|masscan|zgrab|censys|shodan|nuclei|httpx|gobuster|dirbuster)",
re.IGNORECASE,
)
async def block_scanners(request: Request, call_next):
path = request.url.path
ua = request.headers.get("user-agent", "")
if (
path.startswith("/v1/")
or path.startswith("/v2/")
or path in ("/health", "/dados_retorno", "/token")
):
return await call_next(request)
if path in BLOCKED_PATHS or path.endswith((".bak", ".env", ".git", ".php")):
return Response(status_code=status.HTTP_403_FORBIDDEN)
if BLOCKED_UA_PATTERNS.search(ua):
return Response(status_code=status.HTTP_403_FORBIDDEN)
if not ua:
return Response(status_code=status.HTTP_403_FORBIDDEN)
return await call_next(request)

98
app/schemas.py Normal file
View file

@ -0,0 +1,98 @@
"""
Modelos do conector.
Modulo folha: nao importa nada do projeto, so pydantic. Config e contrato
moram aqui; quem le o configs.json e o app/config.py.
"""
from datetime import datetime
from pydantic import BaseModel
# ---------------------------------------------------------------------------
# Config
# ---------------------------------------------------------------------------
class ApiConfig(BaseModel):
port: int
ambiente: str
workers: int
class DocsConfig(BaseModel):
user: str
password: str
class BancoConfig(BaseModel):
"""Oracle do cliente. Fase 2 — o conector ainda nao insere."""
user: str
password: str
dns: str
class VpsConfig(BaseModel):
"""Como a VPS da FLS se identifica e como abrimos o que ela manda."""
# Chave PRIVADA (PEM). Abre o envelope cifrado com a nossa publica.
# Nunca sai daqui — e o que garante que so o conector le o payload.
chave_privada: str
# Segredo compartilhado, esperado no header X-Token. Sem isso a rota
# fica aberta pra quem souber a URL.
token: str
# Allowlist de origem. Vazio = desligado (util em teste local).
ips_permitidos: list[str] = []
class Settings(BaseModel):
api: ApiConfig
docs: DocsConfig
vps: VpsConfig
banco: BancoConfig | None = None
log_level: str = "INFO"
# ---------------------------------------------------------------------------
# Contrato com a API da VPS
# ---------------------------------------------------------------------------
class Envelope(BaseModel):
"""O que chega no corpo do POST, ainda cifrado."""
alg: str
chave: str
nonce: str
dados: str
class Parcela(BaseModel):
valor: float
vencimento: datetime | None = None
class PayloadCompras(BaseModel):
"""
O que sai de dentro do envelope depois de decifrado.
Espelha o PayloadCompras da API da VPS. Se um lado mudar, o outro tem
que mudar junto e o jeito de descobrir e o teste, nao a producao.
"""
# Chave de deduplicacao: id do processo no Holmes. O MERGE no Oracle
# vai por aqui, senao reentrega de webhook vira linha duplicada.
id_processo: str
protocolo: str
cnpj: str
pedido_linx: str
fornecedor: str
tipo: str
nf_entrada: str | None = None
aprovador: str
valor_total: float
parcelas: dict[int, Parcela]

57
app/security.py Normal file
View file

@ -0,0 +1,57 @@
from datetime import datetime, timedelta, timezone
import jwt
from app.config import settings
from app.infra.database import db_instance
from fastapi import Depends, HTTPException, status
from fastapi.security import OAuth2PasswordBearer
ALGORITHM = "HS256"
TOKEN_EXPIRE_MINUTES = 5
oauth2_scheme = OAuth2PasswordBearer(tokenUrl="/token")
def criar_token(data: dict) -> str:
payload = data.copy()
payload["exp"] = datetime.now(timezone.utc) + timedelta(
minutes=TOKEN_EXPIRE_MINUTES
)
return jwt.encode(payload, settings.api.jwt_secret, algorithm=ALGORITHM)
def verificar_credenciais(client_id: str, client_secret: str) -> bool:
db_instance.create_pool()
assert db_instance.pool is not None
conn = db_instance.pool.acquire()
try:
cursor = conn.cursor()
cursor.execute(
"""
SELECT 1
FROM orvel_ti.CAD_USUARIO
WHERE PERFIL_TI = 'operador_ti'
AND ativo = 'S'
AND API = 'S'
AND login = :login
AND senha = :senha
""",
{"login": client_id, "senha": client_secret},
)
return cursor.fetchone() is not None
finally:
db_instance.pool.release(conn)
async def token_valido(token: str = Depends(oauth2_scheme)) -> dict:
try:
payload = jwt.decode(token, settings.api.jwt_secret, algorithms=[ALGORITHM])
return payload
except jwt.ExpiredSignatureError:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED, detail="Token expirado"
)
except jwt.InvalidTokenError:
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED, detail="Token inválido"
)

42
app/seguranca.py Normal file
View file

@ -0,0 +1,42 @@
"""
Quem pode falar com o conector.
Este servico escreve no banco de producao do cliente e alvo de alto valor.
As duas travas aqui sao o minimo: segredo compartilhado e origem conhecida.
Nao e seguranca por obscuridade (URL secreta nao conta).
"""
import logging
import secrets
from app.config import settings
from fastapi import Header, HTTPException, Request, status
def _origem(request: Request) -> str:
"""
IP de origem.
Se um dia entrar proxy na frente, o IP do socket vira o do proxy ai o
uvicorn precisa de --proxy-headers e isto passa a ler X-Forwarded-For.
"""
return request.client.host if request.client else ""
async def autorizar(request: Request, x_token: str = Header(default="")) -> None:
"""Dependencia de rota: barra quem nao for a VPS."""
permitidos = settings.vps.ips_permitidos
if permitidos:
ip = _origem(request)
if ip not in permitidos:
logging.warning(f"Origem recusada: {ip}")
raise HTTPException(
status_code=status.HTTP_403_FORBIDDEN, detail="origem nao autorizada"
)
# compare_digest pra nao vazar o segredo por tempo de resposta.
if not secrets.compare_digest(x_token, settings.vps.token):
logging.warning(f"Token invalido vindo de {_origem(request)}")
raise HTTPException(
status_code=status.HTTP_401_UNAUTHORIZED, detail="token invalido"
)

View file

@ -0,0 +1,33 @@
# app/services/api_tracker.py
import logging
def registrar_contador(
conn, api, endpoint, metodo, status_code, sucesso, duracao_ms, erro=None
):
"""Grava uma linha na tabela contadores_api."""
cursor = conn.cursor()
try:
cursor.execute(
"""
INSERT INTO orvel_ti.contadores_api
(api, endpoint, metodo_http, status_code, sucesso, duracao_ms, mensagem_erro)
VALUES
(:api, :endpoint, :metodo, :status, :sucesso, :duracao, :erro)
""",
{
"api": api,
"endpoint": endpoint,
"metodo": metodo,
"status": status_code,
"sucesso": "S" if sucesso else "N",
"duracao": duracao_ms,
"erro": erro[:500] if erro else None,
},
)
conn.commit()
except Exception as e:
conn.rollback()
logging.error(f"Falha ao gravar contador: {e}")
finally:
cursor.close()

47
app/services/cripto.py Normal file
View file

@ -0,0 +1,47 @@
"""
Abertura do envelope que a API da VPS manda.
Espelho do app/services/cripto.py de la, so que do lado que DECIFRA.
Envelope hibrido: AES-256-GCM nos dados, RSA-OAEP na chave AES.
O GCM autentica: se o corpo for adulterado no caminho, o decrypt levanta
excecao em vez de devolver lixo. Ou seja, decifrou = veio integro.
"""
import base64
import json
from cryptography.hazmat.primitives import hashes, serialization
from cryptography.hazmat.primitives.asymmetric import padding
from cryptography.hazmat.primitives.ciphers.aead import AESGCM
ALGORITMO = "RSA-OAEP-256+A256GCM"
_OAEP = padding.OAEP(
mgf=padding.MGF1(algorithm=hashes.SHA256()),
algorithm=hashes.SHA256(),
label=None,
)
def carregar_chave_privada(pem: str, senha: bytes | None = None):
return serialization.load_pem_private_key(pem.encode(), password=senha)
def descriptografar(envelope: dict, chave_privada) -> dict:
"""
Abre o envelope e devolve o dict original.
Levanta se o algoritmo nao for o esperado, se a chave nao for a par da
publica usada la, ou se o conteudo tiver sido mexido.
"""
if envelope.get("alg") != ALGORITMO:
raise ValueError(f"algoritmo inesperado: {envelope.get('alg')!r}")
chave_aes = chave_privada.decrypt(base64.b64decode(envelope["chave"]), _OAEP)
corpo = AESGCM(chave_aes).decrypt(
base64.b64decode(envelope["nonce"]),
base64.b64decode(envelope["dados"]),
None,
)
return json.loads(corpo)

704
app/services/holmes.py Normal file
View file

@ -0,0 +1,704 @@
import asyncio
import json
import logging
import os
import random
import re
import tempfile
import time
import httpx
from app.config import settings
from app.services.controle_api import registrar_contador
from oracledb import Connection
_TOKEN_FILE = os.path.join(tempfile.gettempdir(), "holmes_token_cache.json")
def _ler_token_arquivo() -> dict | None:
try:
with open(_TOKEN_FILE) as f:
return json.load(f)
except Exception:
return None
def _salvar_token_arquivo(token: dict):
try:
with open(_TOKEN_FILE, "w") as f:
json.dump(token, f)
except Exception:
pass
def _invalidar_token_usuario():
try:
os.remove(_TOKEN_FILE)
except Exception:
pass
async def token_usuario():
cached = _ler_token_arquivo()
if cached:
return cached
await asyncio.sleep(random.uniform(0, 0.5))
cached = _ler_token_arquivo()
if cached:
return cached
url = "https://app-api.holmesdoc.io/v1/session"
body = {"email": settings.holmes.usuario , "password": settings.holmes.senha}
headers = {"Content-Type": "application/json", "User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/150.0.0.0 Safari/537.36"}
async with httpx.AsyncClient() as client:
response = await client.post(url, json=body, headers=headers)
response.raise_for_status()
token = {"authorization": response.json()["token"]}
_salvar_token_arquivo(token)
return token
async def get_header() -> tuple[dict, str]:
if random.randint(1, 1) >= 2:
return {"api_token": settings.holmes.token_api}, "holmes"
header = await token_usuario()
return header, "holmes_user"
async def _holmes_request(
method: str, url: str, **kwargs
) -> tuple[httpx.Response, str]:
api_nome: str = "holmes"
for tentativa in range(2):
header, api_nome = await get_header()
async with httpx.AsyncClient() as client:
response = await getattr(client, method)(url, headers=header, **kwargs)
if (
response.status_code == 401
and api_nome == "holmes_user"
and tentativa == 0
):
_invalidar_token_usuario()
continue
return response, api_nome
raise RuntimeError("Holmes: falha de autenticação após retry")
def extrair_origem(origem: str):
"""Extrair origem da string do Holmes
Args:
origem (str): 1548 Entrada de Freio
Returns:
str: 1548
"""
return str(origem[:4])
def extrair_empresa(unidade: str):
"""Extrai a empresa da string inteira do holmes
Args:
unidade (str): Ex 10.1 Hyundai Teix. Freitas
Returns:
empresa: 10 | None
revenda: 1 | None
"""
match = re.search(r"^(\d+)\.(\d+)", unidade)
if match:
return match.group(1), match.group(2)
return None, None
def extrair_transacao(transacao: str) -> str:
"""Extrai apenas a parte transação do texto do Holmes
Args:
transacao (str): D15 Entrada de Nota
Returns:
str: D15
"""
return transacao[:3]
async def get_holmes_process(id_processo: str, conn: Connection):
"""
Busca um processo no Holmes. Centralizado para Peças e Despesas.
"""
url = f"https://app-api.holmesdoc.io/v1/processes/{id_processo}"
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
try:
response, api_nome = await _holmes_request("get", url)
response.raise_for_status()
status = response.status_code
sucesso = True
return response.json()
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro na API Holmes (ID {id_processo}): {e}")
return None
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao consultar Holmes: {e}")
return None
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/processes/{id}",
metodo="GET",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def get_holmes_process_details(id_processo: str, conn: Connection):
"""
Busca os detalhes (properties) de um processo no Holmes.
"""
url = f"https://app-api.holmesdoc.io/v1/processes/{id_processo}/details"
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
try:
response, api_nome = await _holmes_request("get", url)
response.raise_for_status()
status = response.status_code
sucesso = True
return response.json()
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro na API Holmes details (ID {id_processo}): {e}")
return None
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao consultar Holmes details: {e}")
return None
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/processes/{id}/details",
metodo="GET",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def get_holmes_history(id_processo: str, conn: Connection):
"""
Busca o historico no holmes (importante para pegar o id da ultima task, para avançar futuramente)
Args:
id_processo (str): id do processo no holmes
"""
url = f"https://app-api.holmesdoc.io/v1/processes/{id_processo}/history"
payload = {
"filters": [],
"page": 1,
"per_page": 100,
"sortBy": ["created_at", "desc"],
}
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
try:
response, api_nome = await _holmes_request("post", url, json=payload)
response.raise_for_status()
status = response.status_code
sucesso = True
return response.json()
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro na API Holmes (ID {id_processo}): {e}")
return None
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao consultar Holmes: {e}")
return None
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/processes/{id}/history",
metodo="POST",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def get_holmes_rateio(id_processo: str, conn: Connection):
url = f"https://app-api.holmesdoc.io/v1/processes/{id_processo}/tables/e124b2d0-ee14-11ef-95b4-25dee32fe73f/table_items?page=1&per_page=800"
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
try:
response, api_nome = await _holmes_request("get", url)
response.raise_for_status()
status = response.status_code
sucesso = True
return response.json()
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro na API Holmes (ID {id_processo}): {e}")
return None
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao consultar Holmes: {e}")
return None
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/processes/{id}/table/(rateio)",
metodo="GET",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def task_id_recente(id_processo: str, conn: Connection):
dados_tasks = await get_holmes_history(id_processo, conn)
if (
not dados_tasks
or "histories" not in dados_tasks
or not dados_tasks["histories"]
):
return None
mais_recente = max(dados_tasks["histories"], key=lambda x: x["created_at"])
return mais_recente["properties"]["task_id"]
async def task_mais_recente(id_processo: str, conn: Connection):
dados_tasks = await get_holmes_history(id_processo, conn)
if not dados_tasks or "histories" not in dados_tasks:
return None # Tratamento se a API falhar
# print(dados_tasks)
mais_recente = max(dados_tasks["histories"], key=lambda x: x["created_at"])
return mais_recente
async def historicos_task(id_processo: str, conn: Connection) -> dict | None:
"""_summary_
Args:
id_processo (str): Id do processo no holmes
Returns:
dict | None : dicionario do historico | None
"""
dados_tasks = await get_holmes_history(id_processo, conn)
if not dados_tasks or "histories" not in dados_tasks:
return None # Tratamento se a API falhar
return dados_tasks
async def buscar_processo(
conn: Connection,
chave: str | None = None,
fluxos: list[str] | bool = False,
ativos: bool = True,
payload: dict | bool = False,
) -> dict:
"""Obtem os processos que existem com a sua chave
Args:
chave (str): Chave principal a ser procurada, preferencialmente unica pfvr, ajuda ae po.
ativos (bool) Defaults to True
fluxos (list[str] | bool, optional): _description_. Defaults to False. se quer pegar de um fluxo específico ou geral. Padrão: Geral
Returns:
dict: _description_
"""
if not payload and chave is None:
raise ValueError(
"Chave é obrigatória caso o payload não seja enviado chefia, fica alerta ae rapa"
)
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
url = "https://app-api.holmesdoc.io/v2/search"
if not payload:
payload = {
"query": {
"from": 0,
"size": 200,
"context": "process",
"sort": "updated_at",
"order": "desc",
"groups": [
{
"match_all": True,
"terms": [
{
"value": f"{chave}",
"type": "match_phrase",
"field": "_content",
}
],
}
],
},
"trash": False,
"deleted_by_me": False,
}
try:
response, api_nome = await _holmes_request("post", url, json=payload)
response.raise_for_status()
status = response.status_code
sucesso = True
dados = response.json()
docs = dados.get("docs", [])
if ativos:
docs = [d for d in docs if d.get("status") != "canceled"]
if fluxos:
docs = [d for d in docs if d.get("name") in fluxos]
return {"status": True, "dados": {**dados, "docs": docs, "total": len(docs)}}
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro na API Holmes (ID {chave}): {e}")
return {"status": False, "error": e}
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao consultar Holmes: {e}")
return {"status": False, "error": e}
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/processes/{id}/search/por-chave",
metodo="POST",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def buscar_processo_por_chaves(
conn: Connection,
combinacoes: list[list[str]],
fluxos: list[str] | bool = False,
ativos: bool = True,
) -> dict:
"""Busca processos no Holmes usando combinações de termos.
Cada item de combinacoes é uma lista de valores que juntos identificam
um processo único (ex: [cnpj, numero_nf]). Cada combinação vira um group
separado na query.
Args:
combinacoes: Ex: [["03657256000164", "19"], ["698cd1c570fd0f8f5f8436a4"]]
ativos: Ignora processos cancelados. Padrão: True.
fluxos: Filtra por nome de fluxo. Padrão: False (todos).
"""
url = "https://app-api.holmesdoc.io/v2/search"
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
payload = {
"query": {
"from": 0,
"size": 200,
"context": "process",
"sort": "updated_at",
"order": "desc",
"groups": [
{
"match_all": True,
"terms": [
{"value": termo, "type": "match_phrase", "field": "_content"}
for termo in combinacao
],
}
for combinacao in combinacoes
],
},
"trash": False,
"deleted_by_me": False,
}
try:
response, api_nome = await _holmes_request("post", url, json=payload)
response.raise_for_status()
status = response.status_code
sucesso = True
dados = response.json()
docs = dados.get("docs", [])
if ativos:
docs = [d for d in docs if d.get("status") != "canceled"]
if fluxos:
docs = [d for d in docs if d.get("name") in fluxos]
return {"status": True, "dados": {**dados, "docs": docs, "total": len(docs)}}
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro na API Holmes (combinacoes {combinacoes}): {e}")
return {"status": False, "error": e}
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao consultar Holmes: {e}")
return {"status": False, "error": e}
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/processes/{id}/search/por-chaves",
metodo="POST",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def action(payload: dict, id_task: str, id_processo: str, conn: Connection):
url = f"https://app-api.holmesdoc.io/v1/tasks/{id_task}/action"
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
try:
async with httpx.AsyncClient() as client:
response = await client.post(
url, headers={"api_token": settings.holmes.token_api}, json=payload
)
response.raise_for_status()
status = response.status_code
sucesso = True
return True, response.json()
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro na API Holmes (ID {id_processo}): {e}")
return False, e
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao consultar Holmes: {e}")
return False, e
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/processes/{id}/action",
metodo="POST",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def cria_processo(
id_start: str,
payload: dict,
conn: Connection
) -> tuple[bool, dict | str]:
url = f"https://app-api.holmesdoc.io/v1/workflows/{id_start}/start"
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
try:
response, api_nome = await _holmes_request("post", url, json=payload)
response.raise_for_status()
status = response.status_code
sucesso = True
return True, response.json()
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro na API Holmes - Criar Processo ({payload}): {e}")
return False, str(e)
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao criar processo no Holmes: {e}")
return False, str(e)
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/workflows/{id}/start",
metodo="POST",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def enviar_documento(
id_processo: str,
arquivo: bytes,
nome_arquivo: str,
id_documento: str,
conn: Connection,
) -> tuple[bool, str | dict]:
task_id = await task_id_recente(id_processo, conn)
if not task_id:
return False, "Não foi possível obter a task mais recente do Holmes"
url = f"https://app-api.holmesdoc.io/v1/tasks/{task_id}/documents/{id_documento}"
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
try:
files = {"file": (nome_arquivo, arquivo, "application/pdf")}
async with httpx.AsyncClient() as client:
response = await client.post(
url, headers={"api_token": settings.holmes.token_api}, files=files
)
response.raise_for_status()
status = response.status_code
sucesso = True
return True, {"task_id": task_id}
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro ao enviar documento Holmes (ID {id_processo}): {e}")
return False, str(e)
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao enviar documento Holmes: {e}")
return False, str(e)
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/tasks/{id}/documents/{id_documento}",
metodo="POST",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
# Para testar as funcoes
# async def main():
# from app.infra.database import db_instance
# db_instance.create_pool()
# conn = db_instance.pool.acquire()
# try:
# print(await get_holmes_history('69e65ea83fad950fad5715ed', conn))
# finally:
# db_instance.pool.release(conn)
# if __name__ == '__main__':
# import asyncio
# asyncio.run(main())
async def cancela_processo(
id_processo: str, conn: Connection
) -> tuple[bool, str | dict]:
url = f"https://app-api.holmesdoc.io/v1/processes/{id_processo}/cancel"
payload = {"reason": "Erro na emissão, data de vencimento. Problema na Disal."}
inicio = time.perf_counter()
status = None
sucesso = False
erro = None
api_nome = "holmes"
try:
response, api_nome = await _holmes_request("put", url, json=payload)
response.raise_for_status()
status = response.status_code
sucesso = True
return True, response.json() if response.content else {
"mensagem": "processo cancelado"
}
except httpx.HTTPStatusError as e:
status = e.response.status_code
erro = str(e)
logging.error(f"Erro ao cancelar processo Holmes (ID {id_processo}): {e}")
return False, str(e)
except Exception as e:
erro = str(e)
logging.error(f"Erro inesperado ao cancelar processo Holmes: {e}")
return False, str(e)
finally:
duracao = (time.perf_counter() - inicio) * 1000
registrar_contador(
conn=conn,
api=api_nome,
endpoint="/v1/processes/{id}/cancel",
metodo="PUT",
status_code=status,
sucesso=sucesso,
duracao_ms=duracao,
erro=erro,
)
async def main():
# print(aaaa())
print(await token_usuario())
if __name__ == "__main__":
import asyncio
asyncio.run(main())

7
app/v1/api.py Normal file
View file

@ -0,0 +1,7 @@
from app.v1.compras import router as compras
from fastapi import APIRouter
api_router = APIRouter()
# Rota final: POST /v1/compras/dados
api_router.include_router(compras)

56
app/v1/compras.py Normal file
View file

@ -0,0 +1,56 @@
import logging
from app.config import settings
from app.schemas import Envelope, PayloadCompras
from app.seguranca import autorizar
from app.services import cripto
from fastapi import APIRouter, Depends, HTTPException, status
from pydantic import ValidationError
router = APIRouter(prefix="/compras", tags=["compras"])
# Carregada uma vez no import: parsear PEM a cada request e desperdicio.
_CHAVE = cripto.carregar_chave_privada(settings.vps.chave_privada)
@router.post("/dados", dependencies=[Depends(autorizar)])
async def receber_compras(envelope: Envelope) -> dict:
"""
Recebe o envelope cifrado da API da VPS, abre e valida.
Fase atual: so registra o que chegou. O INSERT no Oracle entra depois
quando entrar, tem que ser MERGE por id_processo, senao reentrega de
webhook duplica linha.
"""
try:
bruto = cripto.descriptografar(envelope.model_dump(), _CHAVE)
except Exception as e:
# Chave errada, envelope adulterado ou algoritmo diferente.
logging.error(f"Falha ao abrir envelope: {type(e).__name__}: {e}")
raise HTTPException(
status_code=status.HTTP_400_BAD_REQUEST, detail="envelope invalido"
)
try:
dados = PayloadCompras.model_validate(bruto)
except ValidationError as e:
# Decifrou mas o formato mudou: os dois lados sairam de sincronia.
logging.error(f"Payload fora do contrato: {e}")
raise HTTPException(
status_code=status.HTTP_422_UNPROCESSABLE_ENTITY,
detail="payload nao bate com o contrato esperado",
)
logging.info(
f"Compra recebida | processo={dados.id_processo} "
f"protocolo={dados.protocolo} cnpj={dados.cnpj} "
f"pedido={dados.pedido_linx} total={dados.valor_total} "
f"parcelas={len(dados.parcelas)}"
)
return {
"status": "recebido",
"id_processo": dados.id_processo,
"parcelas": len(dados.parcelas),
"persistido": False, # vira True quando o Oracle entrar
}