[ADD] ENDPOINTS DE USERS, STUDENTS, PAYMENTPLANS,

parent 2d716325
......@@ -7,24 +7,56 @@ Sincroniza estudiantes (con padres y planes de pago) desde PostgreSQL hacia Mong
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. 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`
4. Horario por pipeline (opcional, `SYNC_<NOMBRE>_...` en mayúsculas, ej. `SYNC_STUDENTS_MODE`, `SYNC_PAYMENT_PLANS_MODE`, `SYNC_PARENTS_MODE`):
- `SYNC_<NOMBRE>_MODE`: `interval` (por defecto) o `daily`.
- `SYNC_<NOMBRE>_INTERVAL`: duración tipo `5m`/`1h` (modo `interval`; si falta, usa `SYNC_INTERVAL`).
- `SYNC_<NOMBRE>_TIME`: hora `HH:MM` (modo `daily`, obligatoria en ese modo). Ej: `SYNC_STUDENTS_MODE=daily` + `SYNC_STUDENTS_TIME=03:00` corre estudiantes una vez al día a las 3am.
5. `go run ./cmd/server`
## Endpoints
- `GET /health` — verifica la conectividad con Postgres y Mongo.
- `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.
Todos los endpoints llevan el prefijo `/api/v1`.
- `GET /api/v1/health` — verifica la conectividad con Postgres y Mongo.
- `POST /api/v1/sync/students/trigger` / `GET /api/v1/sync/students/status` — pipeline de estudiantes: ejecutar ahora / última ejecución, filas sincronizadas (desglose creado/actualizado/sin cambios/eliminado), último error si lo hay.
- `POST /api/v1/sync/payment_plans/trigger` / `GET /api/v1/sync/payment_plans/status` — pipeline de planes de pago (misma forma).
- `POST /api/v1/sync/parents/trigger` / `GET /api/v1/sync/parents/status` — pipeline de padres/apoderados (misma forma).
- `POST /api/v1/sync/users/trigger` / `GET /api/v1/sync/users/status` — pipeline de usuarios (misma forma).
## Estructura del proyecto (por colección/módulo)
Cada colección (`students`, `payment_plans`, `parents`, `users`) es un paquete Go propio bajo `internal/<módulo>` (`internal/students`, `internal/paymentplans`, `internal/parents`, `internal/users`):
- `entity.go` — struct del documento Mongo (implementa `SetUpdatedAt` y `SetRowHash`).
- `query.go``const Query`, el SQL crudo de origen.
- `reader.go``QueryReader`/`NewQueryReader`.
Agregar un módulo nuevo = paquete `internal/<módulo>` + entrada en `internal/config.Pipelines` + case en `cmd/server/main.go` + endpoints en `internal/api`, sin nuevas variables de entorno. Las capas transversales (`sync.Service`, `sync.Scheduler`, `sync.StateStore`, `db.CollectionUpserter`, `db.DeletedMover`) son genéricas y se comparten entre módulos.
## Regla de `_id` en Mongo
**`_id` nunca es el id de negocio.** Todas las colecciones dejan que Mongo genere su propio ObjectID; el id de negocio (`student_id`, `payment_plan_id`, `parent_id`, `user_id`) vive como campo normal del documento, el mismo declarado en `PipelineDef.IDField`. Los upserts y la reconciliación de borrados filtran por ese campo, no por `_id`. Al archivar en `deleted_*` se descarta el `_id` original para que se genere uno nuevo. Requiere índice único sobre `idField` en cada colección Mongo (no lo impone el código).
## 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.
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 `student_id`, con `updated_at` estampado por este servicio. Requiere `student_id` y `enrollment_id` numéricos.
Planes de pago: una consulta SQL propia, un documento Mongo por plan de pago, upsert por `payment_plan_id`.
Padres/apoderados: une `persona.pe_persona` con una subconsulta agregada sobre `matricula.ma_estudiante_apoderado` (una fila por padre/apoderado, con los estudiantes asociados embebidos vía `JSON_AGG`). Upsert por `parent_id`.
Usuarios: una fila por login (`user_id`, `user_login`, `user_password`, `user_creation_date`, `parent_id`, `user_status`). Upsert por `user_id`.
## Comportamiento de sincronización
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:
Cada ciclo tiene dos fases:
1. **Hash + diff (CPU, concurrente)**: un pool de workers acotado (10 en vuelo por defecto) calcula un hash SHA-256 (`row_hash`) de cada fila de Postgres y lo compara contra el `row_hash` ya guardado en Mongo para ese id (traído de antemano en una sola consulta). Las filas cuyo hash coincide se saltan por completo — sin escritura — y se cuentan como `Unchanged`. Solo las filas nuevas o cambiadas reciben `updated_at`/`row_hash` y pasan a la fase 2.
2. **Escritura en lote (I/O, un solo round trip)**: todas las filas cambiadas se mandan en un único `BulkWrite` a Mongo (upsert por `idField`).
- **Registro nuevo**: id aún no en Mongo → insertado, contado como `Created`.
- **Registro existente**: id ya en Mongo → campos sobrescritos, contado como `Updated`.
- **Registro existente, cambiado**: hash distinto al guardado → campos sobrescritos, contado como `Updated`.
- **Registro existente, sin cambios**: hash igual al guardado → no se escribe nada, contado como `Unchanged`.
- **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`.
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`.
El estado persistido (`sync_state`) solo avanza si el ciclo completo (hash/upsert + reconciliación de borrados) tiene éxito; cualquier fallo se corta y se reporta vía `/api/v1/sync/<módulo>/status`.
// Command server es el punto de entrada del servicio y la "raíz de composición"
// (composition root): el ÚNICO lugar donde se crean las implementaciones
// concretas (Postgres, Mongo) y se inyectan en el dominio. Aquí se arma todo
// el árbol de dependencias y se arranca.
//
// Flujo: config -> conexiones a BD -> por cada pipeline armar
// reader/upserter/mover/state/Service/Scheduler -> arrancar schedulers en
// goroutines -> arrancar servidor HTTP -> apagado ordenado con Ctrl+C/SIGTERM.
// Command server es el punto de entrada y raíz de composición del servicio.
package main
import (
......@@ -18,29 +11,28 @@ import (
"syscall"
"time"
"github.com/joho/godotenv" // carga variables desde un archivo .env
"github.com/joho/godotenv"
"intranet-sycronizacion/internal/api"
"intranet-sycronizacion/internal/config"
"intranet-sycronizacion/internal/db"
"intranet-sycronizacion/internal/parents"
"intranet-sycronizacion/internal/paymentplans"
"intranet-sycronizacion/internal/students"
appsync "intranet-sycronizacion/internal/sync"
"intranet-sycronizacion/internal/users"
)
func main() {
// context.Background() es el context "raíz" para las tareas de arranque.
bootCtx := context.Background()
// Carga el archivo .env si existe. El "_ =" descarta el error a propósito:
// en producción las variables pueden venir del entorno, sin .env.
_ = godotenv.Load()
// 1. Cargar y validar configuración. Si falla, log.Fatalf corta el programa.
cfg, err := config.Load()
if err != nil {
log.Fatalf("config error: %v", err)
}
// 2. Conectar a Postgres y preparar la tabla sync_state.
pgPool, err := db.NewPostgresPool(bootCtx, cfg.PGDSN)
if err != nil {
log.Fatalf("postgres connect error: %v", err)
......@@ -49,88 +41,78 @@ func main() {
log.Fatalf("schema bootstrap error: %v", err)
}
// 3. Conectar a Mongo.
mongoClient, err := db.NewMongoClient(bootCtx, cfg.MongoURI)
if err != nil {
log.Fatalf("mongo connect error: %v", err)
}
// svcByName: guarda cada Service por nombre de pipeline, para luego pasar
// el de "students" al router.
svcByName := make(map[string]*appsync.Service, len(config.Pipelines))
// ctx se cancela al recibir Ctrl+C (os.Interrupt) o SIGTERM. Todos los
// schedulers escuchan este ctx para apagarse ordenadamente.
ctx, stop := signal.NotifyContext(context.Background(), os.Interrupt, syscall.SIGTERM)
defer stop()
// 4. Por cada pipeline declarado en config.Pipelines, armar toda su cadena.
// Por cada pipeline declarado en config.Pipelines, armar su cadena reader/upserter/mover/state/Service/Scheduler.
for _, p := range config.Pipelines {
// Colecciones Mongo: la viva y la de archivo.
coll := mongoClient.Database(cfg.MongoDB).Collection(p.MongoCollection)
deletedColl := mongoClient.Database(cfg.MongoDB).Collection(p.MongoDeletedCollection)
// Elegir el lector según cómo esté declarado el pipeline:
// - PGQuery seteado -> SQL crudo (estudiantes)
// - RequireEnrollmentID -> función PG que exige enrollment_id
// - por defecto -> función PG genérica (solo un id)
var reader appsync.PGReader
switch {
case p.PGQuery != "":
reader = db.NewStudentsQueryReader(pgPool, p.PGQuery)
case p.Name == "students":
reader = students.NewQueryReader(pgPool, p.PGQuery)
case p.Name == "payment_plans":
reader = paymentplans.NewQueryReader(pgPool, p.PGQuery)
case p.Name == "parents":
reader = parents.NewQueryReader(pgPool, p.PGQuery)
case p.Name == "users":
reader = users.NewQueryReader(pgPool, p.PGQuery)
case p.RequireEnrollmentID:
reader = db.NewFunctionReader(pgPool, p.PGFunction)
default:
reader = db.NewGenericFunctionReader(pgPool, p.PGFunction, p.IDField)
}
// upserter cumple a la vez MongoUpserter y MongoIDLister (por eso se
// pasa dos veces a NewService).
upserter := db.NewCollectionUpserter(coll)
mover := db.NewDeletedMover(coll, deletedColl)
// upserter cumple a la vez MongoUpserter y MongoIDLister.
upserter := db.NewCollectionUpserter(coll, p.IDField)
mover := db.NewDeletedMover(coll, deletedColl, p.IDField)
state := appsync.NewPGStateStore(pgPool, p.StateID)
// Armar el servicio del dominio inyectándole todas sus dependencias.
// time.Now es la función "dame la hora"; en tests se puede reemplazar.
svc := appsync.NewService(reader, upserter, upserter, mover, state, time.Now)
// Sembrar el último resultado desde el estado persistido, para que
// /sync/students/status muestre algo útil apenas arranca (antes del primer ciclo).
// Sembrar el último resultado desde el estado persistido 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() {
svc.SeedLastResult(lastSynced)
}
// Arrancar el scheduler de este pipeline en su propia goroutine.
scheduler := appsync.NewScheduler(svc, cfg.SyncInterval)
schedule, err := config.ScheduleFor(p.Name, cfg.SyncInterval)
if err != nil {
log.Fatalf("schedule config error for pipeline %s: %v", p.Name, err)
}
scheduler := appsync.NewScheduler(svc, schedule)
go scheduler.Start(ctx)
svcByName[p.Name] = svc
}
// 5. Construir el router HTTP. Los pings se pasan como funciones anónimas
// que envuelven los clientes concretos (así api no depende de pgx/mongo).
router := api.NewRouter(svcByName["students"],
router := api.NewRouter(svcByName["students"], svcByName["payment_plans"], svcByName["parents"], svcByName["users"],
func(pingCtx context.Context) error { return pgPool.Ping(pingCtx) },
func(pingCtx context.Context) error { return mongoClient.Ping(pingCtx, nil) },
)
// 6. Arrancar el servidor HTTP en una goroutine para no bloquear el main.
srv := &http.Server{Addr: ":" + cfg.Port, Handler: router}
go func() {
// ListenAndServe bloquea hasta que el server se cierra. ErrServerClosed
// es el cierre normal (no un error real), por eso se ignora.
if err := srv.ListenAndServe(); err != nil && !errors.Is(err, http.ErrServerClosed) {
log.Fatalf("server error: %v", err)
}
}()
// 7. Esperar la señal de apagado. <-ctx.Done() bloquea hasta Ctrl+C/SIGTERM.
<-ctx.Done()
log.Println("shutdown signal received, shutting down")
// 8. Apagado ordenado: dar hasta 10s para terminar peticiones en curso.
// Apagado ordenado: dar hasta 10s para terminar peticiones en curso.
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if err := srv.Shutdown(shutdownCtx); err != nil {
......
......@@ -9,26 +9,22 @@ import (
appsync "intranet-sycronizacion/internal/sync"
)
// Handlers guarda las dependencias que necesitan los manejadores HTTP.
// Cada método (Health, Trigger, Status) es un "handler" de una ruta.
// Handlers guarda las dependencias de los manejadores HTTP.
type Handlers struct {
studentsSvc *appsync.Service // el servicio de sincronización
pgPing func(context.Context) error // chequeo de salud de Postgres
mongoPing func(context.Context) error // chequeo de salud de Mongo
studentsSvc *appsync.Service
paymentPlansSvc *appsync.Service
parentsSvc *appsync.Service
usersSvc *appsync.Service
pgPing func(context.Context) error
mongoPing func(context.Context) error
}
// Health responde GET /health. Devuelve 200 si ambas bases responden, o 503
// (Service Unavailable) si alguna falla, con el detalle del error.
//
// En Gin, *gin.Context (c) representa la petición y la respuesta. Se usa para
// leer datos de entrada y para escribir la respuesta (c.JSON).
// Health responde GET /health con el estado de Postgres y Mongo.
func (h *Handlers) Health(c *gin.Context) {
// c.Request.Context() propaga cancelación/timeout de la petición HTTP.
pgErr := h.pgPing(c.Request.Context())
mongoErr := h.mongoPing(c.Request.Context())
status := http.StatusOK
// gin.H es simplemente un mapa para armar el cuerpo JSON de respuesta.
body := gin.H{"postgres": "ok", "mongo": "ok"}
if pgErr != nil {
status = http.StatusServiceUnavailable
......@@ -38,17 +34,17 @@ func (h *Handlers) Health(c *gin.Context) {
status = http.StatusServiceUnavailable
body["mongo"] = mongoErr.Error()
}
c.JSON(status, body) // escribe status + cuerpo JSON
c.JSON(status, body)
}
// syncResultJSON arma el cuerpo JSON común para /trigger y /status a partir de
// un SyncResult del dominio. Evita repetir el mismo mapeo en dos lados.
// syncResultJSON arma el cuerpo JSON común de /trigger y /status.
func syncResultJSON(err error, result appsync.SyncResult) gin.H {
resp := gin.H{
"ran_at": result.RanAt,
"rows_synced": result.RowsSynced,
"created": result.Created,
"updated": result.Updated,
"unchanged": result.Unchanged,
"deleted": result.Deleted,
}
if err != nil {
......@@ -57,20 +53,22 @@ func syncResultJSON(err error, result appsync.SyncResult) gin.H {
return resp
}
// 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()))
// triggerHandler arma el handler de POST /sync/<pipeline>/trigger: corre una sincronización ahora.
func triggerHandler(svc *appsync.Service) gin.HandlerFunc {
return func(c *gin.Context) {
err := svc.Run(c.Request.Context())
c.JSON(http.StatusOK, syncResultJSON(err, svc.LastResult()))
}
}
// 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()
// statusHandler arma el handler de GET /sync/<pipeline>/status: devuelve el resultado del último ciclo.
func statusHandler(svc *appsync.Service) gin.HandlerFunc {
return func(c *gin.Context) {
result := svc.LastResult()
resp := syncResultJSON(nil, result)
if result.Err != nil {
resp["last_error"] = result.Err.Error()
}
c.JSON(http.StatusOK, resp)
}
}
// Package api es la capa web (HTTP), construida con el framework Gin.
// Su única responsabilidad es traducir peticiones HTTP <-> llamadas al dominio
// (internal/sync). No contiene lógica de negocio.
// Package api es la capa web (HTTP), construida con Gin.
package api
import (
......@@ -12,24 +10,24 @@ import (
)
// NewRouter construye el router de Gin y registra las rutas.
//
// Recibe sus dependencias inyectadas:
// - studentsSvc: el Service del dominio que hace la sincronización.
// - pgPing / mongoPing: funciones para chequear salud de cada base.
// Se pasan como funciones (no como clientes) para que api no dependa
// directamente de pgx ni de mongo.
//
// Devuelve *gin.Engine, que es el manejador HTTP que main.go pone a escuchar.
func NewRouter(studentsSvc *appsync.Service, pgPing, mongoPing func(context.Context) error) *gin.Engine {
// Handlers agrupa las dependencias que usan los manejadores de rutas.
h := &Handlers{studentsSvc: studentsSvc, pgPing: pgPing, mongoPing: mongoPing}
// gin.Default() crea un router con logger y recuperación de panics incluidos.
func NewRouter(studentsSvc, paymentPlansSvc, parentsSvc, usersSvc *appsync.Service, pgPing, mongoPing func(context.Context) error) *gin.Engine {
h := &Handlers{studentsSvc: studentsSvc, paymentPlansSvc: paymentPlansSvc, parentsSvc: parentsSvc, usersSvc: usersSvc, pgPing: pgPing, mongoPing: mongoPing}
r := gin.Default()
// Registro de rutas: método HTTP + path -> función manejadora.
r.GET("/health", h.Health) // ¿están vivas Postgres y Mongo?
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
v1 := r.Group("/api/v1")
v1.GET("/health", h.Health)
v1.POST("/sync/students/trigger", triggerHandler(studentsSvc))
v1.GET("/sync/students/status", statusHandler(studentsSvc))
v1.POST("/sync/payment_plans/trigger", triggerHandler(paymentPlansSvc))
v1.GET("/sync/payment_plans/status", statusHandler(paymentPlansSvc))
v1.POST("/sync/parents/trigger", triggerHandler(parentsSvc))
v1.GET("/sync/parents/status", statusHandler(parentsSvc))
v1.POST("/sync/users/trigger", triggerHandler(usersSvc))
v1.GET("/sync/users/status", statusHandler(usersSvc))
return r
}
// Package config carga y valida la configuración del servicio.
//
// En Go, un "package" agrupa archivos relacionados (aquí config.go y
// pipelines.go). Todo lo que esté en mayúscula (ej. Config, Load) es
// "exportado": otros paquetes pueden usarlo. Lo que está en minúscula
// (ej. studentsQuery) es privado a este paquete.
package config
import (
"fmt" // formateo de strings y errores
"os" // acceso a variables de entorno (os.Getenv)
"time" // duraciones (time.Duration, time.ParseDuration)
"fmt"
"os"
"strconv"
"strings"
"time"
appsync "intranet-sycronizacion/internal/sync"
)
// Config guarda todos los valores de configuración ya leídos y validados.
//
// Un "struct" es como un objeto/registro: un grupo de campos con nombre.
// Estos campos vienen de variables de entorno (definidas en el archivo .env
// o en el sistema operativo).
// Config guarda los valores de configuración ya leídos y validados.
type Config struct {
PGDSN string // cadena de conexión a PostgreSQL (env PG_DSN)
MongoURI string // cadena de conexión a MongoDB (env MONGO_URI)
MongoDB string // nombre de la base de datos Mongo (env MONGO_DB)
SyncInterval time.Duration // cada cuánto corre la sincronización (env SYNC_INTERVAL)
Port string // puerto del servidor HTTP (env PORT)
PGDSN string
MongoURI string
MongoDB string
SyncInterval time.Duration
Port string
}
// Load lee las variables de entorno, valida las obligatorias y devuelve la
// configuración lista para usar.
//
// En Go una función puede devolver varios valores. Aquí devuelve dos:
// - Config: la configuración construida.
// - error: nil (nulo) si todo salió bien, o un error si algo falló.
//
// El código que llama debe revisar SIEMPRE ese error antes de usar la Config.
// Load lee las variables de entorno, valida las obligatorias y devuelve la configuración.
func Load() (Config, error) {
// Construimos la Config leyendo cada variable de entorno.
// os.Getenv devuelve "" (cadena vacía) si la variable no existe.
cfg := Config{
PGDSN: os.Getenv("PG_DSN"),
MongoURI: os.Getenv("MONGO_URI"),
......@@ -43,9 +29,6 @@ func Load() (Config, error) {
Port: os.Getenv("PORT"),
}
// Estas tres variables son obligatorias. Si falta alguna, no tiene
// sentido arrancar, así que devolvemos un error y el programa se detiene
// ("fail fast": fallar rápido y claro en vez de romper más adelante).
required := map[string]string{
"PG_DSN": cfg.PGDSN,
"MONGO_URI": cfg.MongoURI,
......@@ -53,30 +36,77 @@ func Load() (Config, error) {
}
for name, val := range required {
if val == "" {
// fmt.Errorf crea un error con un mensaje formateado.
return Config{}, fmt.Errorf("missing required env var %s", name)
}
}
// PORT es opcional: si no vino, usamos 8080 por defecto.
if cfg.Port == "" {
cfg.Port = "8080"
}
// SYNC_INTERVAL es opcional: por defecto "5m" (5 minutos).
interval := os.Getenv("SYNC_INTERVAL")
if interval == "" {
interval = "5m"
}
// time.ParseDuration convierte texto como "5m", "30s" o "1h" en una
// duración usable. Si el texto es inválido, devolvemos error.
d, err := time.ParseDuration(interval)
if err != nil {
// %w "envuelve" el error original para no perder su detalle.
return Config{}, fmt.Errorf("invalid SYNC_INTERVAL %q: %w", interval, err)
}
cfg.SyncInterval = d
// Todo válido: devolvemos la config y nil como error.
return cfg, nil
}
// ScheduleFor arma el Schedule de un pipeline a partir de sus variables de entorno opcionales (SYNC_<NOMBRE>_MODE/_INTERVAL/_TIME).
func ScheduleFor(pipelineName string, defaultInterval time.Duration) (appsync.Schedule, error) {
prefix := "SYNC_" + strings.ToUpper(pipelineName)
mode := os.Getenv(prefix + "_MODE")
if mode == "" {
mode = "interval"
}
switch mode {
case "interval":
interval := defaultInterval
if raw := os.Getenv(prefix + "_INTERVAL"); raw != "" {
d, err := time.ParseDuration(raw)
if err != nil {
return nil, fmt.Errorf("invalid %s_INTERVAL %q: %w", prefix, raw, err)
}
interval = d
}
return appsync.IntervalSchedule{Interval: interval}, nil
case "daily":
raw := os.Getenv(prefix + "_TIME")
if raw == "" {
return nil, fmt.Errorf("%s_MODE=daily requires %s_TIME (formato HH:MM)", prefix, prefix)
}
hour, minute, err := parseHHMM(raw)
if err != nil {
return nil, fmt.Errorf("invalid %s_TIME %q: %w", prefix, raw, err)
}
return appsync.DailySchedule{Hour: hour, Minute: minute}, nil
default:
return nil, fmt.Errorf("invalid %s_MODE %q: debe ser \"interval\" o \"daily\"", prefix, mode)
}
}
// parseHHMM parsea una hora de reloj tipo "HH:MM" (24 horas).
func parseHHMM(raw string) (hour, minute int, err error) {
parts := strings.SplitN(raw, ":", 2)
if len(parts) != 2 {
return 0, 0, fmt.Errorf("formato esperado HH:MM")
}
hour, err = strconv.Atoi(parts[0])
if err != nil || hour < 0 || hour > 23 {
return 0, 0, fmt.Errorf("hora inválida")
}
minute, err = strconv.Atoi(parts[1])
if err != nil || minute < 0 || minute > 59 {
return 0, 0, fmt.Errorf("minuto inválido")
}
return hour, minute, nil
}
package config
// PipelineDef describe UNA tubería de sincronización (un "pipeline").
//
// La idea de diseño: en vez de meter cada sincronización nueva a mano en el
// código, se declaran como datos en la lista Pipelines (más abajo). Agregar
// una sincronización nueva = agregar una entrada a esa lista. El arranque
// (cmd/server/main.go) recorre la lista y arma todo automáticamente.
//
// Cada pipeline lee datos de PostgreSQL y los vuelca en una colección de Mongo.
import (
"intranet-sycronizacion/internal/parents"
"intranet-sycronizacion/internal/paymentplans"
"intranet-sycronizacion/internal/students"
"intranet-sycronizacion/internal/users"
)
// PipelineDef describe una tubería de sincronización: origen Postgres -> destino Mongo.
type PipelineDef struct {
Name string // nombre lógico, ej. "students"
PGFunction string // función PostgreSQL a llamar (si no se usa PGQuery)
PGQuery string // SQL crudo; si está seteado, se usa EN VEZ de PGFunction
MongoCollection string // colección Mongo destino, ej. "students"
MongoDeletedCollection string // colección archivo para registros borrados en el origen
IDField string // nombre del campo id, ej. "student_id"
StateID int // id de la fila en la tabla sync_state (students=1)
RequireEnrollmentID bool // caso especial: exigir enrollment_id (solo estudiantes)
Name string
PGFunction string
PGQuery string
MongoCollection string
MongoDeletedCollection string
IDField string
StateID int
RequireEnrollmentID bool
}
// studentsQuery es el SQL crudo de la tubería de estudiantes.
//
// En Go, el texto entre comillas invertidas (`...`) es un string "en crudo":
// puede ocupar varias líneas y no interpreta caracteres de escape. Ideal
// para pegar SQL tal cual.
//
// Qué hace el SQL: junta estudiante + persona + matrícula + plan de pago,
// agrupa por estudiante, y arma dos listas JSON de planes de pago separando
// por periodo_academico_id = 14 (año actual) vs. el resto.
const studentsQuery = `
SELECT e.estudiante_id as student_id,
pp.persona_apellido_paterno as student_paternal_last_name,
pp.persona_apellido_materno as student_maternal_last_name,
pp.persona_nombre as student_name,
pp.persona_numero_documento_identidad AS student_dni,
MAX(m.matricula_id) as enrollment_id,
JSON_AGG(
json_build_object(
'payment_plan_id', cpp.plan_de_pago_id,
'payment_plan_amount', cpp.plan_de_pago_subtotal,
'payment_plan_debt', cpp.plan_de_pago_deuda,
'payment_plan_payment_date', cpp.plan_de_pago_fecha_pago,
'payment_plan_due_date', cpp.plan_de_pago_fecha_ven,
'payment_plan_year', cpp.plan_de_pago_anio
)
) FILTER (WHERE m.periodo_academico_id <> 14) as payment_plans,
JSON_AGG(
json_build_object(
'payment_plan_id', cpp.plan_de_pago_id,
'payment_plan_amount', cpp.plan_de_pago_subtotal,
'payment_plan_debt', cpp.plan_de_pago_deuda,
'payment_plan_payment_date', cpp.plan_de_pago_fecha_pago,
'payment_plan_due_date', cpp.plan_de_pago_fecha_ven,
'payment_plan_year', cpp.plan_de_pago_anio
)
) FILTER (WHERE m.periodo_academico_id = 14) as payment_plans_current_year
FROM matricula.ma_estudiante e
INNER JOIN persona.pe_persona pp on e.persona_id = pp.persona_id
INNER JOIN matricula.ma_matricula m on e.estudiante_id = m.estudiante_id
INNER JOIN caja.ca_plan_de_pago cpp on m.matricula_id = cpp.matricula_id
GROUP BY 1, 2, 3, 4, 5
ORDER BY 1 DESC;
`
// Pipelines es la lista de todas las tuberías activas.
//
// []PipelineDef significa "slice (lista) de PipelineDef". Hoy solo hay una
// entrada (estudiantes). Para sumar otra sincronización, se agrega otra
// entrada aquí; no hace falta tocar el resto del código.
var Pipelines = []PipelineDef{
{
Name: "students",
PGQuery: studentsQuery, // usa SQL crudo (no una función PG)
PGQuery: students.Query,
MongoCollection: "students",
MongoDeletedCollection: "deleted_students",
IDField: "student_id",
StateID: 1,
RequireEnrollmentID: true, // los estudiantes exigen enrollment_id
},
{
Name: "payment_plans",
PGQuery: paymentplans.Query,
MongoCollection: "payment_plans",
MongoDeletedCollection: "deleted_payment_plans",
IDField: "payment_plan_id",
StateID: 2,
},
{
Name: "parents",
PGQuery: parents.Query,
MongoCollection: "parents",
MongoDeletedCollection: "deleted_parents",
IDField: "parent_id",
StateID: 3,
},
{
Name: "users",
PGQuery: users.Query,
MongoCollection: "users",
MongoDeletedCollection: "deleted_users",
IDField: "user_id",
StateID: 4,
},
}
......@@ -6,9 +6,11 @@ import (
"fmt"
"time"
"go.mongodb.org/mongo-driver/bson" // bson: el formato binario de Mongo (como JSON)
"go.mongodb.org/mongo-driver/mongo" // cliente oficial de MongoDB
"go.mongodb.org/mongo-driver/mongo/options" // opciones para operaciones (upsert, projection, etc.)
"go.mongodb.org/mongo-driver/bson"
"go.mongodb.org/mongo-driver/mongo"
"go.mongodb.org/mongo-driver/mongo/options"
appsync "intranet-sycronizacion/internal/sync"
)
// NewMongoClient conecta a MongoDB y verifica con un Ping.
......@@ -23,87 +25,98 @@ func NewMongoClient(ctx context.Context, uri string) (*mongo.Client, error) {
return client, nil
}
// CollectionUpserter escribe en una colección de Mongo. Cumple DOS interfaces
// del dominio a la vez: MongoUpserter (Upsert) y MongoIDLister (ListIDs).
// CollectionUpserter escribe en una colección de Mongo, indexando por idField (no por _id).
type CollectionUpserter struct {
coll *mongo.Collection
idField string
}
func NewCollectionUpserter(coll *mongo.Collection) *CollectionUpserter {
return &CollectionUpserter{coll: coll}
func NewCollectionUpserter(coll *mongo.Collection, idField string) *CollectionUpserter {
return &CollectionUpserter{coll: coll, idField: idField}
}
// Upsert inserta o actualiza el documento por su _id.
// isNew indica si fue una inserción real (no existía) vs. una actualización.
func (c *CollectionUpserter) Upsert(ctx context.Context, id int, doc map[string]interface{}) (bool, error) {
filter := bson.M{"_id": id} // buscar por _id
update := bson.M{"$set": doc} // $set: pisar los campos con los del doc
opts := options.Update().SetUpsert(true) // upsert: si no existe, insertarlo
res, err := c.coll.UpdateOne(ctx, filter, update, opts)
// BulkUpsert manda todas las filas en un único BulkWrite (un round trip) en vez de un UpdateOne
// por fila, evitando la latencia de red repetida de escribir de a una.
func (c *CollectionUpserter) BulkUpsert(ctx context.Context, rows []appsync.SyncRow) (created, updated int, err error) {
if len(rows) == 0 {
return 0, 0, nil
}
models := make([]mongo.WriteModel, len(rows))
for i, row := range rows {
filter := bson.M{c.idField: row.ID}
update := bson.M{"$set": row.Doc}
models[i] = mongo.NewUpdateOneModel().SetFilter(filter).SetUpdate(update).SetUpsert(true)
}
res, err := c.coll.BulkWrite(ctx, models, options.BulkWrite().SetOrdered(false))
if err != nil {
return false, err
return 0, 0, err
}
// UpsertedCount > 0 significa que se creó un documento nuevo.
return res.UpsertedCount > 0, nil
created = int(res.UpsertedCount)
updated = len(rows) - created
return created, updated, nil
}
// ListIDs devuelve todos los _id que hay en la colección. Se usa para detectar
// qué registros se borraron del origen desde el último ciclo.
func (c *CollectionUpserter) ListIDs(ctx context.Context) ([]int, error) {
// Projection {_id: 1}: traer SOLO el campo _id, no el documento entero (más liviano).
cur, err := c.coll.Find(ctx, bson.M{}, options.Find().SetProjection(bson.M{"_id": 1}))
// ListIDsAndHashes devuelve, por cada documento existente, su id de negocio y su row_hash
// guardado (vacío si el documento no tiene hash). Se usa para detectar borrados del origen
// (las claves del mapa) y para saltar upserts de filas sin cambios (comparando el hash).
func (c *CollectionUpserter) ListIDsAndHashes(ctx context.Context) (map[int]string, error) {
cur, err := c.coll.Find(ctx, bson.M{}, options.Find().SetProjection(bson.M{c.idField: 1, "row_hash": 1}))
if err != nil {
return nil, err
}
defer cur.Close(ctx) // cerrar el cursor al terminar
defer cur.Close(ctx)
var ids []int
hashes := make(map[int]string)
for cur.Next(ctx) {
// Decodificamos cada documento en un struct temporal con solo el _id.
var doc struct {
ID int `bson:"_id"`
}
var doc bson.M
if err := cur.Decode(&doc); err != nil {
return nil, err
}
ids = append(ids, doc.ID)
var id int
switch v := doc[c.idField].(type) {
case int32:
id = int(v)
case int64:
id = int(v)
default:
continue
}
hash, _ := doc["row_hash"].(string)
hashes[id] = hash
}
return ids, cur.Err()
return hashes, cur.Err()
}
// DeletedMover archiva un registro que desapareció del origen: lo mueve de la
// colección viva a la de borrados (ej. deleted_students), estampando deleted_at.
// DeletedMover archiva un registro desaparecido del origen a su colección de borrados, estampando deleted_at.
type DeletedMover struct {
source *mongo.Collection // colección viva
deleted *mongo.Collection // colección archivo
source *mongo.Collection
deleted *mongo.Collection
idField string
}
func NewDeletedMover(source, deleted *mongo.Collection) *DeletedMover {
return &DeletedMover{source: source, deleted: deleted}
func NewDeletedMover(source, deleted *mongo.Collection, idField string) *DeletedMover {
return &DeletedMover{source: source, deleted: deleted, idField: idField}
}
// MoveToDeleted copia el documento a la colección de borrados y lo elimina de
// la viva. Es idempotente: si el documento ya no está, no hace nada.
// MoveToDeleted copia el documento a la colección de borrados y lo elimina de la viva (idempotente).
func (m *DeletedMover) MoveToDeleted(ctx context.Context, id int) error {
// 1. Leer el documento original de la colección viva.
var doc bson.M
err := m.source.FindOne(ctx, bson.M{"_id": id}).Decode(&doc)
filter := bson.M{m.idField: id}
err := m.source.FindOne(ctx, filter).Decode(&doc)
if errors.Is(err, mongo.ErrNoDocuments) {
// Ya no existe (quizá lo movió otro ciclo): nada que hacer.
return nil
}
if err != nil {
return fmt.Errorf("find source doc: %w", err)
}
// 2. Marcar cuándo se archivó.
delete(doc, "_id")
doc["deleted_at"] = time.Now().UTC()
filter := bson.M{"_id": id}
// 3. Escribir (upsert) el documento en la colección de borrados.
if _, err := m.deleted.UpdateOne(ctx, filter, bson.M{"$set": doc}, options.Update().SetUpsert(true)); err != nil {
return fmt.Errorf("archive doc: %w", err)
}
// 4. Borrarlo de la colección viva.
if _, err := m.source.DeleteOne(ctx, filter); err != nil {
return fmt.Errorf("delete source doc: %w", err)
}
......
// Package db contiene las implementaciones CONCRETAS de acceso a datos
// (PostgreSQL y MongoDB). Aquí sí se usan las librerías reales (pgx, mongo).
// Estos tipos cumplen las interfaces definidas en internal/sync, y se
// "inyectan" al dominio desde cmd/server/main.go.
// Package db contiene las implementaciones concretas de acceso a datos (PostgreSQL y MongoDB).
package db
import (
......@@ -9,14 +6,12 @@ import (
"encoding/json"
"fmt"
"github.com/jackc/pgx/v5/pgxpool" // pool de conexiones a Postgres
"github.com/jackc/pgx/v5/pgxpool"
appsync "intranet-sycronizacion/internal/sync" // nuestro paquete de dominio (renombrado a appsync)
appsync "intranet-sycronizacion/internal/sync"
)
// Pool envuelve el pool de conexiones de pgx para poder colgarle métodos
// propios. El "*pgxpool.Pool" embebido (sin nombre de campo) hace que Pool
// herede todos los métodos del pool original.
// Pool envuelve el pool de pgx para poder colgarle métodos propios.
type Pool struct {
*pgxpool.Pool
}
......@@ -33,8 +28,7 @@ func NewPostgresPool(ctx context.Context, dsn string) (*Pool, error) {
return &Pool{pool}, nil
}
// QueryRow adapta el QueryRow de pgx a la interfaz sync.Row que espera el
// dominio. Devuelve algo con método Scan.
// QueryRow adapta el QueryRow de pgx a la interfaz sync.Row del dominio.
func (p *Pool) QueryRow(ctx context.Context, sql string, args ...interface{}) appsync.Row {
return p.Pool.QueryRow(ctx, sql, args...)
}
......@@ -45,8 +39,7 @@ func (p *Pool) Exec(ctx context.Context, sql string, args ...interface{}) error
return err
}
// bootstrapSyncStateSQL crea la tabla sync_state si no existe y siembra la
// fila id=1 (estudiantes). Se corre al arrancar, así no hace falta migración manual.
// bootstrapSyncStateSQL crea la tabla sync_state si no existe y siembra las filas iniciales.
const bootstrapSyncStateSQL = `
CREATE TABLE IF NOT EXISTS sync_state (
id INTEGER PRIMARY KEY,
......@@ -54,7 +47,7 @@ CREATE TABLE IF NOT EXISTS sync_state (
);
INSERT INTO sync_state (id, last_synced_at)
VALUES (1, '1970-01-01T00:00:00Z')
VALUES (1, '1970-01-01T00:00:00Z'), (2, '1970-01-01T00:00:00Z'), (3, '1970-01-01T00:00:00Z')
ON CONFLICT (id) DO NOTHING;
`
......@@ -67,15 +60,7 @@ func (p *Pool) BootstrapSchema(ctx context.Context) error {
return nil
}
// ---------------------------------------------------------------------------
// Lectores (readers). Hay tres formas de leer el origen; todas cumplen la
// interfaz sync.PGReader (método FetchAll). main.go elige cuál usar según
// cómo esté declarado el pipeline.
// ---------------------------------------------------------------------------
// FunctionReader (heredado): llama una función de Postgres que devuelve un
// "sobre" JSON, y exige student_id + enrollment_id. Se conserva para pipelines
// aún no migrados a consultas crudas.
// FunctionReader (legado) llama una función Postgres que devuelve un sobre JSON, exige student_id + enrollment_id.
type FunctionReader struct {
pool *Pool
functionName string
......@@ -85,17 +70,14 @@ func NewFunctionReader(pool *Pool, functionName string) *FunctionReader {
return &FunctionReader{pool: pool, functionName: functionName}
}
// functionEnvelope es la forma del JSON que devuelven las funciones de PG:
// { "status": bool, "message": text, "data": [ {...}, {...} ] }.
// Las etiquetas `json:"..."` le dicen a Go cómo mapear cada campo del JSON.
// functionEnvelope es la forma del JSON que devuelven las funciones PG: status/message/data.
type functionEnvelope struct {
Status bool `json:"status"`
Message string `json:"message"`
Data []map[string]interface{} `json:"data"`
}
// parseFunctionEnvelope convierte el JSON crudo en filas listas para sincronizar,
// exigiendo que cada item tenga student_id y enrollment_id numéricos.
// parseFunctionEnvelope convierte el JSON crudo en filas, exigiendo student_id y enrollment_id numéricos.
func parseFunctionEnvelope(raw []byte) ([]appsync.SyncRow, error) {
var env functionEnvelope
if err := json.Unmarshal(raw, &env); err != nil {
......@@ -107,8 +89,6 @@ func parseFunctionEnvelope(raw []byte) ([]appsync.SyncRow, error) {
rows := make([]appsync.SyncRow, 0, len(env.Data))
for _, item := range env.Data {
// El JSON no distingue enteros de decimales: todo número llega como float64.
// Por eso convertimos a int a mano y verificamos el tipo (comma-ok idiom).
idFloat, ok := item["student_id"].(float64)
if !ok {
return nil, fmt.Errorf("data item missing numeric student_id: %v", item)
......@@ -134,9 +114,7 @@ func (f *FunctionReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error
return parseFunctionEnvelope(raw)
}
// parseGenericEnvelope es como parseFunctionEnvelope pero exige un solo campo
// id configurable (idField), sin obligar enrollment_id. Sirve para pipelines
// genéricos que solo necesitan validar un id numérico.
// parseGenericEnvelope es como parseFunctionEnvelope pero exige solo un id configurable (idField).
func parseGenericEnvelope(raw []byte, idField string) ([]appsync.SyncRow, error) {
var env functionEnvelope
if err := json.Unmarshal(raw, &env); err != nil {
......@@ -158,8 +136,7 @@ func parseGenericEnvelope(raw []byte, idField string) ([]appsync.SyncRow, error)
return rows, nil
}
// GenericFunctionReader llama una función de Postgres cuyo sobre solo requiere
// un id numérico (a diferencia de FunctionReader, que también exige enrollment_id).
// GenericFunctionReader llama una función Postgres cuyo sobre solo requiere un id numérico.
type GenericFunctionReader struct {
pool *Pool
functionName string
......@@ -179,85 +156,4 @@ func (f *GenericFunctionReader) FetchAll(ctx context.Context) ([]appsync.SyncRow
return parseGenericEnvelope(raw, f.idField)
}
// StudentsQueryReader ejecuta el SQL crudo de estudiantes directamente (en vez
// de envolver una función que devuelve JSON) y arma una SyncRow por cada fila
// del resultado. Es el que usa el pipeline de estudiantes.
type StudentsQueryReader struct {
pool *Pool
query string
}
func NewStudentsQueryReader(pool *Pool, query string) *StudentsQueryReader {
return &StudentsQueryReader{pool: pool, query: query}
}
// FetchAll corre la consulta y lee fila por fila.
func (f *StudentsQueryReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) {
rows, err := f.pool.Pool.Query(ctx, f.query)
if err != nil {
return nil, fmt.Errorf("query students: %w", err)
}
defer rows.Close() // cerrar el cursor al terminar, pase lo que pase
var result []appsync.SyncRow
for rows.Next() { // avanza a la siguiente fila; false cuando no quedan más
// Los punteros (*string) permiten que el valor sea NULL en la base:
// si la columna es NULL, la variable queda en nil.
var (
studentID int
paternalLastName *string
maternalLastName *string
name *string
dni *string
enrollmentID int
paymentPlans []byte // JSON crudo devuelto por JSON_AGG
paymentPlansCurrentYr []byte
)
// Scan copia cada columna, EN ORDEN, a las variables (por eso van con &).
if err := rows.Scan(&studentID, &paternalLastName, &maternalLastName, &name, &dni,
&enrollmentID, &paymentPlans, &paymentPlansCurrentYr); err != nil {
return nil, fmt.Errorf("scan student row: %w", err)
}
// Convertimos el JSON de planes de pago a una lista usable.
plans, err := decodePaymentPlans(paymentPlans)
if err != nil {
return nil, fmt.Errorf("decode payment_plans for student %d: %w", studentID, err)
}
plansCurrentYr, err := decodePaymentPlans(paymentPlansCurrentYr)
if err != nil {
return nil, fmt.Errorf("decode payment_plans_current_year for student %d: %w", studentID, err)
}
// Armamos el documento que irá a Mongo.
doc := map[string]interface{}{
"student_id": studentID,
"student_paternal_last_name": paternalLastName,
"student_maternal_last_name": maternalLastName,
"student_name": name,
"student_dni": dni,
"enrollment_id": enrollmentID,
"payment_plans": plans,
"payment_plans_current_year": plansCurrentYr,
}
result = append(result, appsync.SyncRow{ID: studentID, Doc: doc})
}
// rows.Err reporta errores ocurridos DURANTE la iteración del cursor.
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate student rows: %w", err)
}
return result, nil
}
// decodePaymentPlans convierte el JSON crudo (bytes) en una lista de planes.
// Si viene NULL (nil), devuelve una lista vacía en vez de error.
func decodePaymentPlans(raw []byte) ([]interface{}, error) {
if raw == nil {
return []interface{}{}, nil
}
var plans []interface{}
if err := json.Unmarshal(raw, &plans); err != nil {
return nil, err
}
return plans, nil
}
// Los readers por módulo (students, payment_plans, parents) viven en internal/<módulo>.
......@@ -7,54 +7,41 @@ import (
)
// Runner es cualquier cosa que sepa ejecutar un ciclo. El Service lo cumple.
// El Scheduler depende de esta interfaz (no del Service concreto), así que se
// lo puede probar con un runner falso.
type Runner interface {
Run(ctx context.Context) error
}
// cycleTimeout: tope máximo de tiempo para un solo ciclo. Si un ciclo se cuelga,
// se corta a los 4 minutos y no bloquea para siempre. Es independiente del
// intervalo entre ciclos.
// cycleTimeout: tope máximo de tiempo para un solo ciclo, independiente del intervalo entre ciclos.
const cycleTimeout = 4 * time.Minute
// Scheduler dispara el Runner cada "interval" de tiempo.
// Scheduler dispara el Runner según su Schedule (intervalo fijo u hora fija del día).
type Scheduler struct {
runner Runner
interval time.Duration
schedule Schedule
}
func NewScheduler(runner Runner, interval time.Duration) *Scheduler {
return &Scheduler{runner: runner, interval: interval}
func NewScheduler(runner Runner, schedule Schedule) *Scheduler {
return &Scheduler{runner: runner, schedule: schedule}
}
// Start corre en bucle disparando un ciclo en cada "tick" del reloj, hasta que
// el context (ctx) se cancele (ej. al apagar el servicio).
//
// Normalmente se llama con "go scheduler.Start(ctx)": la palabra clave "go"
// lo lanza en una goroutine (hilo liviano) para que no bloquee el arranque.
// Start corre en bucle disparando un ciclo cada vez que el Schedule lo indica, hasta cancelar ctx.
func (s *Scheduler) Start(ctx context.Context) {
// ticker.C es un canal que emite un valor cada "interval".
ticker := time.NewTicker(s.interval)
defer ticker.Stop() // liberar el ticker al salir
timer := time.NewTimer(s.schedule.Next(time.Now()))
defer timer.Stop()
for {
// select espera sobre varios canales a la vez y actúa según cuál dispare.
select {
case <-ctx.Done():
// El context se canceló: es hora de apagar. Salimos del bucle.
return
case <-ticker.C:
// Pasó el intervalo: corre un ciclo.
case <-timer.C:
s.runOnce(ctx)
timer.Reset(s.schedule.Next(time.Now()))
}
}
}
// runOnce ejecuta un ciclo con su propio timeout (cycleTimeout) y registra el
// error si falla, sin detener el bucle: el próximo tick vuelve a intentar.
// runOnce ejecuta un ciclo con su propio timeout, registrando el error sin detener el bucle.
func (s *Scheduler) runOnce(ctx context.Context) {
// Creamos un context "hijo" que se cancela solo a los cycleTimeout.
runCtx, cancel := context.WithTimeout(ctx, cycleTimeout)
defer cancel()
if err := s.runner.Run(runCtx); err != nil {
......
// Package sync es el "dominio" del servicio: la lógica central de la
// sincronización, escrita SIN depender de librerías externas de base de datos
// ni del framework web. Solo depende de interfaces pequeñas (definidas aquí),
// y las implementaciones reales (Postgres/Mongo) se inyectan desde afuera.
//
// Esto es el corazón de la "arquitectura limpia": el dominio no sabe qué base
// de datos concreta se usa; solo sabe "necesito algo que pueda leer filas",
// "algo que pueda guardar". Así se puede testear sin bases de datos reales.
// Package sync es el dominio del servicio: lógica de sincronización sin depender de librerías externas de base de datos ni del framework web.
package sync
import (
......@@ -13,51 +6,33 @@ import (
"time"
)
// ---------------------------------------------------------------------------
// Interfaces = "contratos". Definen QUÉ se necesita, no CÓMO se hace.
// El paquete internal/db provee las implementaciones concretas.
// ---------------------------------------------------------------------------
// Row abstrae el resultado de una consulta de una sola fila.
// (En la práctica lo satisface pgx.Row de la librería de Postgres.)
// Row abstrae el resultado de una consulta de una sola fila (lo satisface pgx.Row).
type Row interface {
Scan(dest ...interface{}) error // copia los valores de la fila a las variables dadas
Scan(dest ...interface{}) error
}
// PGExecutor abstrae el acceso mínimo a Postgres que necesita el state store.
// Al ser una interfaz, el state store no depende directamente de pgx.
type PGExecutor interface {
QueryRow(ctx context.Context, sql string, args ...interface{}) Row
Exec(ctx context.Context, sql string, args ...interface{}) error
}
// StateStore guarda y recupera "cuándo fue la última sincronización exitosa".
// Se usa para el endpoint /sync/students/status y para retomar estado tras reiniciar.
// StateStore guarda y recupera cuándo fue la última sincronización exitosa.
type StateStore interface {
Get(ctx context.Context) (time.Time, error)
Set(ctx context.Context, t time.Time) error
}
// ---------------------------------------------------------------------------
// Implementación concreta sobre PostgreSQL.
// ---------------------------------------------------------------------------
// PGStateStore persiste el timestamp de última sincronización en la tabla
// sync_state. Hay una fila por pipeline, identificada por stateID
// (estudiantes = 1).
// PGStateStore persiste el timestamp de última sincronización en la tabla sync_state, una fila por pipeline (stateID).
type PGStateStore struct {
db PGExecutor
stateID int
}
// NewPGStateStore construye un PGStateStore. En Go es convención tener una
// función "NewXxx" que crea e inicializa el struct (no hay constructores como
// en otros lenguajes). Devuelve un puntero (*PGStateStore) para no copiar.
func NewPGStateStore(db PGExecutor, stateID int) *PGStateStore {
return &PGStateStore{db: db, stateID: stateID}
}
// $1 es un parámetro posicional de PostgreSQL (evita inyección SQL).
const selectLastSyncedSQL = `SELECT last_synced_at FROM sync_state WHERE id = $1`
// Get lee el último timestamp sincronizado para este pipeline.
......@@ -65,14 +40,12 @@ func (s *PGStateStore) Get(ctx context.Context) (time.Time, error) {
var t time.Time
row := s.db.QueryRow(ctx, selectLastSyncedSQL, s.stateID)
if err := row.Scan(&t); err != nil {
// time.Time{} es el "valor cero" (fecha vacía). Se devuelve junto al error.
return time.Time{}, err
}
return t, nil
}
// Inserta la fila si no existe, o actualiza last_synced_at si ya existe
// (patrón "upsert" con ON CONFLICT).
// upsertLastSyncedSQL inserta la fila si no existe o actualiza last_synced_at si ya existe.
const upsertLastSyncedSQL = `
INSERT INTO sync_state (id, last_synced_at) VALUES ($1, $2)
ON CONFLICT (id) DO UPDATE SET last_synced_at = EXCLUDED.last_synced_at`
......
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