[ADD] CAMBIOS EN ENDPOINTS Y COMENTARIOS

parent 179e850a
# intranet-sycronizacion
Sincroniza estudiantes (con padres y planes de pago) y apoderados desde PostgreSQL hacia MongoDB mediante sondeo completo cada N minutos, en dos pipelines independientes.
Sincroniza estudiantes (con padres y planes de pago) desde PostgreSQL hacia MongoDB mediante sondeo completo cada N minutos.
## Configuración
1. La tabla `sync_state` se crea/inicializa automáticamente al arrancar (filas `id=1` estudiantes, `id=2` apoderados) — no se requiere paso de migración manual (`migrations/0001_create_sync_state.sql` se conserva como referencia).
1. La tabla `sync_state` se crea/inicializa automáticamente al arrancar (fila `id=1` estudiantes) — no se requiere paso de migración manual (`migrations/0001_create_sync_state.sql` se conserva como referencia).
2. Asegúrate de que la colección `students` exista en Mongo con el validador `$jsonSchema` provisto (requeridos: `student_id` int, `enrollment_id` int). La colección `deleted_students` archiva los estudiantes eliminados en el origen — no requiere configuración de esquema.
3. Asegúrate de que la colección `parents` exista en Mongo con un validador `$jsonSchema` que requiera `parent_id` int. La colección `deleted_parents` archiva los apoderados eliminados en el origen — no requiere configuración de esquema.
4. Define las variables de entorno: `PG_DSN`, `MONGO_URI`, `MONGO_DB=intranet` (o tu base de datos), `SYNC_INTERVAL` (por defecto `5m`, compartido por todos los pipelines), `PORT` (por defecto `8080`). Los nombres de función PG / colección Mongo por pipeline ya no son variables de entorno — se definen en `internal/config/pipelines.go`. Agrega un nuevo pipeline añadiendo una entrada ahí, no variables de entorno.
5. `go run ./cmd/server`
3. Define las variables de entorno: `PG_DSN`, `MONGO_URI`, `MONGO_DB=intranet` (o tu base de datos), `SYNC_INTERVAL` (por defecto `5m`), `PORT` (por defecto `8080`). Los nombres de función PG / colección Mongo por pipeline ya no son variables de entorno — se definen en `internal/config/pipelines.go`. Agrega un nuevo pipeline añadiendo una entrada ahí, no variables de entorno.
4. `go run ./cmd/server`
## Endpoints
- `GET /health` — verifica la conectividad con Postgres y Mongo.
- `POST /sync/trigger` — ejecuta un ciclo de sincronización de estudiantes de inmediato.
- `GET /sync/status` — última ejecución de sincronización de estudiantes, filas sincronizadas (desglose creado/actualizado/eliminado), último error si lo hay.
- `POST /sync/parents/trigger` — ejecuta un ciclo de sincronización de apoderados de inmediato.
- `GET /sync/parents/status` — última ejecución de sincronización de apoderados, filas sincronizadas (desglose creado/actualizado/eliminado), último error si lo hay.
- `POST /sync/students/trigger` — ejecuta un ciclo de sincronización de estudiantes de inmediato.
- `GET /sync/students/status` — última ejecución de sincronización de estudiantes, filas sincronizadas (desglose creado/actualizado/eliminado), último error si lo hay.
## Contratos del origen Postgres
Estudiantes: una consulta SQL cruda (no una función) que une `matricula.ma_estudiante`/`persona.pe_persona`/`matricula.ma_matricula`/`caja.ca_plan_de_pago`, agrupada por estudiante. Cada fila se convierte en un documento Mongo, upsert por `_id = student_id`, con `updated_at` estampado por este servicio. Requiere `student_id` y `enrollment_id` numéricos.
`func_apoderado_listar()` no recibe parámetros y devuelve el conjunto completo como un único envoltorio JSON: `{"status": bool, "message": text, "data": [...]}`. Cada elemento de `data[]` se convierte en un documento Mongo en la colección `parents`, upsert por `_id = parent_id`, con `updated_at` estampado por este servicio. Solo se requiere/valida `parent_id`.
## Comportamiento de sincronización
Cada ciclo (estudiantes y apoderados de forma independiente) corre en concurrencia (pool de workers acotado, 10 en vuelo por defecto) para que una migración completa termine más rápido que un bucle secuencial:
Cada ciclo corre en concurrencia (pool de workers acotado, 10 en vuelo por defecto) para que una migración completa termine más rápido que un bucle secuencial:
- **Registro nuevo**: id aún no en Mongo → insertado, contado como `Created`.
- **Registro existente**: id ya en Mongo → campos sobrescritos, contado como `Updated`.
- **Registro eliminado**: id presente en Mongo pero ya no devuelto por el origen Postgres → el documento se mueve (no solo se elimina) a la colección de archivo del pipeline (`deleted_students` / `deleted_parents`) con una marca `deleted_at`, contado como `Deleted`.
- **Registro eliminado**: id presente en Mongo pero ya no devuelto por el origen Postgres → el documento se mueve (no solo se elimina) a la colección de archivo (`deleted_students`) con una marca `deleted_at`, contado como `Deleted`.
Los dos pipelines son totalmente independientes: tick de scheduler separado, fila `sync_state` separada, dominio de fallos separado — un error del lado de estudiantes no bloquea la sincronización de apoderados ni viceversa.
El estado persistido (`sync_state`) solo avanza si el ciclo completo (upserts + reconciliación de borrados) tiene éxito; cualquier fallo se corta y se reporta vía `/sync/students/status`.
......@@ -95,7 +95,7 @@ func main() {
svc := appsync.NewService(reader, upserter, upserter, mover, state, time.Now)
// Sembrar el último resultado desde el estado persistido, para que
// /sync/status muestre algo útil apenas arranca (antes del primer ciclo).
// /sync/students/status muestre algo útil apenas arranca (antes del primer ciclo).
if lastSynced, err := state.Get(bootCtx); err != nil {
log.Printf("could not load persisted %s sync state at boot, skipping seed: %v", p.Name, err)
} else if !lastSynced.IsZero() {
......
......@@ -57,14 +57,14 @@ func syncResultJSON(err error, result appsync.SyncResult) gin.H {
return resp
}
// Trigger responde POST /sync/trigger: corre una sincronización AHORA y
// Trigger responde POST /sync/students/trigger: corre una sincronización AHORA y
// devuelve cómo salió.
func (h *Handlers) Trigger(c *gin.Context) {
err := h.studentsSvc.Run(c.Request.Context())
c.JSON(http.StatusOK, syncResultJSON(err, h.studentsSvc.LastResult()))
}
// Status responde GET /sync/status: devuelve el resultado del último ciclo
// Status responde GET /sync/students/status: devuelve el resultado del último ciclo
// (sin correr uno nuevo), incluyendo el último error si lo hubo.
func (h *Handlers) Status(c *gin.Context) {
result := h.studentsSvc.LastResult()
......
......@@ -29,7 +29,7 @@ func NewRouter(studentsSvc *appsync.Service, pgPing, mongoPing func(context.Cont
// Registro de rutas: método HTTP + path -> función manejadora.
r.GET("/health", h.Health) // ¿están vivas Postgres y Mongo?
r.POST("/sync/trigger", h.Trigger) // forzar una sincronización ahora
r.GET("/sync/status", h.Status) // ver el resultado del último ciclo
r.POST("/sync/students/trigger", h.Trigger) // forzar una sincronización ahora
r.GET("/sync/students/status", h.Status) // ver el resultado del último ciclo
return r
}
......@@ -136,7 +136,7 @@ func (f *FunctionReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error
// parseGenericEnvelope es como parseFunctionEnvelope pero exige un solo campo
// id configurable (idField), sin obligar enrollment_id. Sirve para pipelines
// genéricos (ej. apoderados con parent_id).
// genéricos que solo necesitan validar un id numérico.
func parseGenericEnvelope(raw []byte, idField string) ([]appsync.SyncRow, error) {
var env functionEnvelope
if err := json.Unmarshal(raw, &env); err != nil {
......
......@@ -50,7 +50,7 @@ type DeletedMover interface {
MoveToDeleted(ctx context.Context, id int) error
}
// SyncResult resume cómo salió el último ciclo. Alimenta /sync/status.
// SyncResult resume cómo salió el último ciclo. Alimenta /sync/students/status.
type SyncResult struct {
RanAt time.Time // cuándo corrió
RowsSynced int // total de filas procesadas
......@@ -219,7 +219,7 @@ func (s *Service) reconcileDeleted(ctx context.Context, pgIDs map[int]struct{})
}
// SeedLastResult inicializa lastResult con un timestamp persistido (ej. justo
// después de arrancar el proceso), para que /sync/status muestre el último
// después de arrancar el proceso), para que /sync/students/status muestre el último
// ciclo exitoso ANTES de que corra un ciclo nuevo. Deja RowsSynced y Err en
// cero a propósito: no tenemos registro de ellos tras un reinicio.
func (s *Service) SeedLastResult(t time.Time) {
......@@ -229,7 +229,7 @@ func (s *Service) SeedLastResult(t time.Time) {
}
// LastResult devuelve una copia del último resultado, de forma segura para
// concurrencia (protegida por el mutex). La lee el endpoint /sync/status.
// concurrencia (protegida por el mutex). La lee el endpoint /sync/students/status.
func (s *Service) LastResult() SyncResult {
s.mu.Lock()
defer s.mu.Unlock()
......
......@@ -32,7 +32,7 @@ type PGExecutor interface {
}
// StateStore guarda y recupera "cuándo fue la última sincronización exitosa".
// Se usa para el endpoint /sync/status y para retomar estado tras reiniciar.
// Se usa para el endpoint /sync/students/status y para retomar estado tras reiniciar.
type StateStore interface {
Get(ctx context.Context) (time.Time, error)
Set(ctx context.Context, t time.Time) error
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment