diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..79f56b5 --- /dev/null +++ b/.env.example @@ -0,0 +1,11 @@ +DB_URL=postgresql://usuario:senha@localhost:5432/nome_do_banco + +CLOUDINARY_CLOUD_NAME= +CLOUDINARY_API_KEY= +CLOUDINARY_API_SECRET= + +ACTA_PG_API_URL=https://acta-pg-api.onrender.com/api/v1 +ACTA_IA_TOKEN=defina-um-token-interno + +REDIS_HOST=localhost +REDIS_PORT=6379 \ No newline at end of file diff --git a/.github/workflows/build-image.yml b/.github/workflows/build-image.yml new file mode 100644 index 0000000..961e30e --- /dev/null +++ b/.github/workflows/build-image.yml @@ -0,0 +1,74 @@ +name: build-image + +on: + push: + branches: [main] + +permissions: + contents: read + packages: write + pull-requests: read + +jobs: + publicar: + runs-on: ubuntu-latest + + steps: + - name: baixar código + uses: actions/checkout@v4 + + - name: preparar docker builds + uses: docker/setup-buildx-action@v3 + + - name: entrar no GHCR + uses: docker/login-action@v3 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ secrets.GITHUB_TOKEN }} + + - name: construir e publicar imagem + uses: docker/build-push-action@v6 + with: + context: . + push: true + tags: ghcr.io/appacta/acta-import-api:sha-${{ github.sha }} + labels: | + org.opencontainers.image.source=https://github.com/AppActa/acta-import-api + org.opencontainers.image.revision=${{ github.sha }} + + - name: baixar manifests do platform + uses: actions/checkout@v4 + with: + repository: AppActa/acta-platform + ref: main + token: ${{ secrets.ACTA_PLATFORM_TOKEN }} + + - name: atualizar imagem no platform + shell: bash + env: + GH_TOKEN: ${{ secrets.GITHUB_TOKEN }} + COMMIT_AUTHOR: ${{ github.event.head_commit.author.name }} + PUSHED_AT: ${{ github.event.head_commit.timestamp }} + run: | + pull_request=$(gh api "repos/${GITHUB_REPOSITORY}/commits/${GITHUB_SHA}/pulls" --jq '.[0].html_url' 2>/dev/null || true) + pull_request=${pull_request:-"não identificada"} + + commit_author="${COMMIT_AUTHOR}" + commit_author=${commit_author:-"não identificado"} + + pushed_at="${PUSHED_AT}" + pushed_at=${pushed_at:-"não informado"} + + sed -i "s|^\([[:space:]]*image: ghcr.io/appacta/acta-import-api\).*$|\1:sha-${GITHUB_SHA}|" apps/import-api.yaml + git diff --quiet -- apps/import-api.yaml && exit 0 + + git config user.name "github-actions[bot]" + git config user.email "41898282+github-actions[bot]@users.noreply.github.com" + git add apps/import-api.yaml + git commit -m "chore: atualizar imagem de acta-import-api" -m "PR: ${pull_request} + SHA: ${GITHUB_SHA} + Autor do commit: ${commit_author} + Usuário do push: ${GITHUB_ACTOR} + Data do push: ${pushed_at}" + git push origin main \ No newline at end of file diff --git a/.github/workflows/validacao.yml b/.github/workflows/validacao.yml new file mode 100644 index 0000000..a12be50 --- /dev/null +++ b/.github/workflows/validacao.yml @@ -0,0 +1,34 @@ +name: validacao + +on: + pull_request: + branches: [main] + push: + branches: [main] + +permissions: + contents: read + +jobs: + validar: + runs-on: ubuntu-latest + + steps: + - uses: actions/checkout@v4 + + - uses: actions/setup-python@v5 + with: + python-version: '3.12' + cache: pip + + - name: instalar dependencias + run: | + python -m pip install --upgrade pip + python -m pip install -r requirements.txt + python -m pip install pytest==9.1.1 + + - name: validar sintaxe + run: python -m compileall -q main.py worker.py app + + - name: executar testes + run: python -m pytest -q -p no:cacheprovider \ No newline at end of file diff --git a/.gitignore b/.gitignore index e69de29..d276bf8 100644 --- a/.gitignore +++ b/.gitignore @@ -0,0 +1,4 @@ +*.env +docs/ +*.pyc +tests/.*/ \ No newline at end of file diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..9678d0b --- /dev/null +++ b/Dockerfile @@ -0,0 +1,16 @@ +FROM python:3.12-slim + +WORKDIR /app + +ENV PYTHONDONTWRITEBYTECODE=1 \ + PYTHONUNBUFFERED=1 + +COPY requirements.txt . +RUN python -m pip install --no-cache-dir -r requirements.txt + +COPY main.py worker.py ./ +COPY app ./app + +EXPOSE 8000 + +CMD ["sh", "-c", "exec uvicorn main:app --host 0.0.0.0 --port ${PORT:-8000}"] \ No newline at end of file diff --git a/README.md b/README.md index 1ca57a0..6e3132b 100644 --- a/README.md +++ b/README.md @@ -1 +1,282 @@ -# acta-template-repository \ No newline at end of file +# 📎 ACTA Import API + +API responsável pelo recebimento e processamento assíncrono de anexos utilizados nos ciclos PDCA do ACTA. + +O serviço valida o arquivo recebido, registra seus metadados no PostgreSQL, salva uma cópia temporária e envia um job para processamento no Redis. O worker realiza o upload para o Cloudinary e atualiza o status do anexo no banco. + +## 📌 Visão geral + +O `acta-import-api` atua como um serviço especializado do ecossistema ACTA: + +1. recebe o arquivo e os dados do anexo; +2. consulta o usuário autenticado no `acta-pg-api`; +3. valida tamanho, extensão e conteúdo; +4. registra os metadados em `pdca.anexo`; +5. enfileira o processamento no Redis; +6. envia o arquivo ao Cloudinary pelo worker; +7. atualiza o status e o caminho do arquivo no PostgreSQL. + +O processamento do upload é assíncrono. Portanto, uma resposta `202 Accepted` confirma o recebimento e o enfileiramento, mas não confirma que o upload no Cloudinary já terminou. + +## 🧰 Redis e fila de processamento + +O Redis funciona como intermediário entre a API e o worker. A API não envia o arquivo diretamente ao Cloudinary durante a requisição; depois de salvar o arquivo temporário e os metadados, publica um `UploadJob` na fila `uploads`. + +O worker executa continuamente essa fila, recupera o job, envia o arquivo ao Cloudinary e atualiza o registro no PostgreSQL. Assim, o processamento pesado acontece fora da requisição HTTP e o endpoint pode responder rapidamente com `202 Accepted`. + +O Redis armazena a fila e o estado necessário para o RQ controlar a execução. Se o processamento falhar, o job possui até três tentativas, com intervalos de 10, 30 e 60 segundos. O arquivo temporário é mantido durante as tentativas e removido quando o processamento termina com sucesso ou falha definitiva. + +O Redis precisa estar acessível tanto pela API, que enfileira os jobs, quanto pelo worker, que os consome. Em ambientes com autenticação, configure também `REDIS_PASSWORD` no ambiente da aplicação e do worker. + +## ✨ Funcionalidades + +- Upload de arquivos para anexos de ciclos PDCA. +- Validação de arquivos PDF, DOCX, PPTX, XLSX e imagens. +- Validação específica para TXT, CSV e SVG. +- Limite de 10 MB por arquivo. +- Autenticação por Bearer Token delegada ao `acta-pg-api`. +- Registro do usuário e da empresa a partir do contexto autenticado. +- Processamento assíncrono com Redis Queue (RQ). +- Upload de arquivos para o Cloudinary. +- Atualização dos status `PROCESSANDO`, `ATIVO` e `ERRO`. + + + +## 🛠️ Tecnologias + + +| Tecnologia | Uso | +| ------------ | ------------------------------------ | +| Python | Linguagem da aplicação | +| FastAPI | API HTTP | +| PostgreSQL | Metadados dos anexos | +| Redis | Fila de processamento | +| RQ | Execução dos jobs assíncronos | +| Cloudinary | Armazenamento dos arquivos | +| `filetype` | Identificação do tipo por assinatura | +| `defusedxml` | Validação segura de SVG | +| Pytest | Testes unitários | + + + + +## ✅ Pré-requisitos + +- Python 3.11 ou superior. +- PostgreSQL acessível com a tabela `pdca.anexo` já criada. +- Redis acessível para a fila `uploads`. +- Conta e credenciais do Cloudinary. +- `acta-pg-api` acessível para autenticação e consulta do usuário. + +As dependências Python usadas pela aplicação devem estar instaladas no ambiente virtual do projeto. + +## ⚙️ Configuração + +Crie um arquivo `.env` local a partir de `.env.example`: + +```env +DB_URL=postgresql://usuario:senha@localhost:5432/nome_do_banco + +CLOUDINARY_CLOUD_NAME=seu_cloud_name +CLOUDINARY_API_KEY=sua_api_key +CLOUDINARY_API_SECRET=seu_api_secret + +ACTA_PG_API_URL=https://acta-pg-api.onrender.com/api/v1 +ACTA_IA_TOKEN=defina-um-token-interno + +REDIS_HOST=localhost +REDIS_PORT=6379 +``` + +Nunca versione senhas, tokens ou credenciais do Cloudinary. O arquivo `.env.example` deve conter somente valores de referência. + +## 🚀 Execução local + +Inicie a API: + +```bash +uvicorn main:app --reload +``` + +Por padrão, o servidor fica disponível em `http://localhost:8000`. + +Em outro terminal, inicie o worker da fila: + +```bash +python worker.py +``` + +A API e o worker precisam estar em execução para que o ciclo completo de importação seja processado. + +## 📚 Documentação + +Com a API em execução: + +- Swagger UI: `http://localhost:8000/docs` +- ReDoc: `http://localhost:8000/redoc` +- OpenAPI: `http://localhost:8000/openapi.json` + + + +## 🔌 Endpoint principal + + + +### Criar anexo + +```http +POST /anexos +Authorization: Bearer +Content-Type: multipart/form-data +``` + +Campos do formulário: + + +| Campo | Tipo | Obrigatório | Descrição | +| ----------- | ------- | ----------- | -------------------------------- | +| `arquivo` | arquivo | Sim | Arquivo que será importado | +| `id_ciclo` | inteiro | Sim | Identificador do ciclo PDCA | +| `id_origem` | inteiro | Sim | Identificador da origem do anexo | +| `categoria` | enum | Sim | Categoria do anexo | +| `descricao` | texto | Não | Descrição do arquivo | + + +Categorias disponíveis: + +`TREINAMENTO`, `PLANO_ACAO`, `CAUSA_RAIZ`, `PROBLEMA`, `META`, `RELATORIO`, `LICAO_APRENDIDA`, `FORMULARIO`, `EVIDENCIA` e `OUTRO`. + +Exemplo: + +```bash +curl -X POST http://localhost:8000/anexos \ + -H "Authorization: Bearer " \ + -F "arquivo=@./evidencia.pdf" \ + -F "id_ciclo=7" \ + -F "id_origem=9" \ + -F "categoria=EVIDENCIA" \ + -F "descricao=Evidência da verificação do resultado" +``` + +Resposta de recebimento: + +```json +{ + "id": 42, + "status": "PROCESSANDO" +} +``` + +Essa resposta não representa a confirmação final do Cloudinary. O status final deve ser consultado no fluxo de anexos do ACTA ou diretamente no banco. + +### Criar anexo pela IA + +```http +POST /anexos/ia +Authorization: Bearer +X-Acta-Usuario-Id: +Content-Type: multipart/form-data +``` + +A IA usa uma credencial própria de serviço e informa apenas o usuário em nome de quem está realizando o upload. A API valida se esse usuário está ativo, obtém a empresa diretamente de `public.usuario_sistema` e verifica se o ciclo pertence à mesma empresa. O `id_empresa` e o `firebase_uid` não são recebidos nessa rota. + +Os campos do formulário e a resposta `202 Accepted` são os mesmos de `POST /anexos`. + +### Consultar processamento do anexo + +Para usuários autenticados: + +```http +GET /anexos/{id_anexo} +Authorization: Bearer +``` + +Para a IA: + +```http +GET /anexos/ia/{id_anexo} +Authorization: Bearer +X-Acta-Usuario-Id: +``` + +A resposta informa o status atual, se o envio terminou com sucesso e a URL gerada pelo Cloudinary: + +```json +{ + "id": 42, + "status": "ATIVO", + "enviado": true, + "url": "https://res.cloudinary.com/..." +} +``` + +## ⚠️ Erros + +O formato confirmado é: + +```json +{ + "mensagens": ["campo: mensagem de validação"], + "httpStatus": 400, + "timestamp": "2026-09-17T10:35:00" +} +``` + +## 📦 Tipos de arquivo + +São aceitos os seguintes formatos: + +- PDF: `.pdf` +- Documentos: `.docx`, `.pptx`, `.xlsx` +- Imagens: `.jpg`, `.jpeg`, `.png`, `.webp`, `.gif`, `.bmp`, `.tif` +- Texto: `.txt`, `.csv`, `.svg` + +Arquivos de texto são validados por conteúdo, incluindo UTF-8, ausência de byte nulo e estrutura básica de CSV ou SVG. + +## 🔄 Processamento assíncrono + +O endpoint grava o anexo com status `PROCESSANDO` e publica um `UploadJob` na fila `uploads`. + +O worker: + +- atualiza o status do anexo; +- envia o arquivo ao Cloudinary; +- salva o `secure_url` retornado; +- altera o status para `ATIVO` em caso de sucesso; +- altera o status para `ERRO` em caso de falha definitiva; +- remove o arquivo temporário após o processamento concluído. + +Falhas de enfileiramento retornam `503`. Arquivos vazios retornam `400` e arquivos acima de 10 MB retornam `413`. + +## 🧪 Testes + +Execute os testes unitários com: + +```bash +python -m pytest -q -p no:cacheprovider +``` + +Os testes usam dublês para banco, Redis e Cloudinary. Portanto, a suíte valida o comportamento interno da aplicação, mas não substitui uma validação de integração com os serviços reais. + +## 📁 Estrutura + +```text +. +├── app/ +│ ├── database/ # Conexões PostgreSQL e Redis +│ ├── models/ # Enums e modelos relacionados +│ ├── routes/ # Endpoints HTTP +│ ├── services/ # Autenticação, fila, jobs e Cloudinary +│ ├── utils/ # Validação de arquivos +│ └── schemas.py # Contratos de entrada e jobs +├── tests/ # Testes unitários +├── main.py # Inicialização da API +└── worker.py # Inicialização do worker RQ +``` + + + +## 🤝 Links e autoria + +- [Repositório](https://github.com/AppActa/acta-import-api) · [Licença MIT](LICENSE) · `acta.institutojef@gmail.com` +- Contribuições: use *issues* e *pull requests*; há um [PULL_REQUEST_TEMPLATE.md](PULL_REQUEST_TEMPLATE.md). \ No newline at end of file diff --git a/app/database/postgres.py b/app/database/postgres.py new file mode 100644 index 0000000..f88f4d2 --- /dev/null +++ b/app/database/postgres.py @@ -0,0 +1,32 @@ +from contextlib import contextmanager +from fastapi import HTTPException +from dotenv import load_dotenv +from os import getenv +from psycopg2 import connect + +load_dotenv() + +DB_URL = getenv('DB_URL') + +def _conectar(): + return connect(DB_URL) + +def get_conn(): + try: + conn = _conectar() + except Exception: + raise HTTPException(status_code=503, detail='Erro no banco de dados') + + try: + yield conn + finally: + conn.close() + +@contextmanager # permite usar def com with +def conn_worker(): + conn = _conectar() + + try: + yield conn + finally: + conn.close() \ No newline at end of file diff --git a/app/database/redis.py b/app/database/redis.py new file mode 100644 index 0000000..1fbfcdb --- /dev/null +++ b/app/database/redis.py @@ -0,0 +1,23 @@ +from fastapi import HTTPException +from redis import Redis, ConnectionPool +from redis.exceptions import RedisError +from dotenv import load_dotenv +from os import getenv +from functools import lru_cache + +load_dotenv() + +REDIS_HOST = getenv('REDIS_HOST') +REDIS_PORT = getenv('REDIS_PORT') +REDIS_DB = getenv('REDIS_DB') +REDIS_PASSWORD = getenv('REDIS_PASSWORD') + +@lru_cache # apenas um pool por sessão +def get_pool() -> ConnectionPool: + return ConnectionPool(host=REDIS_HOST, port=REDIS_PORT, db=REDIS_DB, password=REDIS_PASSWORD, max_connections=20, retry_on_timeout=True) + +def get_conn() -> Redis: + try: + return Redis(connection_pool=get_pool()) + except RedisError as re: + raise HTTPException(status_code=503, detail='Redis indisponível') \ No newline at end of file diff --git a/app/models/anexo.py b/app/models/anexo.py new file mode 100644 index 0000000..97872ad --- /dev/null +++ b/app/models/anexo.py @@ -0,0 +1,23 @@ +from pydantic import BaseModel, Field +from datetime import datetime +from app.models.enums import TipoArquivo, Categoria, TipoImagem, Status + +class AnexoRequest(BaseModel): + id_empresa: int + id_ciclo: int + criado_por: int + id_origem: int + nome_arquivo: str + tipo_arquivo: TipoArquivo | TipoImagem + tamanho_arquivo: int + bucket_arquivo: str + caminho_arquivo: str | None = None + categoria: Categoria + descricao: str | None = None + +class Anexo(AnexoRequest): + id: int + status: Status + criado_em: datetime = Field(default_factory=datetime.now) + atualizado_em: datetime | None = None + excluido_em: datetime | None = None \ No newline at end of file diff --git a/app/models/enums.py b/app/models/enums.py new file mode 100644 index 0000000..22cbb77 --- /dev/null +++ b/app/models/enums.py @@ -0,0 +1,51 @@ +from enum import Enum + +class Categoria(Enum): + TREINAMENTO = 'TREINAMENTO' + PLANO_ACAO = 'PLANO_ACAO' + CAUSA_RAIZ = 'CAUSA_RAIZ' + PROBLEMA = 'PROBLEMA' + META = 'META' + RELATORIO = 'RELATORIO' + LICAO_APRENDIDA = 'LICAO_APRENDIDA' + FORMULARIO = 'FORMULARIO' + EVIDENCIA = 'EVIDENCIA' + OUTRO = 'OUTRO' + +class Status(Enum): + PROCESSANDO = 'PROCESSANDO' + ATIVO = 'ATIVO' + ERRO = 'ERRO' + EXCLUIDO = 'EXCLUIDO' + +class TipoImagem(Enum): + JPEG = 'image/jpeg' + PNG = 'image/png' + WEBP = 'image/webp' + GIF = 'image/gif' + SVG = 'image/svg+xml' + BMP = 'image/bmp' + TIFF = 'image/tiff' + +class TipoArquivo(Enum): + PDF = 'application/pdf' + DOCX = 'application/vnd.openxmlformats-officedocument.wordprocessingml.document' + PPTX = 'application/vnd.openxmlformats-officedocument.presentationml.presentation' + XLSX = 'application/vnd.openxmlformats-officedocument.spreadsheetml.sheet' + CSV = 'text/csv' + TXT = 'text/plain' + +# não incluso svg, txt e csv porque são texto puro e devem ser validados de outra forma +EXTENSOES_PERMITIDAS = set({ + TipoImagem.JPEG: 'jpg', + TipoImagem.JPEG: 'jpeg', + TipoImagem.PNG: 'png', + TipoImagem.WEBP: 'webp', + TipoImagem.GIF: 'gif', + TipoImagem.BMP: 'bmp', + TipoImagem.TIFF: 'tif', + TipoArquivo.PDF: 'pdf', + TipoArquivo.DOCX: 'docx', + TipoArquivo.PPTX: 'pptx', + TipoArquivo.XLSX: 'xlsx' +}.values()) \ No newline at end of file diff --git a/app/routes/anexos.py b/app/routes/anexos.py new file mode 100644 index 0000000..86b8405 --- /dev/null +++ b/app/routes/anexos.py @@ -0,0 +1,24 @@ +from fastapi import APIRouter, Depends, File, UploadFile + +from app.database.postgres import get_conn +from app.schemas import AnexoForm, UsuarioAutenticado +from app.services.anexos import consultar_anexo, processar_anexo +from app.services.autenticacao import obter_contexto_ia, obter_usuario_autenticado + +router = APIRouter(prefix='/anexos', tags=['anexos']) + +@router.post('', status_code=202) +async def post_anexo(arquivo: UploadFile = File(...), dados: AnexoForm = Depends(), conn=Depends(get_conn), usuario: UsuarioAutenticado = Depends(obter_usuario_autenticado)): + return await processar_anexo(arquivo, dados, conn, usuario) + +@router.post('/ia', status_code=202) +async def post_anexo_ia(arquivo: UploadFile = File(...), dados: AnexoForm = Depends(), conn=Depends(get_conn), usuario: UsuarioAutenticado = Depends(obter_contexto_ia)): + return await processar_anexo(arquivo, dados, conn, usuario) + +@router.get('/ia/{id_anexo}') +def get_anexo_ia(id_anexo: int, conn=Depends(get_conn), usuario: UsuarioAutenticado = Depends(obter_contexto_ia)): + return consultar_anexo(conn, id_anexo, usuario.id_empresa) + +@router.get('/{id_anexo}') +def get_anexo(id_anexo: int, conn=Depends(get_conn), usuario: UsuarioAutenticado = Depends(obter_usuario_autenticado)): + return consultar_anexo(conn, id_anexo, usuario.id_empresa) diff --git a/app/schemas.py b/app/schemas.py new file mode 100644 index 0000000..efe6a74 --- /dev/null +++ b/app/schemas.py @@ -0,0 +1,27 @@ +from fastapi import Form +from pydantic import BaseModel, Field +from app.models.enums import Categoria + +class AnexoForm: + def __init__( + self, + id_ciclo: int = Form(...), + id_origem: int = Form(...), + categoria: Categoria = Form(...), + descricao: str | None = Form(None) + ): + self.id_ciclo = id_ciclo + self.id_origem = id_origem + self.categoria = categoria + self.descricao = descricao + +class UploadJob(BaseModel): + id_anexo: int + id_usuario: int + conteudo: bytes + nome_original: str + extensao: str + +class UsuarioAutenticado(BaseModel): + id_usuario: int = Field(alias='idUsuario') + id_empresa: int = Field(alias='idEmpresa') diff --git a/app/services/anexos.py b/app/services/anexos.py new file mode 100644 index 0000000..d93e348 --- /dev/null +++ b/app/services/anexos.py @@ -0,0 +1,99 @@ +from pathlib import Path +from uuid import uuid4 + +from fastapi import HTTPException, UploadFile +from rq import Retry + +from app.models.enums import Status +from app.schemas import AnexoForm, UploadJob, UsuarioAutenticado +from app.services.fila import fila_uploads +from app.services.jobs import processar_upload +from app.utils.validacao import validar_anexo + +DIR_TEMP = Path('/tmp/uploads') +DIR_TEMP.mkdir(parents=True, exist_ok=True) + +async def processar_anexo(arquivo: UploadFile, dados: AnexoForm, conn, usuario: UsuarioAutenticado): + _validar_ciclo_empresa(conn, dados.id_ciclo, usuario.id_empresa) + extensao = await validar_anexo(arquivo) + + id_anexo = _inserir_metadados(conn, arquivo, dados, extensao, usuario) + caminho_temp = DIR_TEMP / f'{id_anexo}_{uuid4().hex}{Path(arquivo.filename).suffix}' + + try: + caminho_temp.write_bytes(await arquivo.read()) + except Exception: + _atualizar_status_erro(conn, id_anexo, usuario.id_usuario) + raise HTTPException(500, 'Falha ao salvar arquivo temporário') + + try: + fila_uploads.enqueue( + processar_upload, UploadJob( + id_anexo=id_anexo, + id_usuario=usuario.id_usuario, + conteudo=caminho_temp.read_bytes(), + nome_original=arquivo.filename, + extensao=extensao), + retry=Retry(max=3, interval=[10, 30, 60]), result_ttl=3600, failure_ttl=86400 + ) + except Exception: + _atualizar_status_erro(conn, id_anexo, usuario.id_usuario) + caminho_temp.unlink(missing_ok=True) + raise HTTPException(503, 'Não foi possível enfileirar o processamento') + + caminho_temp.unlink(missing_ok=True) + return {'id': id_anexo, 'status': Status.PROCESSANDO.value} + +def consultar_anexo(conn, id_anexo: int, id_empresa: int) -> dict: + with conn.cursor() as cursor: + cursor.execute( + """ + SELECT id, status, caminho_arquivo + FROM pdca.anexo + WHERE id = %s AND id_empresa = %s AND excluido_em IS NULL + """, + (id_anexo, id_empresa) + ) + anexo = cursor.fetchone() + + if anexo is None: + raise HTTPException(404, 'Anexo não encontrado') + + return { + 'id': anexo[0], + 'status': anexo[1], + 'enviado': anexo[1] == Status.ATIVO.value and bool(anexo[2]), + 'url': anexo[2] + } + +def _validar_ciclo_empresa(conn, id_ciclo: int, id_empresa: int) -> None: + with conn.cursor() as cursor: + cursor.execute( + 'SELECT 1 FROM pdca.ciclo WHERE id = %s AND id_empresa = %s', + (id_ciclo, id_empresa) + ) + if cursor.fetchone() is None: + raise HTTPException(403, 'Ciclo não pertence à empresa do usuário') + +def _inserir_metadados(conn, arquivo: UploadFile, dados: AnexoForm, extensao: str, usuario: UsuarioAutenticado) -> int: + with conn.cursor() as cursor: + cursor.execute("SELECT set_config('app.current_user_id', %s, true)", (str(usuario.id_usuario),)) + cursor.execute(""" + INSERT INTO pdca.anexo (id_empresa, id_ciclo, criado_por, id_origem, nome_arquivo, tipo_arquivo, tamanho_arquivo, bucket_arquivo, caminho_arquivo, categoria, descricao, status) + VALUES (%s, %s, %s, %s, %s, %s, %s, %s, NULL, %s, %s, %s) + RETURNING id + """, + (usuario.id_empresa, dados.id_ciclo, usuario.id_usuario, dados.id_origem, arquivo.filename, arquivo.content_type or extensao, arquivo.size, 'cloudinary', dados.categoria.value, dados.descricao, Status.PROCESSANDO.value)) + + id_anexo = cursor.fetchone()[0] + conn.commit() + return id_anexo + +def _atualizar_status_erro(conn, id_anexo: int, id_usuario: int) -> None: + with conn.cursor() as cursor: + cursor.execute( + "SELECT set_config('app.current_user_id', %s, true)", + (str(id_usuario),) + ) + cursor.execute('UPDATE pdca.anexo SET status = %s WHERE id = %s', (Status.ERRO.value, id_anexo)) + conn.commit() diff --git a/app/services/autenticacao.py b/app/services/autenticacao.py new file mode 100644 index 0000000..73b6588 --- /dev/null +++ b/app/services/autenticacao.py @@ -0,0 +1,61 @@ +from hmac import compare_digest +from os import getenv +from httpx import AsyncClient, RequestError +from fastapi import Depends, Header, HTTPException, Security +from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer + +from app.database.postgres import get_conn +from app.schemas import UsuarioAutenticado + +ACTA_PG_API_URL = getenv('ACTA_PG_API_URL') +bearer = HTTPBearer() + +async def obter_contexto_ia( + credenciais: HTTPAuthorizationCredentials = Security(bearer), + id_usuario: int = Header(..., alias='X-Acta-Usuario-Id', gt=0), + conn=Depends(get_conn) +) -> UsuarioAutenticado: + token_ia = getenv('ACTA_IA_TOKEN') + if not token_ia or not compare_digest(credenciais.credentials, token_ia): + raise HTTPException(401, 'Token da IA inválido') + + with conn.cursor() as cursor: + cursor.execute( + """ + SELECT id, id_empresa + FROM public.usuario_sistema + WHERE id = %s AND status = 'ATIVO' + """, + (id_usuario,) + ) + usuario = cursor.fetchone() + + if usuario is None: + raise HTTPException(403, 'Usuário da IA não autorizado') + + return UsuarioAutenticado(idUsuario=usuario[0], idEmpresa=usuario[1]) + +async def obter_usuario_autenticado(credenciais: HTTPAuthorizationCredentials = Security(bearer)) -> UsuarioAutenticado: + headers = {'Authorization': f'Bearer {credenciais.credentials}'} + + try: + async with AsyncClient(timeout=60) as cliente: + health = await cliente.get(f'{ACTA_PG_API_URL}/health') + + if health.is_error: + raise HTTPException(503, 'acta-pg-api indisponível') + + resposta = await cliente.get(f'{ACTA_PG_API_URL}/me', headers=headers) + except RequestError: + raise HTTPException(503, 'Não foi possível consultar acta-pg-api') + + if resposta.status_code in (401, 403): + raise HTTPException(401, 'Bearer token inválido') + + if resposta.is_error: + raise HTTPException(502, 'Falha ao consultar usuário autenticado') + + try: + return UsuarioAutenticado.model_validate(resposta.json()) + except (ValueError, TypeError): + raise HTTPException(502, 'Resposta inválida de acta-pg-api') \ No newline at end of file diff --git a/app/services/cloudinary.py b/app/services/cloudinary.py new file mode 100644 index 0000000..64179b2 --- /dev/null +++ b/app/services/cloudinary.py @@ -0,0 +1,33 @@ +from cloudinary import config +from cloudinary.uploader import upload +from os import getenv +from dotenv import load_dotenv + +load_dotenv() + +CLOUD_NAME = getenv('CLOUDINARY_CLOUD_NAME') +API_KEY = getenv('CLOUDINARY_API_KEY') +API_SECRET = getenv('CLOUDINARY_API_SECRET') + +config(cloud_name=CLOUD_NAME, api_key=API_KEY, api_secret=API_SECRET, secure=True) + +# reforçar tipo raw, cloudinary não sabe lidar com esses arquivos +EXTENSOES_RAW = {'docx', 'pptx', 'xlsx', 'csv', 'txt'} + +# erros serão propagados para a fila no redis +def enviar_cloudinary(caminho_arquivo: str, nome_original: str, extensao: str, pasta: str = 'anexos', public_id: str | None = None) -> dict: + resultado = upload( + caminho_arquivo, + folder=pasta, + public_id=public_id, + resource_type='raw' if extensao in EXTENSOES_RAW else 'auto', + filename=nome_original, + use_filename=True, + unique_filename=True + ) + + return { + 'url': resultado['secure_url'], + 'public_id': resultado['public_id'], + 'resource_type': resultado['resource_type'] + } \ No newline at end of file diff --git a/app/services/fila.py b/app/services/fila.py new file mode 100644 index 0000000..0830998 --- /dev/null +++ b/app/services/fila.py @@ -0,0 +1,8 @@ +from rq import Queue +from redis import Redis +from app.database.redis import get_pool + +conn = Redis(connection_pool=get_pool()) + +# 5 minutos de tolerância do job +fila_uploads = Queue('uploads', connection=conn, default_timeout=300) \ No newline at end of file diff --git a/app/services/jobs.py b/app/services/jobs.py new file mode 100644 index 0000000..ab67d07 --- /dev/null +++ b/app/services/jobs.py @@ -0,0 +1,51 @@ +from pathlib import Path +from tempfile import NamedTemporaryFile +from rq import get_current_job +from app.database.postgres import conn_worker +from app.models.enums import Status +from app.schemas import UploadJob +from app.services.cloudinary import enviar_cloudinary + +def processar_upload(job: UploadJob) -> None: + caminho_temp = None + with conn_worker() as conn: + try: + _configurar_usuario(conn, job.id_usuario) + + with NamedTemporaryFile(suffix=f'.{job.extensao}', delete=False) as arquivo_temp: + arquivo_temp.write(job.conteudo) + caminho_temp = arquivo_temp.name + _atualizar_status(conn, job.id_anexo, job.id_usuario, Status.PROCESSANDO) + + resultado = enviar_cloudinary(caminho_temp, job.nome_original, job.extensao, public_id=f'anexos/{job.id_anexo}') + + _atualizar_sucesso(conn, job.id_anexo, job.id_usuario, resultado) + except Exception: + rq_job = get_current_job() + ultima_tentativa = rq_job is None or not rq_job.retries_left + + if ultima_tentativa: + _atualizar_status(conn, job.id_anexo, job.id_usuario, Status.ERRO) + raise + finally: + if caminho_temp: + _limpar_temp(caminho_temp) + +def _configurar_usuario(conn, id_usuario: int) -> None: + with conn.cursor() as cursor: + cursor.execute("SELECT set_config('app.current_user_id', %s, true)", (str(id_usuario),),) + +def _atualizar_status(conn, id_anexo: int, id_usuario: int, status: Status) -> None: + _configurar_usuario(conn, id_usuario) + with conn.cursor() as cursor: + cursor.execute('UPDATE pdca.anexo SET status = %s WHERE id = %s', (status.value, id_anexo)) + conn.commit() + +def _atualizar_sucesso(conn, id_anexo: int, id_usuario: int, resultado: dict) -> None: + _configurar_usuario(conn, id_usuario) + with conn.cursor() as cursor: + cursor.execute('UPDATE pdca.anexo SET status = %s, caminho_arquivo = %s WHERE id = %s', (Status.ATIVO.value, resultado['url'], id_anexo)) + conn.commit() + +def _limpar_temp(caminho_arquivo: str) -> None: + Path(caminho_arquivo).unlink(missing_ok=True) \ No newline at end of file diff --git a/app/utils/sem_magic_bytes.py b/app/utils/sem_magic_bytes.py new file mode 100644 index 0000000..938d3cf --- /dev/null +++ b/app/utils/sem_magic_bytes.py @@ -0,0 +1,51 @@ +from codecs import getincrementaldecoder +from csv import Sniffer, Error +from fastapi import UploadFile, HTTPException +from defusedxml.ElementTree import fromstring + +async def _validar_text(arquivo: UploadFile) -> str: + # lê apenas primeiros 4kb + header = await arquivo.read(4096) + await arquivo.seek(0) + + # verifica se existe byte nulo (exclusivos de arquivos que não são texto puro) + if b'\x00' in header: + raise HTTPException(400, 'Arquivo não é um TXT puro') + + try: + # se houver caractere incompleto no final do chunk não lança exceção + # (pode ter sido fatiado ao pegar os primeiros 4kb) + decoder = getincrementaldecoder('utf-8')() + return decoder.decode(header, final=False) + except UnicodeDecodeError: + raise HTTPException(400, 'Arquivo TXT não está em UTF-8') + +async def validar_plain(arquivo: UploadFile) -> None: + await _validar_text(arquivo) + + +async def validar_csv(arquivo: UploadFile) -> None: + amostra = await _validar_text(arquivo) + + try: + Sniffer().sniff(amostra) + except Error: + raise HTTPException(400, 'Arquivo não é um CSV puro') + +async def validar_svg(arquivo: UploadFile) -> None: + conteudo = await arquivo.read() + await arquivo.seek(0) + + try: + root = fromstring(conteudo) + except Exception: + raise HTTPException(400, 'Arquivo não é um SVG puro') + + if not root.tag.endswith('svg'): + raise HTTPException(400, 'Arquivo não é um SVG puro') + +VALIDADORES_TEXTO ={ + 'txt': validar_plain, + 'csv': validar_csv, + 'svg': validar_svg +} \ No newline at end of file diff --git a/app/utils/validacao.py b/app/utils/validacao.py new file mode 100644 index 0000000..f01b41c --- /dev/null +++ b/app/utils/validacao.py @@ -0,0 +1,40 @@ +import filetype +from fastapi import UploadFile, HTTPException +from app.models.enums import EXTENSOES_PERMITIDAS +from app.utils.sem_magic_bytes import VALIDADORES_TEXTO +from pathlib import Path + +TAMANHO_MAXIMO_MB = 10 +TAMANHO_MAXIMO_BYTES = TAMANHO_MAXIMO_MB * 1024 * 1024 + +async def validar_tipo(arquivo: UploadFile) -> str: + # todos os filetypes estão nos primeiros 261 bytes + # não é necessário ler todo o arquivo para saber o tipo + header = await arquivo.read(261) + await arquivo.seek(0) # cursor volta para início do arquivo + + tipo = filetype.guess(header) + if tipo is not None: + if tipo.extension not in EXTENSOES_PERMITIDAS: + raise HTTPException(400, 'Tipo de arquivo não permitido') + return tipo.extension + + extensao = Path(arquivo.filename).suffix.lstrip('.').lower() + validador = VALIDADORES_TEXTO.get(extensao) + + if validador is None: + raise HTTPException(400, 'Tipo de arquivo não reconhecido') + + await validador(arquivo) + return extensao + +def validar_tamanho(arquivo: UploadFile) -> None: + if not arquivo.size: + raise HTTPException(400, 'Arquivo vazio') + + if arquivo.size > TAMANHO_MAXIMO_BYTES: + raise HTTPException(413, f'Arquivo excede {TAMANHO_MAXIMO_MB}mb') + +async def validar_anexo(arquivo: UploadFile) -> str: + validar_tamanho(arquivo) + return await validar_tipo(arquivo) \ No newline at end of file diff --git a/main.py b/main.py new file mode 100644 index 0000000..d35c383 --- /dev/null +++ b/main.py @@ -0,0 +1,39 @@ +from datetime import datetime +from fastapi import FastAPI, HTTPException, Request +from fastapi.exceptions import RequestValidationError +from fastapi.responses import JSONResponse +from app.routes.anexos import router as anexos_router + +app = FastAPI( + title='ACTA Import API', + description='API de importação de arquivos de .pdf, .docx, .pptx, .txt, .xlsx, .csv e imagens para ciclos PDCA do ACTA' +) + +def _resposta_erro(status_code: int, mensagens: list[str]) -> JSONResponse: + return JSONResponse( + status_code=status_code, + content={ + 'mensagens': mensagens, + 'httpStatus': status_code, + 'timestamp': datetime.now().isoformat(timespec='seconds') + } + ) + +@app.exception_handler(HTTPException) +async def tratar_http_exception(_request: Request, exc: HTTPException): + detalhe = exc.detail if isinstance(exc.detail, list) else [str(exc.detail)] + return _resposta_erro(exc.status_code, [str(mensagem) for mensagem in detalhe]) + +@app.exception_handler(RequestValidationError) +async def tratar_erro_validacao(_request: Request, exc: RequestValidationError): + mensagens = [] + for erro in exc.errors(): + campo = '.'.join(str(parte) for parte in erro['loc'] if parte != 'body') + mensagens.append(f"{campo}: {erro['msg']}" if campo else erro['msg']) + return _resposta_erro(422, mensagens) + +@app.exception_handler(Exception) +async def tratar_erro_interno(_request: Request, exc: Exception): + return _resposta_erro(500, ['Erro interno do servidor']) + +app.include_router(anexos_router) \ No newline at end of file diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..7b9527f --- /dev/null +++ b/requirements.txt @@ -0,0 +1,11 @@ +fastapi==0.141.1 +uvicorn[standard]==0.52.3 +psycopg2-binary==2.9.13 +redis==8.1.0 +rq==2.12.0 +cloudinary==1.46.2 +filetype==1.2.0 +defusedxml==0.7.1 +python-dotenv==1.2.2 +httpx==0.28.1 +python-multipart==0.0.32 \ No newline at end of file diff --git a/tests/test_anexos.py b/tests/test_anexos.py new file mode 100644 index 0000000..a9cc625 --- /dev/null +++ b/tests/test_anexos.py @@ -0,0 +1,95 @@ +from io import BytesIO +import asyncio +from pathlib import Path +from unittest.mock import ANY, AsyncMock, Mock +import pytest +from fastapi import UploadFile +from starlette.datastructures import Headers +from app.models.enums import Categoria, Status +from app.services import anexos +from app.schemas import AnexoForm, UsuarioAutenticado + +class Cursor: + def __init__(self, retorno=(42,)): + self.retorno = retorno + self.executados = [] + + def __enter__(self): + return self + + def __exit__(self, *args): + return False + + def execute(self, sql, parametros): + self.executados.append((sql, parametros)) + + def fetchone(self): + return self.retorno + +class Connection: + def __init__(self): + self.cursor_obj = Cursor() + self.commits = 0 + + def cursor(self): + return self.cursor_obj + + def commit(self): + self.commits += 1 + +def dados(): + return AnexoForm(7, 9, Categoria.EVIDENCIA, "descrição") + +def usuario(): + return UsuarioAutenticado(idUsuario=3, idEmpresa=5) + +def test_insere_metadados_com_usuario_autenticado(): + conn = Connection() + arquivo = UploadFile(filename="arquivo.pdf", file=BytesIO(b"pdf"), headers=Headers({"content-type": "application/pdf"})) + arquivo.size = 3 + + resultado = anexos._inserir_metadados(conn, arquivo, dados(), "pdf", usuario()) + + assert resultado == 42 + assert conn.commits == 1 + assert "INSERT INTO pdca.anexo" in conn.cursor_obj.executados[1][0] + assert conn.cursor_obj.executados[1][1][:4] == (5, 7, 3, 9) + +def test_criar_anexo_enfileira_upload_e_retorna_aceito(monkeypatch): + tmp_path = Path("tests/.tmp_uploads") + tmp_path.mkdir(exist_ok=True) + monkeypatch.setattr(anexos, "DIR_TEMP", tmp_path) + monkeypatch.setattr(anexos, "validar_anexo", AsyncMock(return_value="pdf")) + monkeypatch.setattr(anexos, "_inserir_metadados", Mock(return_value=42)) + enqueue = Mock() + monkeypatch.setattr(anexos.fila_uploads, "enqueue", enqueue) + + arquivo = UploadFile(filename="arquivo.pdf", file=BytesIO(b"pdf")) + arquivo.size = 3 + resposta = asyncio.run(anexos.processar_anexo(arquivo, dados(), Connection(), usuario())) + + assert resposta == {"id": 42, "status": Status.PROCESSANDO.value} + enqueue.assert_called_once() + assert enqueue.call_args.args[1].id_anexo == 42 + assert enqueue.call_args.args[1].extensao == "pdf" + +def test_criar_anexo_marca_erro_se_falhar_ao_enfileirar(monkeypatch): + tmp_path = Path("tests/.tmp_uploads") + tmp_path.mkdir(exist_ok=True) + monkeypatch.setattr(anexos, "DIR_TEMP", tmp_path) + monkeypatch.setattr(anexos, "validar_anexo", AsyncMock(return_value="pdf")) + monkeypatch.setattr(anexos, "_inserir_metadados", Mock(return_value=42)) + monkeypatch.setattr(anexos.fila_uploads, "enqueue", Mock(side_effect=RuntimeError)) + atualizar = Mock() + monkeypatch.setattr(anexos, "_atualizar_status_erro", atualizar) + + with pytest.raises(Exception) as erro: + asyncio.run(anexos.processar_anexo( + UploadFile(filename="arquivo.pdf", file=BytesIO(b"pdf")), + dados(), + Connection(), + usuario() + )) + + assert erro.value.status_code == 503 + atualizar.assert_called_once_with(ANY, 42, 3) diff --git a/tests/test_jobs.py b/tests/test_jobs.py new file mode 100644 index 0000000..a537810 --- /dev/null +++ b/tests/test_jobs.py @@ -0,0 +1,50 @@ +from pathlib import Path +from app.models.enums import Status +from app.schemas import UploadJob +from app.services import jobs + +class Cursor: + def __init__(self): + self.executados = [] + + def __enter__(self): + return self + + def __exit__(self, *args): + return False + + def execute(self, sql, parametros): + self.executados.append((sql, parametros)) + +class Connection: + def __init__(self): + self.cursor_obj = Cursor() + self.commits = 0 + + def cursor(self): + return self.cursor_obj + + def commit(self): + self.commits += 1 + +def test_processa_upload_com_sucesso(monkeypatch): + caminho = Path("tests/.tmp_arquivo.pdf") + caminho.write_bytes(b"pdf") + conexao = Connection() + job = UploadJob(id_anexo=42, id_usuario=3, conteudo=b"pdf", nome_original="arquivo.pdf", extensao="pdf") + monkeypatch.setattr(jobs, "conn_worker", lambda: _contexto(conexao)) + monkeypatch.setattr(jobs, "enviar_cloudinary", lambda *args, **kwargs: {"url": "https://arquivo", "public_id": "42", "resource_type": "image"}) + + jobs.processar_upload(job) + + assert any(Status.ATIVO.value in parametros for _, parametros in conexao.cursor_obj.executados) + +class _contexto: + def __init__(self, valor): + self.valor = valor + + def __enter__(self): + return self.valor + + def __exit__(self, *args): + return False \ No newline at end of file diff --git a/tests/test_validacao.py b/tests/test_validacao.py new file mode 100644 index 0000000..301156d --- /dev/null +++ b/tests/test_validacao.py @@ -0,0 +1,44 @@ +from io import BytesIO +import asyncio +import pytest +from fastapi import HTTPException, UploadFile +from app.utils.validacao import validar_anexo + +def arquivo(nome: str, conteudo: bytes, tamanho: int | None = None) -> UploadFile: + upload = UploadFile(filename=nome, file=BytesIO(conteudo)) + upload.size = len(conteudo) if tamanho is None else tamanho + return upload + +def test_aceita_txt_utf8(): + upload = arquivo("evidencia.txt", "conteúdo válido".encode()) + + assert asyncio.run(validar_anexo(upload)) == "txt" + +def test_rejeita_arquivo_vazio(): + with pytest.raises(HTTPException) as erro: + asyncio.run(validar_anexo(arquivo("vazio.txt", b""))) + + assert erro.value.status_code == 400 + assert erro.value.detail == "Arquivo vazio" + +def test_rejeita_arquivo_maior_que_o_limite(): + upload = arquivo("grande.txt", b"conteudo", tamanho=10 * 1024 * 1024 + 1) + + with pytest.raises(HTTPException) as erro: + asyncio.run(validar_anexo(upload)) + + assert erro.value.status_code == 413 + +def test_rejeita_txt_com_byte_nulo(): + with pytest.raises(HTTPException) as erro: + asyncio.run(validar_anexo(arquivo("invalido.txt", b"texto\x00invalido"))) + + assert erro.value.status_code == 400 + assert erro.value.detail == "Arquivo não é um TXT puro" + +def test_rejeita_extensao_nao_reconhecida(): + with pytest.raises(HTTPException) as erro: + asyncio.run(validar_anexo(arquivo("arquivo.exe", b"conteudo"))) + + assert erro.value.status_code == 400 + assert erro.value.detail == "Tipo de arquivo não reconhecido" \ No newline at end of file diff --git a/worker.py b/worker.py new file mode 100644 index 0000000..0855e1a --- /dev/null +++ b/worker.py @@ -0,0 +1,12 @@ +from redis import Redis +from rq import SimpleWorker +from rq.timeouts import TimerDeathPenalty +from app.database.redis import get_pool +from app.services.fila import fila_uploads + +# timeout adaptado para windows +class WindowsSimpleWorker(SimpleWorker): + death_penalty_class = TimerDeathPenalty + +if __name__ == '__main__': + WindowsSimpleWorker([fila_uploads], connection=Redis(connection_pool=get_pool()),).work() \ No newline at end of file