Refactor: Importación de Evaluaciones (Cosmic Latte) — Referencia técnica - JU-DEV-Bootcamps/ERAS GitHub Wiki
Documento de referencia del refactor del feature de importación de evaluaciones desde Cosmic Latte hacia ERAS. Cubre backend (
ERAS-BE, ASP.NET Core 8 / EF Core 9 / PostgreSQL) y frontend (ERAS-FE, Angular 19). Pensado como guía para mantener y extender el feature.
El proceso de importación pasó de ser síncrono, bloqueante y frágil a una arquitectura asíncrona por estudiante con extracción y confirmación separadas:
| Antes | Después |
|---|---|
GET polls extraía todas las respuestas (~11 s, ~2.3 MB) y bloqueaba la UI |
La extracción corre en background con progreso en vivo |
El "save" reenviaba todo el payload al backend (causaba 413) |
El "confirm" envía solo IDs; el backend persiste desde lo ya extraído |
| Una transacción todo-o-nada; un fallo dejaba datos huérfanos | Una transacción por estudiante; fallos aislados y reintentables |
Full table scans, triple bucle O(n³), SaveChanges por item |
Queries filtradas, lookups O(n), updates set-based |
| 533 líneas de orquestador con estado mutable | Orquestador dividido en colaboradores + contexto explícito |
| Sin visibilidad de progreso |
ImportJob + ImportJobItem con estado por estudiante + polling |
El trabajo se hizo en fases incrementales (0→5). Cada fase es independientemente desplegable.
flowchart LR
subgraph FE["ERAS-FE (Angular 19)"]
EP["Evaluation Process List<br/>(trigger import)"]
IV["Import Status View<br/>(unified, polling)"]
end
subgraph API["ERAS-BE · Eras.Api"]
CC["CosmicLatteController<br/>extract / confirm / status / items / retry"]
end
subgraph APP["Eras.Application"]
IJS["ImportJobService"]
ORC["PollOrchestratorService<br/>+ importers (Phase 3)"]
UOW["IUnitOfWork"]
end
subgraph INFRA["Eras.Infrastructure"]
Q["ImportJobQueue<br/>Channel of int · singleton"]
W["ImportQueueBackgroundService<br/>(BackgroundService)"]
CL["CosmicLatteAPIService<br/>ExtractRespondentsAsync"]
REPO["ImportJob / ImportJobItem<br/>repositories (set-based)"]
end
DB[("PostgreSQL")]
CLAPI[["Cosmic Latte API"]]
EP -->|"POST imports/extract"| CC
IV -->|"poll GET imports/:id (+items)"| CC
IV -->|"POST confirm / retry"| CC
CC --> IJS
IJS -->|"crea ImportJob + encola id"| Q
IJS --> REPO
W -->|"dequeue id"| Q
W -->|"fase Extracting"| CL
W -->|"fase Importing"| ORC
CL -->|"detalle por estudiante"| CLAPI
ORC --> UOW
ORC --> REPO
REPO --> DB
UOW --> DB
Idea central: el ImportJob es una máquina de estados persistida; un único worker la avanza por
fases consumiendo de una cola in-process. El frontend nunca espera; hace polling y muestra progreso.
stateDiagram-v2
[*] --> Extracting: POST imports/extract
Extracting --> Ready: extracción completa
Extracting --> Failed: error de extracción
Ready --> Importing: POST confirm (ids)
Importing --> Completed: todos OK
Importing --> PartiallyCompleted: algunos fallaron
Importing --> Failed: setup falló
PartiallyCompleted --> Importing: POST retry (ids)
Failed --> Importing: POST retry (ids)
Completed --> [*]
Estados por item (ImportJobItem.Status):
Extracting → Extracted → (Skipped | Queued) → Running → Completed | Failed.
-
Extracted: respondiente traído de Cosmic Latte, a la espera de confirmación. -
Skipped: no seleccionado al confirmar. -
Queued/Running/Completed/Failed: ciclo de la importación real (persistencia).
El estado agregado del job se recomputa por SQL (GROUP BY status) contando solo la fase de
import (Pending = Queued+Running, Completed, Failed); Extracted/Skipped no cuentan.
sequenceDiagram
autonumber
participant U as Usuario
participant FE as Import Status View
participant API as CosmicLatteController
participant SVC as ImportJobService
participant Q as ImportJobQueue
participant W as Worker
participant CL as CosmicLatteAPIService
participant ORC as PollOrchestrator
participant DB as PostgreSQL
U->>FE: Importar (modal: poll, fechas, config)
FE->>API: POST imports/extract
API->>SVC: StartExtractionAsync(...)
SVC->>DB: ImportJob(status=Extracting)
SVC->>Q: enqueue(jobId)
API-->>FE: 202 (importJobId)
FE->>FE: navega a import-status/:id
W->>Q: dequeue(jobId)
W->>CL: ExtractRespondentsAsync(...)
CL->>CL: lista + detalle por estudiante (paralelo acotado)
loop por respondiente
CL-->>W: onExtracted(pollDto, alreadyImported)
W->>DB: ImportJobItem(status=Extracted) + extractedCount++
end
W->>DB: job → Ready (totalCount = extraídos)
loop polling cada 3s
FE->>API: GET imports/:id (+items)
API-->>FE: progreso (barra) / lista
end
U->>FE: selecciona y "Import selected"
FE->>API: POST imports/:id/confirm (itemIds)
API->>SVC: ConfirmImportAsync
SVC->>DB: seleccionados=Queued, resto=Skipped, job=Importing
SVC->>Q: enqueue(jobId)
API-->>FE: 202
W->>Q: dequeue(jobId)
W->>ORC: SetupImportStructureAsync (1 transacción)
loop por item Queued
W->>ORC: ProcessStudentAsync(poll) (1 transacción/estudiante)
ORC->>DB: student + poll instance + answers
W->>DB: item → Completed | Failed
end
W->>DB: job → Completed | PartiallyCompleted
erDiagram
import_jobs ||--o{ import_job_items : "1:N"
import_jobs {
int Id PK
int evaluation_id
string status "varchar(20)"
int total_count
int processed_count
int extracted_count
int retry_count
string error_message
string evaluation_set_name
int configuration_id
string start_date
string end_date
jsonb polls_payload "legacy/empty en flujo extract"
timestamptz created_at_utc
timestamptz updated_at_utc
}
import_job_items {
int Id PK
int import_job_id FK
string student_email
string student_name
string cohort
string status "varchar(20)"
int retry_count
bool is_already_imported
string error_message
jsonb poll_payload "PollDTO del estudiante"
timestamptz created_at_utc
timestamptz updated_at_utc
}
- Índices:
ix_import_jobs_evaluation_id,ix_import_job_items_import_job_id_status. -
poll_instances.answers_hash(+ índiceix_poll_instances_student_id_answers_hash): detección de duplicados por hash persistido e indexado (antes se cargaba todo y se hasheaba en memoria). -
Migraciones escritas a mano (
AddAnswersHashToPollInstance,AddImportJobs,AddImportJobItems,AddExtractionFields) para no arrastrar drift preexistente depoll_instances(un FK/índice deEvaluationIdque el modelo configura pero que ninguna migración previa capturó). Se aplican solas al arrancar víaDatabase.Migrate()enProgram.cs.
| Componente | Archivo | Rol |
|---|---|---|
ImportJob / ImportJobItem
|
Eras.Domain/Entities/ |
Entidades de dominio + enum ImportJobStatus
|
ImportJobService |
Eras.Application/Services/ImportJobService.cs |
StartExtractionAsync, ConfirmImportAsync, GetStatusAsync, GetItemsAsync, RetryItemsAsync
|
PollOrchestratorService |
Eras.Application/Services/PollOrchestratorService.cs |
Coordinador: SetupImportStructureAsync (1 tx) + ProcessStudentAsync (1 tx/estudiante) + RebuildPollStructureAsync (estructura desde BD) |
PollStructureImporter / StudentImporter / PollInstanceImporter
|
Eras.Application/Services/ |
Colaboradores con responsabilidad única (Fase 3) |
ImportContext |
Eras.Application/Services/ImportContext.cs |
Estado de versionado explícito (reemplaza campos mutables) |
IUnitOfWork / UnitOfWork
|
Eras.Application/Contracts/... · Eras.Infrastructure/.../UnitOfWork.cs
|
Transacción ambiente (los repos batch la reutilizan en vez de abrir nidadas) |
ImportJobQueue |
Eras.Infrastructure/BackgroundProcessing/ImportJobQueue.cs |
Channel<int> singleton (cola in-process) |
ImportQueueBackgroundService |
Eras.Infrastructure/BackgroundProcessing/ |
Worker: ExtractAsync / ImportAsync / agregado; resiliente (un fallo no tumba el host) |
CosmicLatteAPIService.ExtractRespondentsAsync |
Eras.Infrastructure/External/CosmicLatteClient/ |
Extracción incremental (HTTP paralelo acotado, callback serializado) |
ImportJobRepository / ImportJobItemRepository
|
Eras.Infrastructure/.../Repositories/ |
Updates set-based (ExecuteUpdateAsync) — evitan conflictos de tracking |
| Método | Ruta | Descripción |
|---|---|---|
POST |
imports/extract |
Inicia extracción en background → 202 { importJobId, status:"Extracting" }
|
GET |
imports/{id} |
Estado del job (status, total/processed/extracted counts) |
GET |
imports/{id}/items |
Estado por estudiante (incluye isAlreadyImported) |
POST |
imports/{id}/confirm |
Body { itemIds } → importa seleccionados (resto Skipped) → 202
|
POST |
imports/{id}/retry |
Body { itemIds } → reencola fallidos → 202
|
POST |
polls/{evaluationId} |
Legacy (payload completo). Conservado, no usado por el FE nuevo |
GET |
polls, polls/names, health
|
Preview legacy / utilidades |
-
Una transacción por estudiante (
ProcessStudentAsyncenvuelto enIUnitOfWork.ExecuteInTransactionAsync): un fallo aísla a ese estudiante; los demás se confirman → habilita el reintento selectivo. -
El reintento no depende de estado en memoria:
ProcessStudentAsyncreconstruye la estructura del poll desde la BD (IVariableRepository.GetAllByPollUuidAsync), por eso un retry en otra request funciona. -
Updates set-based (
ExecuteUpdateAsync) en estado de job/item: evitan el conflicto "entidad ya trackeada" que aparece al hacerGetByIdAsync+UpdateAsyncrepetidos en un mismo scope. -
Worker resiliente: cada job se procesa en
try/catch; una excepción nunca detiene el host (BackgroundServicepor defecto haríaStopHost). -
Concurrencia de extracción: HTTP en paralelo acotado (
SemaphoreSlim(8)), pero la persistencia del callback es serializada (SemaphoreSlim(1)) porque elDbContext(scoped) no es thread-safe.
| Componente | Archivo | Rol |
|---|---|---|
CosmicLatteService |
core/services/api/cosmic-latte.service.ts |
startExtraction, confirmImport, getImportStatus, getImportItems, retryImportItems
|
import-job.model.ts |
core/models/ |
ImportJobStatus, ImportJobStatusModel, ImportJobItem, ImportItemRow
|
EvaluationProcessListComponent |
modules/lists/components/evaluacion-process/ |
Trigger: modal → startExtraction → navega a la vista unificada |
ImportStatusComponent |
modules/imports/components/import-status/ |
Vista unificada por fases con polling, barra de progreso, selección y confirmación/retry |
ImportStatusBadgeComponent |
.../import-status/import-status-badge/ |
Badge de estado (mapa estado→clase) |
TableWithActionsComponent |
shared/components/table-with-actions/ |
Reutilizado: grid con multi-select + acciones por fila |
flowchart TD
A["Evaluation Process (modal)"] -->|"startExtraction"| B["POST imports/extract → 202"]
B --> C["navega import-status/:id"]
C --> D{"polling status"}
D -->|"Extracting"| E["barra indeterminada + contador"]
D -->|"Ready"| F["selección de Extracted + Import selected"]
F -->|"confirmImport (itemIds)"| G["POST confirm → 202"]
G --> D
D -->|"Importing"| H["barra determinada processed/total"]
D -->|"Completed / Partial / Failed"| I["resultado + Retry selected"]
I -->|"retryImportItems"| G
-
Polling con
timer(0,3000)+switchMap+takeWhile: solo activo en fasesExtracting/Importing; pausa enReady(espera confirmación) y en terminales. - No bloquea la UI: la ruta es navegable; el job persiste y se puede volver.
-
Selección: en
Readyse seleccionanExtracted(pre-deselecciona ya-importados e inválidos); en fase import, multi-select deFailedpara reintentar.
| Fase | Foco | Cambios principales |
|---|---|---|
| 0 | Correctitud / transaccionalidad |
IUnitOfWork (transacción única), repos batch respetan tx ambiente, propagar excepciones (rollback real), eliminar SaveManyAnswersAsync, hash de duplicados persistido + índice
|
| 1 | Performance | Eliminar full table scans (queries por PollId), aplanar triple bucle O(n³)→O(n) con diccionario, quitar CreatePollAsync redundante, batch de answers |
| 2 | Async / background |
ImportJob + cola Channel<int> + ImportQueueBackgroundService, endpoint 202 + GET status, reintentos |
| 3 | Complejidad |
ImportContext (estado de versionado explícito), dividir orquestador en PollStructureImporter/StudentImporter/PollInstanceImporter
|
| 4 | Granularidad por estudiante |
ImportJobItem, transacción por estudiante, GET items + POST retry, grid de seguimiento (FE) |
| 5 | Extracción como job + save sin payload | Fases Extracting/Ready/Importing, ExtractRespondentsAsync incremental, POST extract/confirm, vista unificada con barra de progreso, retiro de componentes legacy |
-
Nuevo estado de item/job: agregar al enum
ImportJobStatus(la columna esvarchar(20)), mapear enImportStatusBadgeComponent(elRecord<ImportJobStatus,...>debe ser exhaustivo) y en el union del FE. -
Otro paso en la importación de un estudiante: extender
PollOrchestratorService.ProcessStudentAsync(sigue dentro de la transacción por estudiante); reusar los importers. -
Cambiar el paralelismo de extracción: ajustar
SemaphoreSlimenExtractRespondentsAsync(httpGate); mantenerpersistGate(1)para no usar elDbContextconcurrentemente. -
Persistencia/idempotencia entre reinicios: los payloads viven en
ImportJobItem.poll_payload; el worker reprocesa por estado, así que un job puede re-encolarse sin perder datos. -
Migraciones: escribirlas a mano (ver §5) para no empaquetar el drift de
poll_instances; actualizarAppDbContextModelSnapshot.csen consecuencia.
-
Backend unit: worker en fase
Extracting(crea items,extractedCountsube, →Ready);confirmmarcaQueued/Skipped; per-student con un fallo →PartiallyCompleted; retry vuelve aCompleted. -
Frontend (Jasmine/Karma): servicio (URLs), vista unificada (barra por fase, selección por estado,
polling se detiene en
Ready/terminal), badge. -
E2E local (Docker): ver
docker-compose.local.yml(stack para Docker nativo en WSL). Flujo validado: extracción incremental0→N → Ready, confirmar selección →Importing → Completed, retry de fallidos. Notas de despliegue local: proxy nginx/api(same-origin, evita CORS), backend en host networking (issuer de Keycloak coherente),client_max_body_sizeelevado.
- Endpoints/método legacy (
POST polls/{id},GetAllPollsPreview,CosmicLatteAPIService.SavePreviewPolls,importAnswerBySurveyen el FE) quedaron sin uso; candidatos a eliminación. -
Drift de
poll_instances(FK/índice deEvaluationIdno migrado): abordar en una migración dedicada. - Reintentos automáticos con back-off ante fallos transitorios (hoy el retry es manual).
- Tests automatizados de las Fases 4–5 (se validaron E2E manualmente).