initial commit

parents
# CLAUDE.md
This file provides guidance to Claude Code (claude.ai/code) when working with code in this repository.
## What this service does
Runs two independent full-sync pipelines every `SYNC_INTERVAL`, full-syncing the result set into a MongoDB collection: students (raw SQL query joining `matricula.ma_estudiante`/`persona.pe_persona`/`matricula.ma_matricula`/`caja.ca_plan_de_pago`, upsert by `student_id`, collection `students`) and parents (Postgres function `func_apoderado_listar`, upsert by `parent_id`, collection `parents`). Sincroniza estudiantes (con padres y planes de pago) y apoderados desde PostgreSQL hacia MongoDB. No incremental/CDC logic — every cycle re-reads the full dataset for that pipeline and reconciles Mongo against it (upsert what's present, archive what's gone).
## Commands
```
go build ./...
go vet ./...
go run ./cmd/server
```
No Makefile/linter config present, no test suite — use plain `go` toolchain commands. Verify changes with `go build ./...` / `go vet ./...` and by running the service against real/staging Postgres and Mongo.
## Setup / running locally
1. `sync_state` table is auto-bootstrapped at startup (`Pool.BootstrapSchema`, seeds `id=1` students / `id=2` parents rows) — no manual migration step required, though the SQL files under `migrations/` are kept for reference/manual runs.
2. Mongo `students` collection must exist with a `$jsonSchema` validator requiring `student_id` int and `enrollment_id` int. A `deleted_students` collection archives students removed at the source — no schema setup needed for it.
3. Mongo `parents` collection must exist with a `$jsonSchema` validator requiring `parent_id` int. A `deleted_parents` collection archives parents removed at the source — no schema setup needed for it.
4. Required env vars: `PG_DSN`, `MONGO_URI`, `MONGO_DB`. Optional: `SYNC_INTERVAL` (default `5m`, shared by all pipelines), `PORT` (default `8080`). Per-pipeline PG function name (or raw `PGQuery`) / Mongo collection / deleted-collection / id field / `sync_state` id are **not** env vars — they're declared in `internal/config/pipelines.go` (`Pipelines []PipelineDef`). Add a new sync pipeline by adding an entry there; `cmd/server/main.go` loops over `config.Pipelines` to wire reader/upserter/mover/state/service/scheduler for each.
5. `go run ./cmd/server`
## Endpoints
- `GET /health` — checks Postgres and Mongo connectivity.
- `POST /sync/trigger` / `GET /sync/status` — students pipeline: run now / last run time, rows synced (created/updated/deleted breakdown), last error if any.
- `POST /sync/parents/trigger` / `GET /sync/parents/status` — same, for the parents pipeline.
## Postgres source contract
Students: `internal/config.studentsQuery`, a raw SQL query (not a function) joining `matricula.ma_estudiante`/`persona.pe_persona`/`matricula.ma_matricula`/`caja.ca_plan_de_pago`, grouped per student, `payment_plans`/`payment_plans_current_year` built via `JSON_AGG(...) FILTER (...)` split on `periodo_academico_id = 14`. Read directly by `db.StudentsQueryReader` (`pgx` row scan, no JSON envelope) — one Mongo doc per row, upserted by `_id = student_id`, requires numeric `enrollment_id`.
Parents: `PG_APODERADOS_FUNCTION` (`func_apoderado_listar()`), takes no params, returns one JSON envelope: `{"status": bool, "message": text, "data": [...]}`, unwrapped by `db.GenericFunctionReader`. Each `data[]` item → one Mongo doc, upserted by `_id = parent_id`.
`updated_at` is stamped by this service for both pipelines (Postgres does not supply it).
## Sync behavior
Each cycle (students and parents run independently) processes rows concurrently via a bounded worker pool (`defaultConcurrency` = 10, `errgroup`-based) so a full migration finishes faster than a sequential loop:
- **New record**: id not yet in Mongo → inserted, counted as `Created`.
- **Existing record**: id already in Mongo → fields overwritten, counted as `Updated`.
- **Removed record**: id present in Mongo but no longer returned by the Postgres source (query or function) → moved (not just deleted) into the pipeline's archive collection (`deleted_students` / `deleted_parents`) with a `deleted_at` stamp, counted as `Deleted`.
Persisted `sync_state` only advances if the entire cycle (upserts + delete reconciliation) succeeds; any failure short-circuits and is surfaced via `/sync/status`.
## Architecture
Layering, dependency direction is `api` / `sync` → small interfaces, with concrete Postgres/Mongo clients living in `internal/db` and injected from `cmd/server/main.go`:
- `internal/config``config.go`: env var loading/validation (`config.Load()`, only `PGDSN`/`MongoURI`/`MongoDB`/`SyncInterval`/`Port`), fails fast on missing required vars. `pipelines.go`: `Pipelines []PipelineDef` — static Go-defined list of sync pipelines (PG function or raw `PGQuery`, Mongo collection, deleted collection, id field, `sync_state` id, `RequireEnrollmentID` flag). Adding a pipeline means adding an entry here, not new env vars. Also holds `studentsQuery`, the raw SQL for the students pipeline.
- `internal/db``NewPostgresPool`, `NewMongoClient`, `StudentsQueryReader` (students-specific: runs `PGQuery` directly via `pgx` row scan, no JSON envelope, requires `student_id`+`enrollment_id`), `FunctionReader` (legacy: calls a PG function returning a JSON envelope, unwraps into `[]sync.SyncRow`, requires `student_id`+`enrollment_id` — kept for pipelines not yet migrated to raw queries), `GenericFunctionReader` (parents-and-beyond: same envelope unwrap, validates only a configurable id field), `CollectionUpserter` (Mongo upsert by int id, also implements `ListIDs` for delete reconciliation), `DeletedMover` (moves a doc from a source collection into an archive collection, stamping `deleted_at`).
- `internal/sync` — core domain, framework-free:
- `service.go`: `Service.Run(ctx)` does one full cycle — fetch all rows, concurrently upsert each with `updated_at` stamped (tracking created vs. updated counts via `MongoUpserter`), then reconcile deletions by diffing Mongo's existing ids (`MongoIDLister`) against the current Postgres id set and archiving the difference (`DeletedMover`). Advances persisted state only on full success; stops and records the error in `lastResult` otherwise. `LastResult()`/`SeedLastResult()` back the `/sync/status` endpoint and boot-time state hydration.
- `scheduler.go`: `Scheduler.Start(ctx)` ticks every `interval` and calls `Run` with a hard `cycleTimeout` (4 min) sub-context per cycle, independent of the ticker interval.
- `state.go`: `StateStore` interface + `PGStateStore` (constructed with an explicit `stateID`: students=1, parents=2), persists last-synced timestamp in the `sync_state` table so `/sync/status`, `/sync/parents/status`, and boot seeding survive restarts for each pipeline.
- Depends only on small interfaces (`PGReader`, `MongoUpserter`, `MongoIDLister`, `DeletedMover`, `StateStore`) — no real DB/Mongo needed to exercise it.
- `internal/api` — Gin router (`NewRouter(studentsSvc, parentsSvc, pgPing, mongoPing)`) wiring `/health` (pings Postgres + Mongo), `POST /sync/trigger` + `GET /sync/status` (students), `POST /sync/parents/trigger` + `GET /sync/parents/status` (parents), each reading its own `Service.LastResult()`.
- `cmd/server/main.go` — composition root: builds config → DB clients → two independent reader/upserter/mover/state/`Service`/`Scheduler` sets (students, parents) → seeds each `Service`'s last result from its persisted state → starts both `Scheduler`s in goroutines → starts HTTP server → graceful shutdown on SIGINT/SIGTERM (10s timeout). The two pipelines share nothing at runtime beyond the Postgres pool and Mongo client — separate scheduler tick, separate `sync_state` row, separate failure domain.
# 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.
## 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).
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`
## 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.
## 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:
- **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`.
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.
// 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.
package main
import (
"context"
"errors"
"log"
"net/http"
"os"
"os/signal"
"syscall"
"time"
"github.com/joho/godotenv" // carga variables desde un archivo .env
"intranet-sycronizacion/internal/api"
"intranet-sycronizacion/internal/config"
"intranet-sycronizacion/internal/db"
appsync "intranet-sycronizacion/internal/sync"
)
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)
}
if err := pgPool.BootstrapSchema(bootCtx); err != nil {
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.
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.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)
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/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() {
svc.SeedLastResult(lastSynced)
}
// Arrancar el scheduler de este pipeline en su propia goroutine.
scheduler := appsync.NewScheduler(svc, cfg.SyncInterval)
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"],
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.
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if err := srv.Shutdown(shutdownCtx); err != nil {
log.Printf("server shutdown error: %v", err)
}
}
module intranet-sycronizacion
go 1.25.0
require (
github.com/gin-gonic/gin v1.12.0
github.com/jackc/pgx/v5 v5.10.0
github.com/joho/godotenv v1.5.1
github.com/stretchr/testify v1.11.1
go.mongodb.org/mongo-driver v1.17.9
golang.org/x/sync v0.19.0
)
require (
github.com/bytedance/gopkg v0.1.3 // indirect
github.com/bytedance/sonic v1.15.0 // indirect
github.com/bytedance/sonic/loader v0.5.0 // indirect
github.com/cloudwego/base64x v0.1.6 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/gabriel-vasile/mimetype v1.4.12 // indirect
github.com/gin-contrib/sse v1.1.0 // indirect
github.com/go-playground/locales v0.14.1 // indirect
github.com/go-playground/universal-translator v0.18.1 // indirect
github.com/go-playground/validator/v10 v10.30.1 // indirect
github.com/goccy/go-json v0.10.5 // indirect
github.com/goccy/go-yaml v1.19.2 // indirect
github.com/golang/snappy v0.0.4 // indirect
github.com/jackc/pgpassfile v1.0.0 // indirect
github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect
github.com/jackc/puddle/v2 v2.2.2 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/klauspost/compress v1.17.6 // indirect
github.com/klauspost/cpuid/v2 v2.3.0 // indirect
github.com/leodido/go-urn v1.4.0 // indirect
github.com/mattn/go-isatty v0.0.20 // indirect
github.com/modern-go/concurrent v0.0.0-20180306012644-bacd9c7ef1dd // indirect
github.com/modern-go/reflect2 v1.0.2 // indirect
github.com/montanaflynn/stats v0.7.1 // indirect
github.com/pelletier/go-toml/v2 v2.2.4 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/quic-go/qpack v0.6.0 // indirect
github.com/quic-go/quic-go v0.59.0 // indirect
github.com/twitchyliquid64/golang-asm v0.15.1 // indirect
github.com/ugorji/go/codec v1.3.1 // indirect
github.com/xdg-go/pbkdf2 v1.0.0 // indirect
github.com/xdg-go/scram v1.2.0 // indirect
github.com/xdg-go/stringprep v1.0.4 // indirect
github.com/youmark/pkcs8 v0.0.0-20240726163527-a2c0da244d78 // indirect
go.mongodb.org/mongo-driver/v2 v2.5.0 // indirect
golang.org/x/arch v0.22.0 // indirect
golang.org/x/crypto v0.48.0 // indirect
golang.org/x/net v0.51.0 // indirect
golang.org/x/sys v0.41.0 // indirect
golang.org/x/text v0.34.0 // indirect
google.golang.org/protobuf v1.36.10 // indirect
gopkg.in/yaml.v3 v3.0.1 // indirect
)
This diff is collapsed. Click to expand it.
package api
import (
"context"
"net/http"
"github.com/gin-gonic/gin"
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.
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
}
// 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).
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
body["postgres"] = pgErr.Error()
}
if mongoErr != nil {
status = http.StatusServiceUnavailable
body["mongo"] = mongoErr.Error()
}
c.JSON(status, body) // escribe status + cuerpo JSON
}
// 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.
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,
"deleted": result.Deleted,
}
if err != nil {
resp["error"] = err.Error()
}
return resp
}
// Trigger responde POST /sync/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
// (sin correr uno nuevo), incluyendo el último error si lo hubo.
func (h *Handlers) Status(c *gin.Context) {
result := h.studentsSvc.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
import (
"context"
"github.com/gin-gonic/gin"
appsync "intranet-sycronizacion/internal/sync"
)
// 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.
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/trigger", h.Trigger) // forzar una sincronización ahora
r.GET("/sync/status", h.Status) // ver el resultado del último ciclo
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)
)
// 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).
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)
}
// 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.
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"),
MongoDB: os.Getenv("MONGO_DB"),
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,
"MONGO_DB": cfg.MongoDB,
}
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
}
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.
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)
}
// 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)
MongoCollection: "students",
MongoDeletedCollection: "deleted_students",
IDField: "student_id",
StateID: 1,
RequireEnrollmentID: true, // los estudiantes exigen enrollment_id
},
}
package db
import (
"context"
"errors"
"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.)
)
// NewMongoClient conecta a MongoDB y verifica con un Ping.
func NewMongoClient(ctx context.Context, uri string) (*mongo.Client, error) {
client, err := mongo.Connect(ctx, options.Client().ApplyURI(uri))
if err != nil {
return nil, err
}
if err := client.Ping(ctx, nil); err != nil {
return nil, err
}
return client, nil
}
// CollectionUpserter escribe en una colección de Mongo. Cumple DOS interfaces
// del dominio a la vez: MongoUpserter (Upsert) y MongoIDLister (ListIDs).
type CollectionUpserter struct {
coll *mongo.Collection
}
func NewCollectionUpserter(coll *mongo.Collection) *CollectionUpserter {
return &CollectionUpserter{coll: coll}
}
// 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)
if err != nil {
return false, err
}
// UpsertedCount > 0 significa que se creó un documento nuevo.
return res.UpsertedCount > 0, 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}))
if err != nil {
return nil, err
}
defer cur.Close(ctx) // cerrar el cursor al terminar
var ids []int
for cur.Next(ctx) {
// Decodificamos cada documento en un struct temporal con solo el _id.
var doc struct {
ID int `bson:"_id"`
}
if err := cur.Decode(&doc); err != nil {
return nil, err
}
ids = append(ids, doc.ID)
}
return ids, 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.
type DeletedMover struct {
source *mongo.Collection // colección viva
deleted *mongo.Collection // colección archivo
}
func NewDeletedMover(source, deleted *mongo.Collection) *DeletedMover {
return &DeletedMover{source: source, deleted: deleted}
}
// 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.
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)
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ó.
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)
}
return nil
}
// 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
import (
"context"
"encoding/json"
"fmt"
"github.com/jackc/pgx/v5/pgxpool" // pool de conexiones a Postgres
appsync "intranet-sycronizacion/internal/sync" // nuestro paquete de dominio (renombrado a appsync)
)
// 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.
type Pool struct {
*pgxpool.Pool
}
// NewPostgresPool abre el pool y verifica la conexión con un Ping.
func NewPostgresPool(ctx context.Context, dsn string) (*Pool, error) {
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
return nil, err
}
if err := pool.Ping(ctx); err != nil {
return nil, err
}
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.
func (p *Pool) QueryRow(ctx context.Context, sql string, args ...interface{}) appsync.Row {
return p.Pool.QueryRow(ctx, sql, args...)
}
// Exec ejecuta una sentencia que no devuelve filas (INSERT/UPDATE/DDL).
func (p *Pool) Exec(ctx context.Context, sql string, args ...interface{}) error {
_, err := p.Pool.Exec(ctx, sql, args...)
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.
const bootstrapSyncStateSQL = `
CREATE TABLE IF NOT EXISTS sync_state (
id INTEGER PRIMARY KEY,
last_synced_at TIMESTAMPTZ NOT NULL
);
INSERT INTO sync_state (id, last_synced_at)
VALUES (1, '1970-01-01T00:00:00Z')
ON CONFLICT (id) DO NOTHING;
`
// BootstrapSchema prepara la tabla sync_state al inicio del servicio.
func (p *Pool) BootstrapSchema(ctx context.Context) error {
_, err := p.Pool.Exec(ctx, bootstrapSyncStateSQL)
if err != nil {
return fmt.Errorf("bootstrap sync_state schema: %w", err)
}
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.
type FunctionReader struct {
pool *Pool
functionName string
}
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.
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.
func parseFunctionEnvelope(raw []byte) ([]appsync.SyncRow, error) {
var env functionEnvelope
if err := json.Unmarshal(raw, &env); err != nil {
return nil, fmt.Errorf("unmarshal function envelope: %w", err)
}
if !env.Status {
return nil, fmt.Errorf("function reported failure: %s", env.Message)
}
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)
}
enrollmentFloat, ok := item["enrollment_id"].(float64)
if !ok {
return nil, fmt.Errorf("data item missing numeric enrollment_id: %v", item)
}
item["student_id"] = int(idFloat)
item["enrollment_id"] = int(enrollmentFloat)
rows = append(rows, appsync.SyncRow{ID: int(idFloat), Doc: item})
}
return rows, nil
}
// FetchAll llama la función PG y devuelve todas las filas.
func (f *FunctionReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) {
sql := "SELECT json FROM " + f.functionName + "()"
var raw []byte
if err := f.pool.Pool.QueryRow(ctx, sql).Scan(&raw); err != nil {
return nil, fmt.Errorf("query %s: %w", f.functionName, err)
}
return parseFunctionEnvelope(raw)
}
// 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).
func parseGenericEnvelope(raw []byte, idField string) ([]appsync.SyncRow, error) {
var env functionEnvelope
if err := json.Unmarshal(raw, &env); err != nil {
return nil, fmt.Errorf("unmarshal function envelope: %w", err)
}
if !env.Status {
return nil, fmt.Errorf("function reported failure: %s", env.Message)
}
rows := make([]appsync.SyncRow, 0, len(env.Data))
for _, item := range env.Data {
idFloat, ok := item[idField].(float64)
if !ok {
return nil, fmt.Errorf("data item missing numeric %s: %v", idField, item)
}
item[idField] = int(idFloat)
rows = append(rows, appsync.SyncRow{ID: int(idFloat), Doc: item})
}
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).
type GenericFunctionReader struct {
pool *Pool
functionName string
idField string
}
func NewGenericFunctionReader(pool *Pool, functionName, idField string) *GenericFunctionReader {
return &GenericFunctionReader{pool: pool, functionName: functionName, idField: idField}
}
func (f *GenericFunctionReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) {
sql := "SELECT json FROM " + f.functionName + "()"
var raw []byte
if err := f.pool.Pool.QueryRow(ctx, sql).Scan(&raw); err != nil {
return nil, fmt.Errorf("query %s: %w", f.functionName, err)
}
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
}
package sync
import (
"context"
"log"
"time"
)
// 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.
const cycleTimeout = 4 * time.Minute
// Scheduler dispara el Runner cada "interval" de tiempo.
type Scheduler struct {
runner Runner
interval time.Duration
}
func NewScheduler(runner Runner, interval time.Duration) *Scheduler {
return &Scheduler{runner: runner, interval: interval}
}
// 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.
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
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.
s.runOnce(ctx)
}
}
}
// 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.
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 {
log.Printf("sync cycle failed: %v", err)
}
}
package sync
import (
"context"
"fmt"
"sync" // primitivas de concurrencia (Mutex). OJO: es el paquete estándar,
"sync/atomic" // no confundir con NUESTRO paquete que también se llama "sync".
"time"
"golang.org/x/sync/errgroup" // grupo de goroutines con manejo de errores
)
// defaultConcurrency: cuántas filas se procesan a la vez (en paralelo).
// Un pool acotado: más rápido que ir de a una, sin saturar la base.
const defaultConcurrency = 10
// SyncRow es una fila lista para volcar a Mongo: su id y su documento.
// map[string]interface{} = un mapa de "nombre de campo" -> "valor de cualquier tipo".
type SyncRow struct {
ID int
Doc map[string]interface{}
}
// ---------------------------------------------------------------------------
// Interfaces del dominio. El Service depende solo de estos contratos, no de
// Postgres/Mongo concretos. Las implementaciones viven en internal/db.
// ---------------------------------------------------------------------------
// PGReader lee TODAS las filas del origen (Postgres) en un ciclo.
type PGReader interface {
FetchAll(ctx context.Context) ([]SyncRow, error)
}
// MongoUpserter inserta o actualiza un documento por id.
// isNew indica si el registro NO existía antes (inserción real) vs. si ya
// existía (actualización). Sirve para contar creados vs. actualizados.
type MongoUpserter interface {
Upsert(ctx context.Context, id int, doc map[string]interface{}) (isNew bool, err error)
}
// MongoIDLister lista todos los ids que hay actualmente en la colección Mongo,
// para detectar registros que desaparecieron de Postgres desde el último ciclo.
type MongoIDLister interface {
ListIDs(ctx context.Context) ([]int, error)
}
// DeletedMover archiva un registro que ya no está en Postgres: lo saca de la
// colección viva y lo mete en la colección de borrados (deleted_students).
type DeletedMover interface {
MoveToDeleted(ctx context.Context, id int) error
}
// SyncResult resume cómo salió el último ciclo. Alimenta /sync/status.
type SyncResult struct {
RanAt time.Time // cuándo corrió
RowsSynced int // total de filas procesadas
Created int // insertados nuevos
Updated int // actualizados existentes
Deleted int // archivados por desaparecer del origen
Err error // error si el ciclo falló, o nil si fue exitoso
}
// Service orquesta un ciclo completo de sincronización de UN pipeline.
type Service struct {
reader PGReader
mongo MongoUpserter
lister MongoIDLister
mover DeletedMover
state StateStore
now func() time.Time // función que da "la hora actual"; inyectable para tests
concurrency int
// mu protege el acceso a lastResult desde varias goroutines a la vez.
// Un Mutex es un candado: solo una goroutine puede tenerlo tomado.
mu sync.Mutex
lastResult SyncResult
}
// NewService arma un Service inyectándole todas sus dependencias (interfaces).
// Este estilo (inyección de dependencias) es lo que hace testeable al dominio.
func NewService(reader PGReader, mongo MongoUpserter, lister MongoIDLister, mover DeletedMover, state StateStore, now func() time.Time) *Service {
return &Service{
reader: reader,
mongo: mongo,
lister: lister,
mover: mover,
state: state,
now: now,
concurrency: defaultConcurrency,
}
}
// Run ejecuta UN ciclo completo de sincronización:
// 1. lee todas las filas de Postgres,
// 2. las inserta/actualiza en Mongo en paralelo (contando creados/actualizados),
// 3. archiva los registros que ya no están en Postgres,
// 4. avanza el estado persistido SOLO si todo salió bien.
//
// Devuelve error si algo falló (y lo deja registrado en lastResult para /status).
func (s *Service) Run(ctx context.Context) error {
// Tomamos el candado: garantiza que no corran dos ciclos a la vez y
// protege la escritura de lastResult. defer = ejecutar al salir de la función.
s.mu.Lock()
defer s.mu.Unlock()
runAt := s.now().UTC()
result := SyncResult{RanAt: runAt}
// --- Paso 1: leer todo desde Postgres ---
rows, err := s.reader.FetchAll(ctx)
if err != nil {
result.Err = fmt.Errorf("fetch all: %w", err)
s.lastResult = result
return result.Err
}
// Guardamos en un set (mapa) los ids que SÍ vinieron de Postgres.
// struct{} es un valor vacío (ocupa 0 bytes): el mapa se usa solo como conjunto.
pgIDs := make(map[int]struct{}, len(rows))
for _, row := range rows {
pgIDs[row.ID] = struct{}{}
}
// --- Paso 2: upsert en paralelo ---
// Contadores atómicos: se incrementan de forma segura desde varias goroutines.
var created, updated int64
// errgroup lanza goroutines y si alguna devuelve error, cancela el resto.
g, gctx := errgroup.WithContext(ctx)
g.SetLimit(s.concurrency) // máximo defaultConcurrency corriendo a la vez
for _, row := range rows {
row := row // copia local: necesaria para que cada goroutine vea SU fila
g.Go(func() error {
// Estampamos updated_at (Postgres no lo provee; lo pone este servicio).
row.Doc["updated_at"] = runAt
isNew, err := s.mongo.Upsert(gctx, row.ID, row.Doc)
if err != nil {
return fmt.Errorf("upsert row %d: %w", row.ID, err)
}
if isNew {
atomic.AddInt64(&created, 1)
} else {
atomic.AddInt64(&updated, 1)
}
return nil
})
}
// g.Wait espera a que terminen todas y devuelve el primer error (si hubo).
if err := g.Wait(); err != nil {
result.Err = err
s.lastResult = result
return result.Err
}
result.Created = int(created)
result.Updated = int(updated)
result.RowsSynced = len(rows)
// --- Paso 3: archivar los que desaparecieron del origen ---
deleted, err := s.reconcileDeleted(ctx, pgIDs)
if err != nil {
result.Err = fmt.Errorf("reconcile deleted: %w", err)
s.lastResult = result
return result.Err
}
result.Deleted = deleted
// --- Paso 4: avanzar el estado persistido (solo si TODO salió bien) ---
if err := s.state.Set(ctx, runAt); err != nil {
result.Err = fmt.Errorf("advance sync state: %w", err)
s.lastResult = result
return result.Err
}
s.lastResult = result
return nil
}
// reconcileDeleted encuentra los ids que están en Mongo pero YA NO en Postgres
// (registros borrados en el origen) y los archiva en la colección de borrados.
func (s *Service) reconcileDeleted(ctx context.Context, pgIDs map[int]struct{}) (int, error) {
// Si este pipeline no tiene con qué listar o mover, no hay nada que hacer.
if s.lister == nil || s.mover == nil {
return 0, nil
}
existingIDs, err := s.lister.ListIDs(ctx)
if err != nil {
return 0, fmt.Errorf("list existing ids: %w", err)
}
// toDelete = ids que están en Mongo pero no en el set de Postgres.
var toDelete []int
for _, id := range existingIDs {
if _, ok := pgIDs[id]; !ok {
toDelete = append(toDelete, id)
}
}
if len(toDelete) == 0 {
return 0, nil
}
// Archivamos en paralelo, igual que los upserts.
g, gctx := errgroup.WithContext(ctx)
g.SetLimit(s.concurrency)
for _, id := range toDelete {
id := id
g.Go(func() error {
if err := s.mover.MoveToDeleted(gctx, id); err != nil {
return fmt.Errorf("move student %d to deleted: %w", id, err)
}
return nil
})
}
if err := g.Wait(); err != nil {
return 0, err
}
return len(toDelete), nil
}
// SeedLastResult inicializa lastResult con un timestamp persistido (ej. justo
// después de arrancar el proceso), para que /sync/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) {
s.mu.Lock()
defer s.mu.Unlock()
s.lastResult = SyncResult{RanAt: t}
}
// LastResult devuelve una copia del último resultado, de forma segura para
// concurrencia (protegida por el mutex). La lee el endpoint /sync/status.
func (s *Service) LastResult() SyncResult {
s.mu.Lock()
defer s.mu.Unlock()
return s.lastResult
}
// 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
import (
"context"
"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.)
type Row interface {
Scan(dest ...interface{}) error // copia los valores de la fila a las variables dadas
}
// 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/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
}
// ---------------------------------------------------------------------------
// 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).
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.
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).
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`
// Set guarda un nuevo timestamp de última sincronización exitosa.
func (s *PGStateStore) Set(ctx context.Context, t time.Time) error {
return s.db.Exec(ctx, upsertLastSyncedSQL, s.stateID, t)
}
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