feat: endpoint de upsert puntual por id

Agrega POST /api/v1/sync/<módulo>/upsert (hoy solo users) para reflejar en
Mongo una fila concreta de Postgres sin esperar al ciclo programado, con su
caso de uso y handler. Elimina los comentarios explicativos del código.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
parent 5bfa78f5
......@@ -2,6 +2,10 @@
This file provides guidance to Claude Code (claude.ai/code) when working with code in this repository.
## Code comments
No comments in code unless explicitly requested by the user. Identifiers should be self-explanatory; don't explain what code does.
## Documentation language
**Toda la documentación del proyecto debe estar en español**: comentarios de código, README, docstrings y cualquier texto explicativo nuevo o modificado. Los identificadores de código (nombres de funciones, variables, tipos) siguen en inglés; solo la prosa explicativa va en español.
......@@ -46,6 +50,7 @@ All endpoints prefixed with `/api/v1`.
- `POST /api/v1/sync/payment_plans/trigger` / `GET /api/v1/sync/payment_plans/status` — payment_plans pipeline (same shape).
- `POST /api/v1/sync/parents/trigger` / `GET /api/v1/sync/parents/status` — parents pipeline (same shape).
- `POST /api/v1/sync/users/trigger` / `GET /api/v1/sync/users/status` — users pipeline (same shape).
- `POST /api/v1/sync/users/upsert` — upsert puntual de un solo usuario, sin esperar al ciclo programado. Body `{"id": <user_id>}`; relee esa fila de `intranet.users` en Postgres y la refleja en la colección `users` de Mongo. Responde `{"status": true, "module": "users", "id": N, "outcome": "created|updated|unchanged"}`, `404` si el id no existe en el origen. Lo consume Trismegisto al crear usuario o restablecer contraseña desde `padresIntranet.jsp`. Solo se registra para los módulos con lector por id en `pipeline.rowReaderFor` (hoy `users`); agregar un `case` ahí más un `FetchByID` en el reader habilita el endpoint para otro módulo.
## Postgres source contract
......@@ -79,6 +84,19 @@ Each cycle has two phases:
Every entity implements `SetRowHash(h string)` (alongside `SetUpdatedAt`) to get the skip-unchanged optimization — a `RowHash string `bson:"row_hash,omitempty" json:"-"`` field, stamped by the service, never read back from Postgres.
### Orden determinista de arrays anidados (obligatorio)
El `row_hash` se calcula sobre el JSON del documento, así que **cualquier array anidado debe salir siempre en el mismo orden** o el ciclo lo detecta como cambio falso y reescribe el documento en cada corrida. Regla: **ordenar todo array anidado por su id principal** (el id de negocio del elemento: `student_id`, `payment_plan_id`, etc.).
- En SQL: `ORDER BY` explícito en el `SELECT` externo (no dentro de un CTE ni de una subconsulta — Postgres no garantiza que se preserve). Para agregaciones usar `JSONB_AGG(DISTINCT ... ORDER BY ...)` — `JSON_AGG` no sirve porque el tipo `json` no tiene operadores `=`/orden; `DISTINCT` además elimina duplicados que generan los joins.
- En Go: `sort.Slice` por el id principal en el reader antes de armar el documento, como defensa adicional.
Implementado hoy en:
- `parent_query.go` — `JSONB_AGG(DISTINCT jsonb_build_object('student_id', ...) ORDER BY ...)`, más `sort.Slice` por `StudentID` en `parent_reader.go`.
- `payment_plan_query.go` — `ORDER BY payment_plan_order, payment_plan_id` en el `SELECT` externo (orden funcional por número de orden, con el id como desempate determinista); los años de `payment_plans` se ordenan descendente en `buildPaymentPlanYears`.
Cualquier módulo nuevo con arrays anidados debe cumplir esta regla antes de entrar a `config.Pipelines`.
## Estructura del proyecto (arquitectura limpia)
Misma arquitectura por capas que `intranet-so-be`: `domain` (entidades + puertos), `usecase` (lógica de aplicación), `infrastructure` (Postgres, Mongo, HTTP, config), `pkg` (helpers reutilizables), `cmd/server` (raíz de composición).
......
// Command server es el punto de entrada y raíz de composición del servicio.
package main
import (
......@@ -56,6 +55,7 @@ func main() {
r := router.New(router.Config{
SyncUseCases: useCases,
UpsertUseCases: pipeline.BuildUpserters(cfg, pool, mongoClient),
PGPing: func(pingCtx context.Context) error { return pool.Ping(pingCtx) },
MongoPing: func(pingCtx context.Context) error { return mongoClient.Ping(pingCtx, nil) },
})
......@@ -70,7 +70,6 @@ func main() {
<-ctx.Done()
log.Println("shutdown signal received, shutting down")
// 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 {
......
......@@ -2,12 +2,25 @@ package entities
import "time"
// StudentRef es un estudiante asociado a un padre/apoderado, embebido dentro de Parent.
type StudentRef struct {
StudentID int `bson:"student_id" json:"student_id"`
StudentCode *string `bson:"student_code" json:"student_code"`
StudentInternalCode *string `bson:"student_internal_code" json:"student_internal_code"`
PaternalLastName *string `bson:"student_paternal_last_name" json:"student_paternal_last_name"`
MaternalLastName *string `bson:"student_maternal_last_name" json:"student_maternal_last_name"`
Name *string `bson:"student_name" json:"student_name"`
DNI *string `bson:"student_dni" json:"student_dni"`
Birthday *time.Time `bson:"student_birthday" json:"student_birthday"`
Email *string `bson:"student_email" json:"student_email"`
Genre *string `bson:"student_genre" json:"student_genre"`
Branch *string `bson:"student_branch" json:"student_branch"`
Level *string `bson:"student_level" json:"student_level"`
Grade *string `bson:"student_grade" json:"student_grade"`
ClassroomAttendanceID *int `bson:"student_classroom_attendance_id" json:"student_classroom_attendance_id"`
PeriodID int `bson:"student_period_id" json:"student_period_id"`
CicleID int `bson:"student_cicle_id" json:"student_cicle_id"`
}
// Parent es la entidad de la colección "parents".
type Parent struct {
ParentID int `bson:"parent_id"`
DNI *string `bson:"parent_dni"`
......@@ -16,17 +29,21 @@ type Parent struct {
Name *string `bson:"parent_name"`
Email *string `bson:"parent_email"`
Phone *string `bson:"parent_phone"`
Gender *int `bson:"parent_gender"`
Students []StudentRef `bson:"students"`
CreatedAt time.Time `bson:"created_at"`
UpdatedAt time.Time `bson:"updated_at"`
RowHash string `bson:"row_hash,omitempty" json:"-"`
}
// SetUpdatedAt estampa la hora de sincronización.
func (r *Parent) SetUpdatedAt(t time.Time) {
r.UpdatedAt = t
}
// SetRowHash guarda el hash de contenido usado para saltar upserts sin cambios.
func (r *Parent) SetRowHash(h string) {
r.RowHash = h
}
func (r *Parent) SetCreatedAt(t time.Time) {
r.CreatedAt = t
}
......@@ -2,23 +2,25 @@ package entities
import "time"
// Professor es la entidad de la colección "professors".
type Professor struct {
ProfessorID int `bson:"professor_id"`
ProfessorPaternalLastName *string `bson:"professor_paternal_last_name"`
ProfessorMaternalLastName *string `bson:"professor_maternal_last_name"`
ProfessorName *string `bson:"professor_name"`
ProfessorDNI *string `bson:"professor_dni"`
CreatedAt time.Time `bson:"created_at"`
UpdatedAt time.Time `bson:"updated_at"`
RowHash string `bson:"row_hash,omitempty" json:"-"`
}
// SetUpdatedAt estampa la hora de sincronización.
func (r *Professor) SetUpdatedAt(t time.Time) {
r.UpdatedAt = t
}
// SetRowHash guarda el hash de contenido usado para saltar upserts sin cambios.
func (r *Professor) SetRowHash(h string) {
r.RowHash = h
}
func (r *Professor) SetCreatedAt(t time.Time) {
r.CreatedAt = t
}
// Package entities contiene las entidades de dominio del servicio.
package entities
import "time"
// Student es la entidad de la colección "students".
type Student struct {
StudentID int `bson:"student_id"`
StudentCode *string `bson:"student_code"`
......@@ -14,23 +12,24 @@ type Student struct {
DNI *string `bson:"student_dni"`
Birthday *time.Time `bson:"student_birthday"`
Email *string `bson:"student_email"`
Genre *string `bson:"genre"`
Branch *string `bson:"branch"`
Level *string `bson:"level"`
Grade *string `bson:"grade"`
ClassroomAttendanceID *int `bson:"classroom_attendance_id"`
Genre *string `bson:"student_genre"`
Branch *string `bson:"student_branch"`
Level *string `bson:"student_level"`
Grade *string `bson:"student_grade"`
ClassroomAttendanceID *int `bson:"student_classroom_attendance_id"`
PeriodID int `bson:"student_period_id"`
CicleID int `bson:"student_cicle_id"`
PaymentPlans []PaymentPlanYear `bson:"payment_plans"`
CreatedAt time.Time `bson:"created_at"`
UpdatedAt time.Time `bson:"updated_at"`
RowHash string `bson:"row_hash,omitempty" json:"-"`
}
// PaymentPlanYear agrupa los planes de pago de un estudiante por año académico.
type PaymentPlanYear struct {
Year int `bson:"year"`
Payments []Payment `bson:"payments"`
}
// Payment es un plan de pago individual embebido dentro de un Student.
type Payment struct {
PaymentPlanID int `bson:"payment_plan_id"`
EnrollmentID int `bson:"enrollment_id"`
......@@ -43,12 +42,14 @@ type Payment struct {
Debt *int `bson:"payment_plan_debt"`
}
// SetUpdatedAt estampa la hora de sincronización.
func (s *Student) SetUpdatedAt(t time.Time) {
s.UpdatedAt = t
}
// SetRowHash guarda el hash de contenido usado para saltar upserts sin cambios.
func (s *Student) SetRowHash(h string) {
s.RowHash = h
}
func (s *Student) SetCreatedAt(t time.Time) {
s.CreatedAt = t
}
......@@ -2,13 +2,11 @@ package entities
import "time"
// SyncRow es una fila lista para volcar a Mongo: su id de negocio y su documento.
type SyncRow struct {
ID int
Doc interface{}
}
// SyncResult resume cómo salió el último ciclo de un pipeline.
type SyncResult struct {
RanAt time.Time
RowsSynced int
......@@ -19,12 +17,14 @@ type SyncResult struct {
Err error
}
// UpdatedAtSetter lo implementan las entidades que quieren que el caso de uso les estampe updated_at.
type UpdatedAtSetter interface {
SetUpdatedAt(t time.Time)
}
// RowHashSetter lo implementan las entidades que quieren que el caso de uso les estampe row_hash.
type RowHashSetter interface {
SetRowHash(h string)
}
type CreatedAtSetter interface {
SetCreatedAt(t time.Time)
}
......@@ -2,7 +2,6 @@ package entities
import "time"
// User es la entidad de la colección "users".
type User struct {
UserID int `bson:"user_id"`
Login *string `bson:"user_login"`
......@@ -10,16 +9,19 @@ type User struct {
CreationDate *time.Time `bson:"user_creation_date"`
ParentID *int `bson:"parent_id"`
Status *int `bson:"user_status"`
CreatedAt time.Time `bson:"created_at"`
UpdatedAt time.Time `bson:"updated_at"`
RowHash string `bson:"row_hash,omitempty" json:"-"`
}
// SetUpdatedAt estampa la hora de sincronización.
func (r *User) SetUpdatedAt(t time.Time) {
r.UpdatedAt = t
}
// SetRowHash guarda el hash de contenido usado para saltar upserts sin cambios.
func (r *User) SetRowHash(h string) {
r.RowHash = h
}
func (r *User) SetCreatedAt(t time.Time) {
r.CreatedAt = t
}
......@@ -2,8 +2,6 @@ package repositories
import "context"
// DocumentFetcher trae los documentos actuales de la colección destino por id de negocio.
// Se usa para registrar en el log de cambios el JSON anterior de las filas que se actualizan.
type DocumentFetcher interface {
FindByIDs(ctx context.Context, ids []int) (map[int]map[string]interface{}, error)
}
// Package repositories declara los contratos (puertos) que la capa de infraestructura implementa.
package repositories
import (
......@@ -8,29 +7,26 @@ import (
"intranet-synchronizer/internal/domain/entities"
)
// SourceReader lee todas las filas del origen (Postgres) en un ciclo.
type SourceReader interface {
FetchAll(ctx context.Context) ([]entities.SyncRow, error)
}
// DocumentUpserter inserta o actualiza un lote de documentos en una sola operación bulk,
// devolviendo cuántos fueron inserciones nuevas (created) y cuántos actualizaciones (updated).
type SourceRowReader interface {
FetchByID(ctx context.Context, id int) (entities.SyncRow, bool, error)
}
type DocumentUpserter interface {
BulkUpsert(ctx context.Context, rows []entities.SyncRow) (created, updated int, err error)
}
// DocumentLister lista los ids actuales en la colección destino junto a su row_hash guardado.
// Las claves del mapa detectan borrados; los valores permiten saltar upserts sin cambios.
type DocumentLister interface {
ListIDsAndHashes(ctx context.Context) (map[int]string, error)
}
// DeletedArchiver archiva un registro que ya no está en el origen a la colección de borrados.
type DeletedArchiver interface {
MoveToDeleted(ctx context.Context, id int) error
}
// SyncStateRepository guarda y recupera cuándo fue la última sincronización exitosa.
type SyncStateRepository interface {
Get(ctx context.Context) (time.Time, error)
Set(ctx context.Context, t time.Time) error
......
// Package config carga y valida la configuración del servicio.
package config
import (
......@@ -11,7 +10,6 @@ import (
"intranet-synchronizer/internal/usecase"
)
// Config guarda los valores de configuración ya leídos y validados.
type Config struct {
PGDSN string
MongoURI string
......@@ -20,7 +18,6 @@ type Config struct {
Port string
}
// Load lee las variables de entorno, valida las obligatorias y devuelve la configuración.
func Load() (Config, error) {
cfg := Config{
PGDSN: os.Getenv("PG_DSN"),
......@@ -57,7 +54,6 @@ func Load() (Config, 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) (usecase.Schedule, error) {
prefix := "SYNC_" + strings.ToUpper(pipelineName)
......@@ -94,7 +90,6 @@ func ScheduleFor(pipelineName string, defaultInterval time.Duration) (usecase.Sc
}
}
// 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 {
......
......@@ -2,7 +2,6 @@ package config
import "intranet-synchronizer/internal/infrastructure/database/postgres"
// PipelineDef describe una tubería de sincronización: origen Postgres -> destino Mongo.
type PipelineDef struct {
Name string
PGQuery string
......@@ -12,7 +11,6 @@ type PipelineDef struct {
StateID int
}
// Pipelines es la lista de todas las tuberías activas.
var Pipelines = []PipelineDef{
{
Name: "students",
......
// Package mongo contiene el cliente concreto de MongoDB.
package mongo
import (
......@@ -8,7 +7,6 @@ import (
"go.mongodb.org/mongo-driver/mongo/options"
)
// NewClient conecta a MongoDB y verifica con un Ping.
func NewClient(ctx context.Context, uri string) (*mongo.Client, error) {
client, err := mongo.Connect(ctx, options.Client().ApplyURI(uri))
if err != nil {
......
......@@ -13,7 +13,6 @@ import (
"intranet-synchronizer/internal/domain/repositories"
)
// DeletedRepository archiva un registro desaparecido del origen a su colección de borrados, estampando deleted_at.
type DeletedRepository struct {
source *mongo.Collection
deleted *mongo.Collection
......@@ -24,7 +23,6 @@ func NewDeletedRepository(source, deleted *mongo.Collection, idField string) rep
return &DeletedRepository{source: source, deleted: deleted, idField: idField}
}
// MoveToDeleted copia el documento a la colección de borrados y lo elimina de la viva (idempotente).
func (m *DeletedRepository) MoveToDeleted(ctx context.Context, id int) error {
var doc bson.M
filter := bson.M{m.idField: id}
......
// Package repository implementa, sobre MongoDB, los puertos de escritura del dominio.
package repository
import (
"context"
"fmt"
"go.mongodb.org/mongo-driver/bson"
"go.mongodb.org/mongo-driver/mongo"
......@@ -11,7 +11,6 @@ import (
"intranet-synchronizer/internal/domain/entities"
)
// DocumentRepository escribe en una colección de Mongo, indexando por idField (no por _id).
type DocumentRepository struct {
coll *mongo.Collection
idField string
......@@ -21,8 +20,6 @@ func NewDocumentRepository(coll *mongo.Collection, idField string) *DocumentRepo
return &DocumentRepository{coll: coll, idField: idField}
}
// 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 *DocumentRepository) BulkUpsert(ctx context.Context, rows []entities.SyncRow) (created, updated int, err error) {
if len(rows) == 0 {
return 0, 0, nil
......@@ -31,7 +28,24 @@ func (c *DocumentRepository) BulkUpsert(ctx context.Context, rows []entities.Syn
models := make([]mongo.WriteModel, len(rows))
for i, row := range rows {
filter := bson.M{c.idField: row.ID}
update := bson.M{"$set": row.Doc}
// created_at va aparte por $setOnInsert: así solo se estampa en el alta y un update
// nunca pisa el created_at real ya guardado en Mongo.
fields, err := bson.Marshal(row.Doc)
if err != nil {
return 0, 0, fmt.Errorf("marshal doc %d: %w", row.ID, err)
}
var set bson.M
if err := bson.Unmarshal(fields, &set); err != nil {
return 0, 0, fmt.Errorf("unmarshal doc %d: %w", row.ID, err)
}
createdAt, hasCreatedAt := set["created_at"]
delete(set, "created_at")
update := bson.M{"$set": set}
if hasCreatedAt {
update["$setOnInsert"] = bson.M{"created_at": createdAt}
}
models[i] = mongo.NewUpdateOneModel().SetFilter(filter).SetUpdate(update).SetUpsert(true)
}
......@@ -44,9 +58,6 @@ func (c *DocumentRepository) BulkUpsert(ctx context.Context, rows []entities.Syn
return created, updated, nil
}
// 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 *DocumentRepository) 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 {
......@@ -75,7 +86,6 @@ func (c *DocumentRepository) ListIDsAndHashes(ctx context.Context) (map[int]stri
return hashes, cur.Err()
}
// FindByIDs devuelve los documentos actuales de los ids dados, indexados por id de negocio.
func (c *DocumentRepository) FindByIDs(ctx context.Context, ids []int) (map[int]map[string]interface{}, error) {
docs := make(map[int]map[string]interface{}, len(ids))
if len(ids) == 0 {
......
// Package postgres contiene el acceso concreto a PostgreSQL: pool de conexiones,
// readers del origen y persistencia del estado de sincronización.
package postgres
import (
......@@ -9,12 +7,10 @@ import (
"github.com/jackc/pgx/v5/pgxpool"
)
// Pool envuelve el pool de pgx para poder colgarle métodos propios.
type Pool struct {
*pgxpool.Pool
}
// NewPool abre el pool y verifica la conexión con un Ping.
func NewPool(ctx context.Context, dsn string) (*Pool, error) {
pool, err := pgxpool.New(ctx, dsn)
if err != nil {
......@@ -26,7 +22,6 @@ func NewPool(ctx context.Context, dsn string) (*Pool, error) {
return &Pool{pool}, nil
}
// 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,
......@@ -38,7 +33,6 @@ VALUES (1, '1970-01-01T00:00:00Z'), (3, '1970-01-01T00:00:00Z'), (4, '1970-01-01
ON CONFLICT (id) DO NOTHING;
`
// BootstrapSchema prepara la tabla sync_state al inicio del servicio.
func (p *Pool) BootstrapSchema(ctx context.Context) error {
if _, err := p.Pool.Exec(ctx, bootstrapSyncStateSQL); err != nil {
return fmt.Errorf("bootstrap sync_state schema: %w", err)
......
package postgres
// ParentQuery es el SQL crudo de la tubería de padres/apoderados, con estudiantes a cargo embebidos.
const ParentQuery = `
SELECT pp.persona_id as parent_id,
pp.persona_numero_documento_identidad as parent_dni,
......@@ -9,11 +8,16 @@ SELECT pp.persona_id as parent_id,
public.to_camel_case(pp.persona_nombre) as parent_name,
pp.persona_correo as parent_email,
pp.persona_telefono as parent_phone,
pp.sexo_id as parent_gender,
est.students
FROM persona.pe_persona pp
INNER JOIN (SELECT apoderado.persona_id,
JSON_AGG(json_build_object('student_id', apoderado.estudiante_id)) as students
JSONB_AGG(DISTINCT jsonb_build_object('student_id', apoderado.estudiante_id)
ORDER BY jsonb_build_object('student_id', apoderado.estudiante_id)) as students
FROM matricula.ma_estudiante_apoderado apoderado
WHERE apoderado.estudiante_apoderado_estado = '1'
INNER JOIN matricula.ma_matricula matricula on apoderado.estudiante_apoderado_id = matricula.estudiante_apoderado_id
INNER JOIN academico.ac_apertura apertura on matricula.apertura_id = apertura.apertura_id
WHERE apoderado.estudiante_apoderado_estado = '1' AND NOT matricula.matricula_anulada AND NOT matricula.matricula_retirado
AND apertura.periodo_academico_id = 14 AND apertura.ciclo_id = 28
GROUP BY apoderado.persona_id) est ON est.persona_id = pp.persona_id;
`
......@@ -3,32 +3,118 @@ package postgres
import (
"context"
"fmt"
"sort"
"golang.org/x/sync/errgroup"
"intranet-synchronizer/internal/domain/entities"
"intranet-synchronizer/internal/domain/repositories"
"intranet-synchronizer/pkg/jsonutil"
)
// ParentReader ejecuta Query directamente y arma una SyncRow por fila del resultado.
type ParentReader struct {
pool *Pool
query string
}
// NewParentReader construye el reader de parents sobre el pool dado.
func NewParentReader(pool *Pool, query string) repositories.SourceReader {
return &ParentReader{pool: pool, query: query}
}
// FetchAll corre la consulta y lee fila por fila.
type parentRow struct {
parentID int
dni *string
paternalLastName *string
maternalLastName *string
name *string
email *string
phone *string
gender *int
studentIDs []entities.StudentRef
}
func (f *ParentReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error) {
var (
parentRows []parentRow
students map[int]*entities.Student
)
g, gctx := errgroup.WithContext(ctx)
g.Go(func() error {
rows, err := f.fetchParents(gctx)
if err != nil {
return err
}
parentRows = rows
return nil
})
g.Go(func() error {
byID, _, err := FetchStudentBasics(gctx, f.pool, StudentQuery)
if err != nil {
return err
}
students = byID
return nil
})
if err := g.Wait(); err != nil {
return nil, err
}
result := make([]entities.SyncRow, 0, len(parentRows))
for _, p := range parentRows {
studentRefs := make([]entities.StudentRef, 0, len(p.studentIDs))
for _, ref := range p.studentIDs {
s, ok := students[ref.StudentID]
if !ok {
continue
}
studentRefs = append(studentRefs, entities.StudentRef{
StudentID: s.StudentID,
StudentCode: s.StudentCode,
StudentInternalCode: s.StudentInternalCode,
PaternalLastName: s.PaternalLastName,
MaternalLastName: s.MaternalLastName,
Name: s.Name,
DNI: s.DNI,
Birthday: s.Birthday,
Email: s.Email,
Genre: s.Genre,
Branch: s.Branch,
Level: s.Level,
Grade: s.Grade,
ClassroomAttendanceID: s.ClassroomAttendanceID,
PeriodID: s.PeriodID,
CicleID: s.CicleID,
})
}
sort.Slice(studentRefs, func(i, j int) bool {
return studentRefs[i].StudentID < studentRefs[j].StudentID
})
doc := &entities.Parent{
ParentID: p.parentID,
DNI: p.dni,
PaternalLastName: p.paternalLastName,
MaternalLastName: p.maternalLastName,
Name: p.name,
Email: p.email,
Phone: p.phone,
Gender: p.gender,
Students: studentRefs,
}
result = append(result, entities.SyncRow{ID: p.parentID, Doc: doc})
}
return result, nil
}
func (f *ParentReader) fetchParents(ctx context.Context) ([]parentRow, error) {
rows, err := f.pool.Pool.Query(ctx, f.query)
if err != nil {
return nil, fmt.Errorf("query parents: %w", err)
}
defer rows.Close()
var result []entities.SyncRow
var result []parentRow
for rows.Next() {
var (
parentID int
......@@ -38,28 +124,29 @@ func (f *ParentReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error)
name *string
email *string
phone *string
students []byte
gender *int
studentsRaw []byte
)
if err := rows.Scan(&parentID, &dni, &paternalLastName, &maternalLastName, &name, &email, &phone, &students); err != nil {
if err := rows.Scan(&parentID, &dni, &paternalLastName, &maternalLastName, &name, &email, &phone, &gender, &studentsRaw); err != nil {
return nil, fmt.Errorf("scan parent row: %w", err)
}
studentRefs, err := jsonutil.DecodeList[entities.StudentRef](students)
studentIDs, err := jsonutil.DecodeList[entities.StudentRef](studentsRaw)
if err != nil {
return nil, fmt.Errorf("decode students for parent %d: %w", parentID, err)
}
doc := &entities.Parent{
ParentID: parentID,
DNI: dni,
PaternalLastName: paternalLastName,
MaternalLastName: maternalLastName,
Name: name,
Email: email,
Phone: phone,
Students: studentRefs,
}
result = append(result, entities.SyncRow{ID: parentID, Doc: doc})
result = append(result, parentRow{
parentID: parentID,
dni: dni,
paternalLastName: paternalLastName,
maternalLastName: maternalLastName,
name: name,
email: email,
phone: phone,
gender: gender,
studentIDs: studentIDs,
})
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate parent rows: %w", err)
......
package postgres
// PaymentPlanQuery es el SQL crudo de la tubería de planes de pago: junta caja.ca_plan_de_pago con matricula.ma_matricula.
const PaymentPlanQuery = `
WITH tb_lista_plan_pago AS (select c.plan_de_pago_id as payment_plan_id,
ma.estudiante_id as student_id,
......@@ -22,7 +21,8 @@ WITH tb_lista_plan_pago AS (select c.plan_de_pago_id
and ma.matricula_anulada = false
and ac.ciclo_id = 28
and ac.periodo_academico_id = 14
ORDER BY to_char(plan_de_pago_fecha_ven, 'YYYY-mm-dd'), plan_de_pago_id)
)
SELECT *
FROM tb_lista_plan_pago;
FROM tb_lista_plan_pago
ORDER BY payment_plan_order, payment_plan_id;
`
package postgres
// ProfessorQuery es el SQL crudo de la tubería de profesores.
const ProfessorQuery = `
SELECT profesor_id as professor_id,
public.to_camel_case(p.persona_apellido_paterno) as professor_paternal_last_name,
......
......@@ -8,18 +8,15 @@ import (
"intranet-synchronizer/internal/domain/repositories"
)
// ProfessorReader ejecuta Query directamente y arma una SyncRow por fila del resultado.
type ProfessorReader struct {
pool *Pool
query string
}
// NewProfessorReader construye el reader de professors sobre el pool dado.
func NewProfessorReader(pool *Pool, query string) repositories.SourceReader {
return &ProfessorReader{pool: pool, query: query}
}
// FetchAll corre la consulta y lee fila por fila.
func (f *ProfessorReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error) {
rows, err := f.pool.Pool.Query(ctx, f.query)
if err != nil {
......
......@@ -14,11 +14,13 @@ SELECT e.estudiante_id as student_id,
pp.persona_numero_documento_identidad as student_dni,
pp.persona_fecha_nacimiento as student_birthday,
pp.persona_correo as student_email,
sexo.catalogo_siiaa_nombre as genre,
public.to_camel_case(se.sede_nombre) as branch,
public.to_camel_case(nivel.catalogo_siiaa_nombre) as level,
public.to_camel_case(grado.grado_sigla) as grade,
ma.aula_id_asiste as classroom_attendance_id
sexo.catalogo_siiaa_nombre as student_genre,
public.to_camel_case(se.sede_nombre) as student_branch,
public.to_camel_case(nivel.catalogo_siiaa_nombre) as student_level,
public.to_camel_case(grado.grado_sigla) as student_grade,
ma.aula_id_asiste as student_classroom_attendance_id,
14 as student_period_id,
28 as student_cicle_id
FROM matricula.ma_estudiante e
INNER JOIN persona.pe_persona pp ON e.persona_id = pp.persona_id
INNER JOIN general.gn_catalogo_siiaa sexo on pp.sexo_id = sexo.catalogo_siiaa_id
......@@ -27,8 +29,7 @@ FROM matricula.ma_estudiante e
INNER JOIN administracion.ad_sede se on ac.sede_id = se.sede_id
INNER JOIN general.gn_catalogo_siiaa nivel on ac.nivel_id = nivel.catalogo_siiaa_id
INNER JOIN academico.ac_grado grado on ac.grado_id = grado.grado_id
WHERE ac.periodo_academico_id >= 14
and ac.periodo_academico_id = 14
WHERE ac.periodo_academico_id = 14
and ac.ciclo_id = 28
and ma.matricula_retirado = false
and ma.matricula_anulada = false;
......
......@@ -10,21 +10,17 @@ import (
"intranet-synchronizer/internal/domain/repositories"
)
// StudentReader ejecuta Query + PaymentPlanQuery y arma una SyncRow por estudiante,
// anidando sus planes de pago agrupados por año.
type StudentReader struct {
pool *Pool
query string
}
// NewStudentReader construye el reader de estudiantes sobre el pool dado.
func NewStudentReader(pool *Pool, query string) repositories.SourceReader {
return &StudentReader{pool: pool, query: query}
}
// FetchAll corre la consulta de estudiantes y la de planes de pago, y anida la segunda dentro de la primera.
func (f *StudentReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error) {
students, order, err := f.fetchStudents(ctx)
students, order, err := FetchStudentBasics(ctx, f.pool, f.query)
if err != nil {
return nil, err
}
......@@ -43,9 +39,8 @@ func (f *StudentReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error
return result, nil
}
// fetchStudents lee la fila base de cada estudiante (sin planes de pago).
func (f *StudentReader) fetchStudents(ctx context.Context) (map[int]*entities.Student, []int, error) {
rows, err := f.pool.Pool.Query(ctx, f.query)
func FetchStudentBasics(ctx context.Context, pool *Pool, query string) (map[int]*entities.Student, []int, error) {
rows, err := pool.Pool.Query(ctx, query)
if err != nil {
return nil, nil, fmt.Errorf("query students: %w", err)
}
......@@ -69,10 +64,12 @@ func (f *StudentReader) fetchStudents(ctx context.Context) (map[int]*entities.St
level *string
grade *string
classroomAttendanceID *int
periodID int
cicleID int
)
if err := rows.Scan(&studentID, &studentCode, &studentInternalCode, &paternalLastName, &maternalLastName,
&name, &dni, &birthday, &email, &genre, &branch, &level,
&grade, &classroomAttendanceID); err != nil {
&grade, &classroomAttendanceID, &periodID, &cicleID); err != nil {
return nil, nil, fmt.Errorf("scan student row: %w", err)
}
......@@ -91,6 +88,8 @@ func (f *StudentReader) fetchStudents(ctx context.Context) (map[int]*entities.St
Level: level,
Grade: grade,
ClassroomAttendanceID: classroomAttendanceID,
PeriodID: periodID,
CicleID: cicleID,
}
order = append(order, studentID)
}
......@@ -100,13 +99,11 @@ func (f *StudentReader) fetchStudents(ctx context.Context) (map[int]*entities.St
return students, order, nil
}
// paymentRow es una fila de PaymentPlanQuery antes de agruparse por año.
type paymentRow struct {
entities.Payment
year int
}
// fetchPaymentPlans lee todos los planes de pago y los agrupa por student_id.
func (f *StudentReader) fetchPaymentPlans(ctx context.Context) (map[int][]paymentRow, error) {
rows, err := f.pool.Pool.Query(ctx, PaymentPlanQuery)
if err != nil {
......@@ -155,7 +152,6 @@ func (f *StudentReader) fetchPaymentPlans(ctx context.Context) (map[int][]paymen
return plansByStudent, nil
}
// buildPaymentPlanYears agrupa los pagos de un estudiante por año, orden descendente.
func buildPaymentPlanYears(payments []paymentRow) []entities.PaymentPlanYear {
byYear := make(map[int][]entities.Payment)
for _, p := range payments {
......
......@@ -7,21 +7,17 @@ import (
"intranet-synchronizer/internal/domain/repositories"
)
// SyncStateRepository persiste el timestamp de última sincronización en la tabla sync_state,
// una fila por pipeline (stateID).
type SyncStateRepository struct {
pool *Pool
stateID int
}
// NewSyncStateRepository construye el repositorio de estado de un pipeline.
func NewSyncStateRepository(pool *Pool, stateID int) repositories.SyncStateRepository {
return &SyncStateRepository{pool: pool, stateID: stateID}
}
const selectLastSyncedSQL = `SELECT last_synced_at FROM sync_state WHERE id = $1`
// Get lee el último timestamp sincronizado para este pipeline.
func (s *SyncStateRepository) Get(ctx context.Context) (time.Time, error) {
var t time.Time
if err := s.pool.Pool.QueryRow(ctx, selectLastSyncedSQL, s.stateID).Scan(&t); err != nil {
......@@ -30,12 +26,10 @@ func (s *SyncStateRepository) Get(ctx context.Context) (time.Time, error) {
return t, nil
}
// 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`
// Set guarda un nuevo timestamp de última sincronización exitosa.
func (s *SyncStateRepository) Set(ctx context.Context, t time.Time) error {
_, err := s.pool.Pool.Exec(ctx, upsertLastSyncedSQL, s.stateID, t)
return err
......
package postgres
// UserQuery es el SQL crudo de la tubería de usuarios.
const UserQuery = `
SELECT user_id, user_login, user_password, user_creation_date, parent_id, user_status
FROM intranet.users;
`
const UserByIDQuery = `
SELECT user_id, user_login, user_password, user_creation_date, parent_id, user_status
FROM intranet.users
WHERE user_id = $1;
`
......@@ -2,34 +2,30 @@ package postgres
import (
"context"
"errors"
"fmt"
"time"
"github.com/jackc/pgx/v5"
"intranet-synchronizer/internal/domain/entities"
"intranet-synchronizer/internal/domain/repositories"
)
// UserReader ejecuta Query directamente y arma una SyncRow por fila del resultado.
type UserReader struct {
pool *Pool
query string
}
// NewUserReader construye el reader de users sobre el pool dado.
func NewUserReader(pool *Pool, query string) repositories.SourceReader {
return &UserReader{pool: pool, query: query}
}
// FetchAll corre la consulta y lee fila por fila.
func (f *UserReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error) {
rows, err := f.pool.Pool.Query(ctx, f.query)
if err != nil {
return nil, fmt.Errorf("query users: %w", err)
}
defer rows.Close()
func NewUserRowReader(pool *Pool) repositories.SourceRowReader {
return &UserReader{pool: pool, query: UserByIDQuery}
}
var result []entities.SyncRow
for rows.Next() {
func scanUserRow(row pgx.CollectableRow) (entities.SyncRow, error) {
var (
userID int
login *string
......@@ -38,8 +34,8 @@ func (f *UserReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error) {
parentID *int
status *int
)
if err := rows.Scan(&userID, &login, &password, &creationDate, &parentID, &status); err != nil {
return nil, fmt.Errorf("scan user row: %w", err)
if err := row.Scan(&userID, &login, &password, &creationDate, &parentID, &status); err != nil {
return entities.SyncRow{}, fmt.Errorf("scan user row: %w", err)
}
doc := &entities.User{
......@@ -50,10 +46,34 @@ func (f *UserReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error) {
ParentID: parentID,
Status: status,
}
result = append(result, entities.SyncRow{ID: userID, Doc: doc})
return entities.SyncRow{ID: userID, Doc: doc}, nil
}
func (f *UserReader) FetchAll(ctx context.Context) ([]entities.SyncRow, error) {
rows, err := f.pool.Pool.Query(ctx, f.query)
if err != nil {
return nil, fmt.Errorf("query users: %w", err)
}
if err := rows.Err(); err != nil {
result, err := pgx.CollectRows(rows, scanUserRow)
if err != nil {
return nil, fmt.Errorf("iterate user rows: %w", err)
}
return result, nil
}
func (f *UserReader) FetchByID(ctx context.Context, id int) (entities.SyncRow, bool, error) {
rows, err := f.pool.Pool.Query(ctx, UserByIDQuery, id)
if err != nil {
return entities.SyncRow{}, false, fmt.Errorf("query user %d: %w", id, err)
}
row, err := pgx.CollectExactlyOneRow(rows, scanUserRow)
if errors.Is(err, pgx.ErrNoRows) {
return entities.SyncRow{}, false, nil
}
if err != nil {
return entities.SyncRow{}, false, fmt.Errorf("read user %d: %w", id, err)
}
return row, true, nil
}
// Package handlers contiene los manejadores HTTP (capa web, Gin).
package handlers
import (
......@@ -8,18 +7,15 @@ import (
"github.com/gin-gonic/gin"
)
// HealthHandler responde el chequeo de salud del servicio.
type HealthHandler struct {
pgPing func(context.Context) error
mongoPing func(context.Context) error
}
// NewHealthHandler construye el handler de salud con los pings de cada motor.
func NewHealthHandler(pgPing, mongoPing func(context.Context) error) *HealthHandler {
return &HealthHandler{pgPing: pgPing, mongoPing: mongoPing}
}
// Health responde GET /health con el estado de Postgres y Mongo.
func (h *HealthHandler) Health(c *gin.Context) {
pgErr := h.pgPing(c.Request.Context())
mongoErr := h.mongoPing(c.Request.Context())
......
......@@ -9,17 +9,14 @@ import (
"intranet-synchronizer/internal/usecase"
)
// SyncHandler expone los endpoints trigger/status de UN pipeline de sincronización.
type SyncHandler struct {
runSync usecase.RunSyncPipelineUseCase
}
// NewSyncHandler construye el handler de un pipeline sobre su caso de uso.
func NewSyncHandler(runSync usecase.RunSyncPipelineUseCase) *SyncHandler {
return &SyncHandler{runSync: runSync}
}
// syncResultJSON arma el cuerpo JSON común de /trigger y /status.
func syncResultJSON(err error, result entities.SyncResult) gin.H {
resp := gin.H{
"ran_at": result.RanAt,
......@@ -35,13 +32,11 @@ func syncResultJSON(err error, result entities.SyncResult) gin.H {
return resp
}
// Trigger atiende POST /sync/<pipeline>/trigger: corre una sincronización ahora.
func (h *SyncHandler) Trigger(c *gin.Context) {
err := h.runSync.Run(c.Request.Context())
c.JSON(http.StatusOK, syncResultJSON(err, h.runSync.LastResult()))
}
// Status atiende GET /sync/<pipeline>/status: devuelve el resultado del último ciclo.
func (h *SyncHandler) Status(c *gin.Context) {
result := h.runSync.LastResult()
resp := syncResultJSON(nil, result)
......
package handlers
import (
"bytes"
"io"
"log"
"net/http"
"github.com/gin-gonic/gin"
"intranet-synchronizer/internal/usecase"
)
type UpsertHandler struct {
name string
upsert usecase.UpsertDocumentUseCase
}
func NewUpsertHandler(name string, upsert usecase.UpsertDocumentUseCase) *UpsertHandler {
return &UpsertHandler{name: name, upsert: upsert}
}
type upsertRequest struct {
ID int `json:"id" binding:"required"`
}
func (h *UpsertHandler) Upsert(c *gin.Context) {
body, _ := io.ReadAll(c.Request.Body)
log.Printf("[%s upsert] payload=%s", h.name, string(body))
c.Request.Body = io.NopCloser(bytes.NewReader(body))
var req upsertRequest
if err := c.ShouldBindJSON(&req); err != nil {
log.Printf("[%s upsert] payload inválido: %v", h.name, err)
c.JSON(http.StatusBadRequest, gin.H{"status": false, "message": "id es obligatorio"})
return
}
outcome, err := h.upsert.Upsert(c.Request.Context(), req.ID)
if err != nil {
log.Printf("[%s upsert] id=%d error: %v", h.name, req.ID, err)
c.JSON(http.StatusInternalServerError, gin.H{"status": false, "message": err.Error()})
return
}
log.Printf("[%s upsert] id=%d outcome=%s", h.name, req.ID, outcome)
if outcome == usecase.UpsertNotFound {
c.JSON(http.StatusNotFound, gin.H{
"status": false,
"message": "no existe el registro en el origen",
"module": h.name,
"id": req.ID,
})
return
}
c.JSON(http.StatusOK, gin.H{
"status": true,
"module": h.name,
"id": req.ID,
"outcome": outcome,
})
}
// Package router arma el router HTTP y conecta handlers con casos de uso.
package router
import (
......@@ -10,15 +9,13 @@ import (
"intranet-synchronizer/internal/usecase"
)
// Config son las dependencias ya construidas que el router necesita.
type Config struct {
// SyncUseCases mapea el nombre del pipeline a su caso de uso (ej. "students").
SyncUseCases map[string]usecase.RunSyncPipelineUseCase
UpsertUseCases map[string]usecase.UpsertDocumentUseCase
PGPing func(context.Context) error
MongoPing func(context.Context) error
}
// New construye el router de Gin y registra las rutas de /api/v1.
func New(cfg Config) *gin.Engine {
r := gin.Default()
......@@ -27,11 +24,15 @@ func New(cfg Config) *gin.Engine {
v1 := r.Group("/api/v1")
v1.GET("/health", health.Health)
// Un par trigger/status por pipeline declarado.
for name, uc := range cfg.SyncUseCases {
h := handlers.NewSyncHandler(uc)
v1.POST("/sync/"+name+"/trigger", h.Trigger)
v1.GET("/sync/"+name+"/status", h.Status)
}
for name, uc := range cfg.UpsertUseCases {
h := handlers.NewUpsertHandler(name, uc)
v1.POST("/sync/"+name+"/upsert", h.Upsert)
}
return r
}
// Package pipeline arma, por cada pipeline declarado en config.Pipelines, la cadena
// completa reader -> upserter/lister -> archiver -> estado -> caso de uso -> scheduler.
package pipeline
import (
......@@ -17,7 +15,6 @@ import (
"intranet-synchronizer/pkg/changelog"
)
// readerFor elige el reader de Postgres que corresponde al pipeline.
func readerFor(name string, pool *postgres.Pool, query string) (repositories.SourceReader, error) {
switch name {
case "students":
......@@ -33,8 +30,38 @@ func readerFor(name string, pool *postgres.Pool, query string) (repositories.Sou
}
}
// Build arma los casos de uso y schedulers de todos los pipelines declarados.
// Devuelve los casos de uso por nombre (para el router) y los schedulers listos para arrancar.
// rowReaderFor devuelve el lector por id de los módulos que aceptan upsert bajo demanda.
// Los módulos sin lector puntual simplemente no exponen el endpoint /upsert.
func rowReaderFor(name string, pool *postgres.Pool) repositories.SourceRowReader {
switch name {
case "users":
return postgres.NewUserRowReader(pool)
default:
return nil
}
}
// BuildUpserters arma un caso de uso de upsert puntual por cada pipeline que tenga lector por id.
func BuildUpserters(
cfg config.Config,
pool *postgres.Pool,
mongoClient *mongodriver.Client,
) map[string]usecase.UpsertDocumentUseCase {
useCases := make(map[string]usecase.UpsertDocumentUseCase)
for _, p := range config.Pipelines {
reader := rowReaderFor(p.Name, pool)
if reader == nil {
continue
}
docRepo := repository.NewDocumentRepository(mongoClient.Database(cfg.MongoDB).Collection(p.MongoCollection), p.IDField)
useCases[p.Name] = usecase.NewUpsertDocumentUseCase(reader, docRepo, docRepo, time.Now)
}
return useCases
}
func Build(
ctx context.Context,
cfg config.Config,
......@@ -52,13 +79,10 @@ func Build(
return nil, nil, err
}
// docRepo cumple a la vez DocumentUpserter y DocumentLister.
docRepo := repository.NewDocumentRepository(db.Collection(p.MongoCollection), p.IDField)
archiver := repository.NewDeletedRepository(db.Collection(p.MongoCollection), db.Collection(p.MongoDeletedCollection), p.IDField)
state := postgres.NewSyncStateRepository(pool, p.StateID)
// changeLog deja en /opt/logs/sync_<pipeline>.log el JSON de cada alta y el par
// anterior/nuevo de cada actualización.
changeLog, err := changelog.New(changelog.DefaultDir, p.Name)
if err != nil {
return nil, nil, fmt.Errorf("change log for pipeline %s: %w", p.Name, err)
......@@ -66,7 +90,6 @@ func Build(
uc := usecase.NewRunSyncPipelineUseCase(p.Name, reader, docRepo, docRepo, archiver, state, docRepo, changeLog, time.Now)
// Sembrar el último resultado desde el estado persistido antes del primer ciclo.
if lastSynced, err := state.Get(ctx); err != nil {
return nil, nil, fmt.Errorf("load persisted %s sync state: %w", p.Name, err)
} else if !lastSynced.IsZero() {
......
// Package usecase contiene la lógica de aplicación: orquesta los puertos del dominio
// sin depender de librerías de base de datos ni del framework web.
package usecase
import (
......@@ -20,10 +18,8 @@ import (
"intranet-synchronizer/internal/domain/repositories"
)
// defaultConcurrency: cuántas filas se procesan a la vez en paralelo.
const defaultConcurrency = 10
// stampUpdatedAt estampa updated_at en el documento, sea map genérico o entidad tipada.
func stampUpdatedAt(doc interface{}, t time.Time) {
switch v := doc.(type) {
case map[string]interface{}:
......@@ -33,7 +29,6 @@ func stampUpdatedAt(doc interface{}, t time.Time) {
}
}
// stampRowHash guarda el hash de contenido en el documento, sea map genérico o entidad tipada.
func stampRowHash(doc interface{}, h string) {
switch v := doc.(type) {
case map[string]interface{}:
......@@ -43,6 +38,17 @@ func stampRowHash(doc interface{}, h string) {
}
}
// stampCreatedAt estampa created_at en el documento; el repositorio Mongo lo escribe con
// $setOnInsert, así que este valor solo queda en el alta y no pisa el created_at de un update.
func stampCreatedAt(doc interface{}, t time.Time) {
switch v := doc.(type) {
case map[string]interface{}:
v["created_at"] = t
case entities.CreatedAtSetter:
v.SetCreatedAt(t)
}
}
// computeRowHash calcula un hash de contenido del documento, ANTES de estampar updated_at/row_hash,
// para poder comparar contra el hash guardado en Mongo y saltar el upsert si la fila no cambió.
func computeRowHash(doc interface{}) (string, error) {
......@@ -54,21 +60,17 @@ func computeRowHash(doc interface{}) (string, error) {
return hex.EncodeToString(sum[:]), nil
}
// RunSyncPipelineUseCase ejecuta un ciclo completo de sincronización de UN pipeline
// y expone el resultado del último ciclo.
type RunSyncPipelineUseCase interface {
Run(ctx context.Context) error
SeedLastResult(t time.Time)
LastResult() entities.SyncResult
}
// ChangeLogger registra en un archivo el detalle de cada documento insertado o actualizado.
type ChangeLogger interface {
Insert(id int, doc interface{}) error
Update(id int, before, after interface{}) error
}
// runSyncPipeline es la implementación de RunSyncPipelineUseCase.
type runSyncPipeline struct {
name string
fetcher repositories.DocumentFetcher
......@@ -81,12 +83,10 @@ type runSyncPipeline struct {
now func() time.Time
concurrency int
// mu protege lastResult entre goroutines concurrentes.
mu sync.Mutex
lastResult entities.SyncResult
}
// NewRunSyncPipelineUseCase arma el caso de uso inyectándole todos sus puertos.
func NewRunSyncPipelineUseCase(
name string,
reader repositories.SourceReader,
......@@ -112,8 +112,6 @@ func NewRunSyncPipelineUseCase(
}
}
// Run ejecuta un ciclo completo: lee Postgres, upsertea en Mongo, archiva borrados,
// avanza el estado persistido solo si todo salió bien.
func (s *runSyncPipeline) Run(ctx context.Context) error {
s.mu.Lock()
defer s.mu.Unlock()
......@@ -133,8 +131,6 @@ func (s *runSyncPipeline) Run(ctx context.Context) error {
pgIDs[row.ID] = struct{}{}
}
// existingHashes trae, en una sola pasada, el row_hash ya guardado por id (para saltar
// upserts de filas sin cambios) y sirve además como base para detectar borrados abajo.
var existingHashes map[int]string
if s.lister != nil {
var err error
......@@ -146,12 +142,9 @@ func (s *runSyncPipeline) Run(ctx context.Context) error {
}
}
// Fase 1 (CPU, concurrente): calcular el hash de cada fila y descartar las que no cambiaron
// frente a lo guardado en Mongo. No hay I/O acá, solo sirve para no mandar filas de más al bulk.
changed := make([]entities.SyncRow, len(rows))
var changedCount, unchanged int64
// createdIDs/updatedIDs se acumulan para poder listarlos en consola al cerrar el ciclo.
var (
idsMu sync.Mutex
createdIDs []int
......@@ -183,6 +176,7 @@ func (s *runSyncPipeline) Run(ctx context.Context) error {
stampUpdatedAt(row.Doc, runAt)
stampRowHash(row.Doc, hash)
stampCreatedAt(row.Doc, runAt)
idx := atomic.AddInt64(&changedCount, 1) - 1
changed[idx] = row
return nil
......@@ -207,7 +201,6 @@ func (s *runSyncPipeline) Run(ctx context.Context) error {
}
}
// Fase 2 (I/O, un solo round trip): escribir todas las filas que cambiaron en un único BulkWrite.
created, updated, err := s.upserter.BulkUpsert(ctx, changed)
if err != nil {
result.Err = fmt.Errorf("bulk upsert: %w", err)
......@@ -243,7 +236,6 @@ func (s *runSyncPipeline) Run(ctx context.Context) error {
return nil
}
// logChangedIDs imprime en consola los ids creados y actualizados del ciclo.
func (s *runSyncPipeline) logChangedIDs(createdIDs, updatedIDs []int) {
log.Printf("[%s] created=%d updated=%d", s.name, len(createdIDs), len(updatedIDs))
for _, id := range createdIDs {
......@@ -254,8 +246,6 @@ func (s *runSyncPipeline) logChangedIDs(createdIDs, updatedIDs []int) {
}
}
// writeChangeLog vuelca al archivo de cambios el JSON insertado de cada alta y el par
// anterior/nuevo de cada actualización. Un fallo de escritura se registra pero no corta el ciclo.
func (s *runSyncPipeline) writeChangeLog(changed []entities.SyncRow, createdIDs, updatedIDs []int, beforeDocs map[int]map[string]interface{}) {
if s.changeLog == nil {
return
......@@ -278,7 +268,6 @@ func (s *runSyncPipeline) writeChangeLog(changed []entities.SyncRow, createdIDs,
}
}
// reconcileDeleted encuentra ids en Mongo (existingHashes) que ya no están en Postgres y los archiva.
func (s *runSyncPipeline) reconcileDeleted(ctx context.Context, pgIDs map[int]struct{}, existingHashes map[int]string) (int, error) {
if s.lister == nil || s.archiver == nil {
return 0, nil
......@@ -311,14 +300,12 @@ func (s *runSyncPipeline) reconcileDeleted(ctx context.Context, pgIDs map[int]st
return len(toDelete), nil
}
// SeedLastResult inicializa lastResult con un timestamp persistido tras arrancar el proceso.
func (s *runSyncPipeline) SeedLastResult(t time.Time) {
s.mu.Lock()
defer s.mu.Unlock()
s.lastResult = entities.SyncResult{RanAt: t}
}
// LastResult devuelve una copia del último resultado, protegida por el mutex.
func (s *runSyncPipeline) LastResult() entities.SyncResult {
s.mu.Lock()
defer s.mu.Unlock()
......
......@@ -2,13 +2,10 @@ package usecase
import "time"
// Schedule decide cuándo debe correr el próximo ciclo de un pipeline.
type Schedule interface {
// Next devuelve cuánto falta para el próximo disparo desde "now".
Next(now time.Time) time.Duration
}
// IntervalSchedule dispara cada "Interval" de tiempo, sin importar la hora.
type IntervalSchedule struct {
Interval time.Duration
}
......@@ -17,8 +14,6 @@ func (s IntervalSchedule) Next(now time.Time) time.Duration {
return s.Interval
}
// DailySchedule dispara una vez al día a la hora Hour:Minute (hora local).
// Si esa hora ya pasó hoy, el próximo disparo es mañana a esa misma hora.
type DailySchedule struct {
Hour int
Minute int
......
......@@ -6,15 +6,12 @@ import (
"time"
)
// Runner es cualquier cosa que sepa ejecutar un ciclo. El Service lo cumple.
type Runner interface {
Run(ctx context.Context) error
}
// cycleTimeout: tope máximo de tiempo para un solo ciclo, independiente del intervalo entre ciclos.
const cycleTimeout = 4 * time.Minute
// Scheduler dispara el Runner según su Schedule (intervalo fijo u hora fija del día).
type Scheduler struct {
runner Runner
schedule Schedule
......@@ -24,7 +21,6 @@ func NewScheduler(runner Runner, schedule Schedule) *Scheduler {
return &Scheduler{runner: runner, schedule: schedule}
}
// Start corre en bucle disparando un ciclo cada vez que el Schedule lo indica, hasta cancelar ctx.
func (s *Scheduler) Start(ctx context.Context) {
timer := time.NewTimer(s.schedule.Next(time.Now()))
defer timer.Stop()
......@@ -40,7 +36,6 @@ func (s *Scheduler) Start(ctx context.Context) {
}
}
// runOnce ejecuta un ciclo con su propio timeout, registrando el error sin detener el bucle.
func (s *Scheduler) runOnce(ctx context.Context) {
runCtx, cancel := context.WithTimeout(ctx, cycleTimeout)
defer cancel()
......
package usecase
import (
"context"
"fmt"
"time"
"intranet-synchronizer/internal/domain/entities"
"intranet-synchronizer/internal/domain/repositories"
)
type UpsertOutcome string
const (
UpsertCreated UpsertOutcome = "created"
UpsertUpdated UpsertOutcome = "updated"
UpsertUnchanged UpsertOutcome = "unchanged"
UpsertNotFound UpsertOutcome = "not_found"
)
// UpsertDocumentUseCase sincroniza un solo documento bajo demanda (sin esperar al ciclo
// programado), releyéndolo del origen Postgres por id y reflejándolo en Mongo.
type UpsertDocumentUseCase interface {
Upsert(ctx context.Context, id int) (UpsertOutcome, error)
}
type upsertDocument struct {
reader repositories.SourceRowReader
upserter repositories.DocumentUpserter
fetcher repositories.DocumentFetcher
now func() time.Time
}
func NewUpsertDocumentUseCase(
reader repositories.SourceRowReader,
upserter repositories.DocumentUpserter,
fetcher repositories.DocumentFetcher,
now func() time.Time,
) UpsertDocumentUseCase {
return &upsertDocument{reader: reader, upserter: upserter, fetcher: fetcher, now: now}
}
func (s *upsertDocument) Upsert(ctx context.Context, id int) (UpsertOutcome, error) {
row, found, err := s.reader.FetchByID(ctx, id)
if err != nil {
return "", fmt.Errorf("fetch by id: %w", err)
}
if !found {
return UpsertNotFound, nil
}
hash, err := computeRowHash(row.Doc)
if err != nil {
return "", err
}
existing, err := s.fetcher.FindByIDs(ctx, []int{id})
if err != nil {
return "", fmt.Errorf("fetch current doc: %w", err)
}
current, exists := existing[id]
if exists {
if storedHash, _ := current["row_hash"].(string); storedHash == hash {
return UpsertUnchanged, nil
}
}
runAt := s.now().UTC()
stampUpdatedAt(row.Doc, runAt)
stampRowHash(row.Doc, hash)
stampCreatedAt(row.Doc, runAt)
if _, _, err := s.upserter.BulkUpsert(ctx, []entities.SyncRow{row}); err != nil {
return "", fmt.Errorf("upsert: %w", err)
}
if exists {
return UpsertUpdated, nil
}
return UpsertCreated, nil
}
// Package changelog escribe un archivo de texto con el detalle de cada cambio aplicado
// por un ciclo de sincronización (JSON insertado, o JSON anterior/nuevo en las actualizaciones).
package changelog
import (
......@@ -13,16 +11,13 @@ import (
"go.mongodb.org/mongo-driver/bson"
)
// DefaultDir es la carpeta donde se dejan los archivos de log de cambios.
const DefaultDir = "/opt/logs"
// Writer agrega entradas al archivo <dir>/sync_<name>.log de un pipeline.
type Writer struct {
mu sync.Mutex
path string
}
// New crea el Writer de un pipeline, asegurando que la carpeta exista.
func New(dir, name string) (*Writer, error) {
if dir == "" {
dir = DefaultDir
......@@ -33,22 +28,18 @@ func New(dir, name string) (*Writer, error) {
return &Writer{path: filepath.Join(dir, fmt.Sprintf("sync_%s.log", name))}, nil
}
// Path devuelve la ruta del archivo de log.
func (w *Writer) Path() string {
return w.path
}
// Insert registra el documento insertado.
func (w *Writer) Insert(id int, doc interface{}) error {
return w.write(fmt.Sprintf("INSERT id=%d\n new: %s\n", id, marshal(doc)))
}
// Update registra el documento anterior y el nuevo.
func (w *Writer) Update(id int, before, after interface{}) error {
return w.write(fmt.Sprintf("UPDATE id=%d\n old: %s\n new: %s\n", id, marshal(before), marshal(after)))
}
// write agrega una entrada con su timestamp al archivo, en modo append.
func (w *Writer) write(entry string) error {
w.mu.Lock()
defer w.mu.Unlock()
......@@ -63,10 +54,6 @@ func (w *Writer) write(entry string) error {
return err
}
// marshal serializa a JSON en una línea usando los MISMOS nombres de campo que la colección Mongo.
// Las entidades tienen tags bson (no json), así que se pasa por un round trip bson -> map para que
// el documento nuevo y el anterior (leído de Mongo) sean comparables campo a campo con cualquier diff.
// Si falla, devuelve el error como texto para no perder la entrada.
func marshal(v interface{}) string {
doc, err := toCollectionShape(v)
if err != nil {
......@@ -79,8 +66,6 @@ func marshal(v interface{}) string {
return string(b)
}
// toCollectionShape normaliza cualquier documento (entidad tipada o bson.M leído de Mongo)
// a un map con los nombres de campo de la colección, descartando `_id` (interno de Mongo).
func toCollectionShape(v interface{}) (map[string]interface{}, error) {
if v == nil {
return nil, nil
......
// Package jsonutil trae helpers chicos para decodificar JSON crudo devuelto
// por Postgres (JSON_AGG), compartidos entre los módulos de sincronización.
package jsonutil
import "encoding/json"
// DecodeList convierte bytes JSON (ej. de JSON_AGG) a []T. Si raw es nil
// (columna NULL), devuelve una lista vacía en vez de error.
func DecodeList[T any](raw []byte) ([]T, error) {
if raw == nil {
return []T{}, nil
......
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