chore: WIP snapshot antes de reestructurar a arquitectura limpia

parent 01fc5a13
# intranet-sycronizacion # intranet-synchronizer
Sincroniza estudiantes (con padres y planes de pago) desde PostgreSQL hacia MongoDB mediante sondeo completo cada N minutos. Sincroniza estudiantes (con padres y planes de pago) desde PostgreSQL hacia MongoDB mediante sondeo completo cada N minutos.
......
...@@ -13,15 +13,14 @@ import ( ...@@ -13,15 +13,14 @@ import (
"github.com/joho/godotenv" "github.com/joho/godotenv"
"intranet-sycronizacion/internal/api" "intranet-synchronizer/internal/api"
"intranet-sycronizacion/internal/config" "intranet-synchronizer/internal/config"
"intranet-sycronizacion/internal/db" "intranet-synchronizer/internal/db"
"intranet-sycronizacion/internal/parents" "intranet-synchronizer/internal/parents"
"intranet-sycronizacion/internal/paymentplans" "intranet-synchronizer/internal/professors"
"intranet-sycronizacion/internal/professors" "intranet-synchronizer/internal/students"
"intranet-sycronizacion/internal/students" appsync "intranet-synchronizer/internal/sync"
appsync "intranet-sycronizacion/internal/sync" "intranet-synchronizer/internal/users"
"intranet-sycronizacion/internal/users"
) )
func main() { func main() {
...@@ -58,21 +57,17 @@ func main() { ...@@ -58,21 +57,17 @@ func main() {
deletedColl := mongoClient.Database(cfg.MongoDB).Collection(p.MongoDeletedCollection) deletedColl := mongoClient.Database(cfg.MongoDB).Collection(p.MongoDeletedCollection)
var reader appsync.PGReader var reader appsync.PGReader
switch { switch p.Name {
case p.Name == "students": case "students":
reader = students.NewQueryReader(pgPool, p.PGQuery) reader = students.NewQueryReader(pgPool, p.PGQuery)
case p.Name == "payment_plans": case "parents":
reader = paymentplans.NewQueryReader(pgPool, p.PGQuery)
case p.Name == "parents":
reader = parents.NewQueryReader(pgPool, p.PGQuery) reader = parents.NewQueryReader(pgPool, p.PGQuery)
case p.Name == "users": case "users":
reader = users.NewQueryReader(pgPool, p.PGQuery) reader = users.NewQueryReader(pgPool, p.PGQuery)
case p.Name == "professors": case "professors":
reader = professors.NewQueryReader(pgPool, p.PGQuery) reader = professors.NewQueryReader(pgPool, p.PGQuery)
case p.RequireEnrollmentID:
reader = db.NewFunctionReader(pgPool, p.PGFunction)
default: default:
reader = db.NewGenericFunctionReader(pgPool, p.PGFunction, p.IDField) log.Fatalf("no reader wired for pipeline %q", p.Name)
} }
// upserter cumple a la vez MongoUpserter y MongoIDLister. // upserter cumple a la vez MongoUpserter y MongoIDLister.
...@@ -100,7 +95,7 @@ func main() { ...@@ -100,7 +95,7 @@ func main() {
svcByName[p.Name] = svc svcByName[p.Name] = svc
} }
router := api.NewRouter(svcByName["students"], svcByName["payment_plans"], svcByName["parents"], svcByName["users"], svcByName["professors"], router := api.NewRouter(svcByName["students"], svcByName["parents"], svcByName["users"], svcByName["professors"],
func(pingCtx context.Context) error { return pgPool.Ping(pingCtx) }, func(pingCtx context.Context) error { return pgPool.Ping(pingCtx) },
func(pingCtx context.Context) error { return mongoClient.Ping(pingCtx, nil) }, func(pingCtx context.Context) error { return mongoClient.Ping(pingCtx, nil) },
) )
......
module intranet-sycronizacion module intranet-synchronizer
go 1.25.0 go 1.25.0
......
...@@ -48,10 +48,6 @@ github.com/klauspost/compress v1.17.6 h1:60eq2E/jlfwQXtvZEeBUYADs+BwKBWURIY+Gj2e ...@@ -48,10 +48,6 @@ github.com/klauspost/compress v1.17.6 h1:60eq2E/jlfwQXtvZEeBUYADs+BwKBWURIY+Gj2e
github.com/klauspost/compress v1.17.6/go.mod h1:/dCuZOvVtNoHsyb+cuJD3itjs3NbnF6KH9zAO4BDxPM= github.com/klauspost/compress v1.17.6/go.mod h1:/dCuZOvVtNoHsyb+cuJD3itjs3NbnF6KH9zAO4BDxPM=
github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y= github.com/klauspost/cpuid/v2 v2.3.0 h1:S4CRMLnYUhGeDFDqkGriYKdfoFlDnMtqTiI/sFzhA9Y=
github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0= github.com/klauspost/cpuid/v2 v2.3.0/go.mod h1:hqwkgyIinND0mEev00jJYCxPNVRVXFQeu1XKlok6oO0=
github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
github.com/leodido/go-urn v1.4.0 h1:WT9HwE9SGECu3lg4d/dIA+jxlljEa1/ffXKmRjqdmIQ= github.com/leodido/go-urn v1.4.0 h1:WT9HwE9SGECu3lg4d/dIA+jxlljEa1/ffXKmRjqdmIQ=
github.com/leodido/go-urn v1.4.0/go.mod h1:bvxc+MVxLKB4z00jd1z+Dvzr47oO32F/QSNjSBOlFxI= github.com/leodido/go-urn v1.4.0/go.mod h1:bvxc+MVxLKB4z00jd1z+Dvzr47oO32F/QSNjSBOlFxI=
github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY= github.com/mattn/go-isatty v0.0.20 h1:xfD0iDuEKnDkl03q4limB+vH+GxLEtL/jb4xVJSWWEY=
...@@ -71,8 +67,6 @@ github.com/quic-go/qpack v0.6.0 h1:g7W+BMYynC1LbYLSqRt8PBg5Tgwxn214ZZR34VIOjz8= ...@@ -71,8 +67,6 @@ github.com/quic-go/qpack v0.6.0 h1:g7W+BMYynC1LbYLSqRt8PBg5Tgwxn214ZZR34VIOjz8=
github.com/quic-go/qpack v0.6.0/go.mod h1:lUpLKChi8njB4ty2bFLX2x4gzDqXwUpaO1DP9qMDZII= github.com/quic-go/qpack v0.6.0/go.mod h1:lUpLKChi8njB4ty2bFLX2x4gzDqXwUpaO1DP9qMDZII=
github.com/quic-go/quic-go v0.59.0 h1:OLJkp1Mlm/aS7dpKgTc6cnpynnD2Xg7C1pwL6vy/SAw= github.com/quic-go/quic-go v0.59.0 h1:OLJkp1Mlm/aS7dpKgTc6cnpynnD2Xg7C1pwL6vy/SAw=
github.com/quic-go/quic-go v0.59.0/go.mod h1:upnsH4Ju1YkqpLXC305eW3yDZ4NfnNbmQRCMWS58IKU= github.com/quic-go/quic-go v0.59.0/go.mod h1:upnsH4Ju1YkqpLXC305eW3yDZ4NfnNbmQRCMWS58IKU=
github.com/rogpeppe/go-internal v1.10.0 h1:TMyTOH3F/DB16zRVcYyreMH6GnZZrwQVAoYjRBZyWFQ=
github.com/rogpeppe/go-internal v1.10.0/go.mod h1:UQnix2H7Ngw/k4C5ijL5+65zddjncjaFoBhdsK/akog=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw= github.com/stretchr/objx v0.4.0/go.mod h1:YvHI0jy2hoMjB+UWwv71VJQ9isScKT/TqJzVSSt89Yw=
github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo= github.com/stretchr/objx v0.5.0/go.mod h1:Yh+to48EsGEfYuaHDzXPcE3xhTkx73EhmCGUpEOglKo=
...@@ -143,8 +137,6 @@ golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8T ...@@ -143,8 +137,6 @@ golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8T
google.golang.org/protobuf v1.36.10 h1:AYd7cD/uASjIL6Q9LiTjz8JLcrh/88q5UObnmY3aOOE= google.golang.org/protobuf v1.36.10 h1:AYd7cD/uASjIL6Q9LiTjz8JLcrh/88q5UObnmY3aOOE=
google.golang.org/protobuf v1.36.10/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= google.golang.org/protobuf v1.36.10/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk=
gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q=
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
...@@ -6,13 +6,12 @@ import ( ...@@ -6,13 +6,12 @@ import (
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
appsync "intranet-sycronizacion/internal/sync" appsync "intranet-synchronizer/internal/sync"
) )
// Handlers guarda las dependencias de los manejadores HTTP. // Handlers guarda las dependencias de los manejadores HTTP.
type Handlers struct { type Handlers struct {
studentsSvc *appsync.Service studentsSvc *appsync.Service
paymentPlansSvc *appsync.Service
parentsSvc *appsync.Service parentsSvc *appsync.Service
usersSvc *appsync.Service usersSvc *appsync.Service
professorsSvc *appsync.Service professorsSvc *appsync.Service
......
...@@ -6,12 +6,12 @@ import ( ...@@ -6,12 +6,12 @@ import (
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
appsync "intranet-sycronizacion/internal/sync" appsync "intranet-synchronizer/internal/sync"
) )
// NewRouter construye el router de Gin y registra las rutas. // NewRouter construye el router de Gin y registra las rutas.
func NewRouter(studentsSvc, paymentPlansSvc, parentsSvc, usersSvc, professorsSvc *appsync.Service, pgPing, mongoPing func(context.Context) error) *gin.Engine { func NewRouter(studentsSvc, parentsSvc, usersSvc, professorsSvc *appsync.Service, pgPing, mongoPing func(context.Context) error) *gin.Engine {
h := &Handlers{studentsSvc: studentsSvc, paymentPlansSvc: paymentPlansSvc, parentsSvc: parentsSvc, usersSvc: usersSvc, professorsSvc: professorsSvc, pgPing: pgPing, mongoPing: mongoPing} h := &Handlers{studentsSvc: studentsSvc, parentsSvc: parentsSvc, usersSvc: usersSvc, professorsSvc: professorsSvc, pgPing: pgPing, mongoPing: mongoPing}
r := gin.Default() r := gin.Default()
...@@ -21,9 +21,6 @@ func NewRouter(studentsSvc, paymentPlansSvc, parentsSvc, usersSvc, professorsSvc ...@@ -21,9 +21,6 @@ func NewRouter(studentsSvc, paymentPlansSvc, parentsSvc, usersSvc, professorsSvc
v1.POST("/sync/students/trigger", triggerHandler(studentsSvc)) v1.POST("/sync/students/trigger", triggerHandler(studentsSvc))
v1.GET("/sync/students/status", statusHandler(studentsSvc)) v1.GET("/sync/students/status", statusHandler(studentsSvc))
v1.POST("/sync/payment_plans/trigger", triggerHandler(paymentPlansSvc))
v1.GET("/sync/payment_plans/status", statusHandler(paymentPlansSvc))
v1.POST("/sync/parents/trigger", triggerHandler(parentsSvc)) v1.POST("/sync/parents/trigger", triggerHandler(parentsSvc))
v1.GET("/sync/parents/status", statusHandler(parentsSvc)) v1.GET("/sync/parents/status", statusHandler(parentsSvc))
......
...@@ -8,7 +8,7 @@ import ( ...@@ -8,7 +8,7 @@ import (
"strings" "strings"
"time" "time"
appsync "intranet-sycronizacion/internal/sync" appsync "intranet-synchronizer/internal/sync"
) )
// Config guarda los valores de configuración ya leídos y validados. // Config guarda los valores de configuración ya leídos y validados.
......
package config package config
import ( import (
"intranet-sycronizacion/internal/parents" "intranet-synchronizer/internal/parents"
"intranet-sycronizacion/internal/paymentplans" "intranet-synchronizer/internal/professors"
"intranet-sycronizacion/internal/professors" "intranet-synchronizer/internal/students"
"intranet-sycronizacion/internal/students" "intranet-synchronizer/internal/users"
"intranet-sycronizacion/internal/users"
) )
// PipelineDef describe una tubería de sincronización: origen Postgres -> destino Mongo. // PipelineDef describe una tubería de sincronización: origen Postgres -> destino Mongo.
type PipelineDef struct { type PipelineDef struct {
Name string Name string
PGFunction string
PGQuery string PGQuery string
MongoCollection string MongoCollection string
MongoDeletedCollection string MongoDeletedCollection string
IDField string IDField string
StateID int StateID int
RequireEnrollmentID bool
} }
// Pipelines es la lista de todas las tuberías activas. // Pipelines es la lista de todas las tuberías activas.
...@@ -31,14 +28,6 @@ var Pipelines = []PipelineDef{ ...@@ -31,14 +28,6 @@ var Pipelines = []PipelineDef{
StateID: 1, StateID: 1,
}, },
{ {
Name: "payment_plans",
PGQuery: paymentplans.Query,
MongoCollection: "payment_plans",
MongoDeletedCollection: "deleted_payment_plans",
IDField: "payment_plan_id",
StateID: 2,
},
{
Name: "parents", Name: "parents",
PGQuery: parents.Query, PGQuery: parents.Query,
MongoCollection: "parents", MongoCollection: "parents",
......
...@@ -10,7 +10,7 @@ import ( ...@@ -10,7 +10,7 @@ import (
"go.mongodb.org/mongo-driver/mongo" "go.mongodb.org/mongo-driver/mongo"
"go.mongodb.org/mongo-driver/mongo/options" "go.mongodb.org/mongo-driver/mongo/options"
appsync "intranet-sycronizacion/internal/sync" appsync "intranet-synchronizer/internal/sync"
) )
// NewMongoClient conecta a MongoDB y verifica con un Ping. // NewMongoClient conecta a MongoDB y verifica con un Ping.
......
...@@ -3,12 +3,11 @@ package db ...@@ -3,12 +3,11 @@ package db
import ( import (
"context" "context"
"encoding/json"
"fmt" "fmt"
"github.com/jackc/pgx/v5/pgxpool" "github.com/jackc/pgx/v5/pgxpool"
appsync "intranet-sycronizacion/internal/sync" appsync "intranet-synchronizer/internal/sync"
) )
// Pool envuelve el pool de pgx para poder colgarle métodos propios. // Pool envuelve el pool de pgx para poder colgarle métodos propios.
...@@ -47,7 +46,7 @@ CREATE TABLE IF NOT EXISTS sync_state ( ...@@ -47,7 +46,7 @@ CREATE TABLE IF NOT EXISTS sync_state (
); );
INSERT INTO sync_state (id, last_synced_at) INSERT INTO sync_state (id, last_synced_at)
VALUES (1, '1970-01-01T00:00:00Z'), (2, '1970-01-01T00:00:00Z'), (3, '1970-01-01T00:00:00Z') VALUES (1, '1970-01-01T00:00:00Z'), (3, '1970-01-01T00:00:00Z'), (4, '1970-01-01T00:00:00Z'), (5, '1970-01-01T00:00:00Z')
ON CONFLICT (id) DO NOTHING; ON CONFLICT (id) DO NOTHING;
` `
...@@ -60,100 +59,4 @@ func (p *Pool) BootstrapSchema(ctx context.Context) error { ...@@ -60,100 +59,4 @@ func (p *Pool) BootstrapSchema(ctx context.Context) error {
return nil return nil
} }
// FunctionReader (legado) llama una función Postgres que devuelve un sobre JSON, exige student_id + enrollment_id. // Los readers por módulo (students, parents, users, professors) viven en internal/<módulo>.
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 PG: status/message/data.
type functionEnvelope struct {
Status bool `json:"status"`
Message string `json:"message"`
Data []map[string]interface{} `json:"data"`
}
// parseFunctionEnvelope convierte el JSON crudo en filas, exigiendo student_id y enrollment_id numéricos.
func parseFunctionEnvelope(raw []byte) ([]appsync.SyncRow, error) {
var env functionEnvelope
if err := json.Unmarshal(raw, &env); err != nil {
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["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 solo un id configurable (idField).
func parseGenericEnvelope(raw []byte, idField string) ([]appsync.SyncRow, error) {
var env functionEnvelope
if err := json.Unmarshal(raw, &env); err != nil {
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 Postgres cuyo sobre solo requiere un id numérico.
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)
}
// Los readers por módulo (students, payment_plans, parents) viven en internal/<módulo>.
...@@ -4,9 +4,9 @@ import ( ...@@ -4,9 +4,9 @@ import (
"context" "context"
"fmt" "fmt"
"intranet-sycronizacion/internal/db" "intranet-synchronizer/internal/db"
"intranet-sycronizacion/internal/jsonutil" "intranet-synchronizer/internal/jsonutil"
appsync "intranet-sycronizacion/internal/sync" appsync "intranet-synchronizer/internal/sync"
) )
// QueryReader ejecuta Query directamente y arma una SyncRow por fila del resultado. // QueryReader ejecuta Query directamente y arma una SyncRow por fila del resultado.
......
// Package paymentplans contiene la entidad, query SQL y reader de la colección "payment_plans".
package paymentplans
import "time"
// Record es la entidad de la colección "payment_plans": un documento por plan de pago.
type Record struct {
PaymentPlanID int `bson:"payment_plan_id"`
StudentID int `bson:"student_id"`
EnrollmentID int `bson:"enrollment_id"`
Order int `bson:"payment_plan_order"`
Concept *string `bson:"payment_plan_concept"`
Year int `bson:"payment_plan_year"`
Date *string `bson:"payment_plan_date"`
AdditionalConcepts *string `bson:"payment_plan_additional_concepts"`
IssueDate *string `bson:"payment_plan_issue_date"`
Total *float64 `bson:"payment_plan_total"`
Debt *int `bson:"payment_plan_debt"`
UpdatedAt time.Time `bson:"updated_at"`
RowHash string `bson:"row_hash,omitempty" json:"-"`
}
// SetUpdatedAt estampa la hora de sincronización.
func (r *Record) SetUpdatedAt(t time.Time) {
r.UpdatedAt = t
}
// SetRowHash guarda el hash de contenido usado para saltar upserts sin cambios.
func (r *Record) SetRowHash(h string) {
r.RowHash = h
}
package paymentplans
import (
"context"
"fmt"
"intranet-sycronizacion/internal/db"
appsync "intranet-sycronizacion/internal/sync"
)
// QueryReader ejecuta Query directamente y arma una SyncRow por fila del resultado.
type QueryReader struct {
pool *db.Pool
query string
}
// NewQueryReader construye el reader de payment_plans sobre el pool dado.
func NewQueryReader(pool *db.Pool, query string) *QueryReader {
return &QueryReader{pool: pool, query: query}
}
// FetchAll corre la consulta y lee fila por fila.
func (f *QueryReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) {
rows, err := f.pool.Pool.Query(ctx, f.query)
if err != nil {
return nil, fmt.Errorf("query payment_plans: %w", err)
}
defer rows.Close()
var result []appsync.SyncRow
for rows.Next() {
var (
paymentPlanID int
studentID int
enrollmentID int
order int
concept *string
year int
date *string
additionalConcepts *string
issueDate *string
total *float64
debt *int
)
if err := rows.Scan(&paymentPlanID, &studentID, &enrollmentID, &order, &concept, &year,
&date, &additionalConcepts, &issueDate, &total, &debt); err != nil {
return nil, fmt.Errorf("scan payment_plan row: %w", err)
}
doc := &Record{
PaymentPlanID: paymentPlanID,
StudentID: studentID,
EnrollmentID: enrollmentID,
Order: order,
Concept: concept,
Year: year,
Date: date,
AdditionalConcepts: additionalConcepts,
IssueDate: issueDate,
Total: total,
Debt: debt,
}
result = append(result, appsync.SyncRow{ID: paymentPlanID, Doc: doc})
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate payment_plan rows: %w", err)
}
return result, nil
}
...@@ -4,8 +4,8 @@ import ( ...@@ -4,8 +4,8 @@ import (
"context" "context"
"fmt" "fmt"
"intranet-sycronizacion/internal/db" "intranet-synchronizer/internal/db"
appsync "intranet-sycronizacion/internal/sync" appsync "intranet-synchronizer/internal/sync"
) )
// QueryReader ejecuta Query directamente y arma una SyncRow por fila del resultado. // QueryReader ejecuta Query directamente y arma una SyncRow por fila del resultado.
......
...@@ -23,10 +23,30 @@ type Student struct { ...@@ -23,10 +23,30 @@ type Student struct {
GradeID *int `bson:"grade_id"` GradeID *int `bson:"grade_id"`
Grade *string `bson:"grade"` Grade *string `bson:"grade"`
ClassroomAttendanceID *int `bson:"classroom_attendance_id"` ClassroomAttendanceID *int `bson:"classroom_attendance_id"`
PaymentPlans []PaymentPlanYear `bson:"payment_plans"`
UpdatedAt time.Time `bson:"updated_at"` UpdatedAt time.Time `bson:"updated_at"`
RowHash string `bson:"row_hash,omitempty" json:"-"` 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"`
Order int `bson:"payment_plan_order"`
Concept *string `bson:"payment_plan_concept"`
Date *string `bson:"payment_plan_date"`
AdditionalConcepts *string `bson:"payment_plan_additional_concepts"`
IssueDate *string `bson:"payment_plan_issue_date"`
Total *float64 `bson:"payment_plan_total"`
Debt *int `bson:"payment_plan_debt"`
}
// SetUpdatedAt estampa la hora de sincronización. // SetUpdatedAt estampa la hora de sincronización.
func (s *Student) SetUpdatedAt(t time.Time) { func (s *Student) SetUpdatedAt(t time.Time) {
s.UpdatedAt = t s.UpdatedAt = t
......
package students package students
// Query es el SQL crudo de la tubería de estudiantes: junta estudiante + persona + matrícula + apertura (sede/nivel/grado). // Query es el SQL crudo de la tubería de estudiantes: junta estudiante + persona + matrícula + apertura (sede/nivel/grado).
// El reader corre además paymentplans.Query por separado y anida sus filas en Student.PaymentPlans por student_id/año
// (ver reader.go) — separado por rendimiento, no se puede traer todo en un solo query sin duplicar filas de estudiante
// por cada plan de pago.
const Query = ` const Query = `
SELECT e.estudiante_id as student_id, SELECT e.estudiante_id as student_id,
e.estudiante_codigo as student_code, e.estudiante_codigo as student_code,
...@@ -11,13 +14,9 @@ SELECT e.estudiante_id as student_id, ...@@ -11,13 +14,9 @@ SELECT e.estudiante_id as student_id,
pp.persona_numero_documento_identidad as student_dni, pp.persona_numero_documento_identidad as student_dni,
pp.persona_fecha_nacimiento as student_birthday, pp.persona_fecha_nacimiento as student_birthday,
pp.persona_correo as student_email, pp.persona_correo as student_email,
sexo.catalogo_siiaa_id as genre_id,
sexo.catalogo_siiaa_nombre as genre, sexo.catalogo_siiaa_nombre as genre,
ac.sede_id as branch_id,
public.to_camel_case(se.sede_nombre) as branch, public.to_camel_case(se.sede_nombre) as branch,
ac.nivel_id as level_id,
public.to_camel_case(nivel.catalogo_siiaa_nombre) as level, public.to_camel_case(nivel.catalogo_siiaa_nombre) as level,
ac.grado_id as grade_id,
public.to_camel_case(grado.grado_sigla) as grade, public.to_camel_case(grado.grado_sigla) as grade,
ma.aula_id_asiste as classroom_attendance_id ma.aula_id_asiste as classroom_attendance_id
FROM matricula.ma_estudiante e FROM matricula.ma_estudiante e
......
...@@ -3,13 +3,16 @@ package students ...@@ -3,13 +3,16 @@ package students
import ( import (
"context" "context"
"fmt" "fmt"
"sort"
"time" "time"
"intranet-sycronizacion/internal/db" "intranet-synchronizer/internal/db"
appsync "intranet-sycronizacion/internal/sync" "intranet-synchronizer/internal/paymentplans"
appsync "intranet-synchronizer/internal/sync"
) )
// QueryReader ejecuta Query directamente y arma una SyncRow por fila del resultado. // QueryReader ejecuta Query + paymentplans.Query y arma una SyncRow por estudiante,
// anidando sus planes de pago agrupados por año.
type QueryReader struct { type QueryReader struct {
pool *db.Pool pool *db.Pool
query string query string
...@@ -20,15 +23,37 @@ func NewQueryReader(pool *db.Pool, query string) *QueryReader { ...@@ -20,15 +23,37 @@ func NewQueryReader(pool *db.Pool, query string) *QueryReader {
return &QueryReader{pool: pool, query: query} return &QueryReader{pool: pool, query: query}
} }
// FetchAll corre la consulta y lee fila por fila. // FetchAll corre la consulta de estudiantes y la de planes de pago, y anida la segunda dentro de la primera.
func (f *QueryReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) { func (f *QueryReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) {
students, order, err := f.fetchStudents(ctx)
if err != nil {
return nil, err
}
plansByStudent, err := f.fetchPaymentPlans(ctx)
if err != nil {
return nil, err
}
result := make([]appsync.SyncRow, 0, len(order))
for _, studentID := range order {
doc := students[studentID]
doc.PaymentPlans = buildPaymentPlanYears(plansByStudent[studentID])
result = append(result, appsync.SyncRow{ID: studentID, Doc: doc})
}
return result, nil
}
// fetchStudents lee la fila base de cada estudiante (sin planes de pago).
func (f *QueryReader) fetchStudents(ctx context.Context) (map[int]*Student, []int, error) {
rows, err := f.pool.Pool.Query(ctx, f.query) rows, err := f.pool.Pool.Query(ctx, f.query)
if err != nil { if err != nil {
return nil, fmt.Errorf("query students: %w", err) return nil, nil, fmt.Errorf("query students: %w", err)
} }
defer rows.Close() defer rows.Close()
var result []appsync.SyncRow students := make(map[int]*Student)
var order []int
for rows.Next() { for rows.Next() {
var ( var (
studentID int studentID int
...@@ -53,10 +78,10 @@ func (f *QueryReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) { ...@@ -53,10 +78,10 @@ func (f *QueryReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) {
if err := rows.Scan(&studentID, &studentCode, &studentInternalCode, &paternalLastName, &maternalLastName, if err := rows.Scan(&studentID, &studentCode, &studentInternalCode, &paternalLastName, &maternalLastName,
&name, &dni, &birthday, &email, &genreID, &genre, &branchID, &branch, &levelID, &level, &name, &dni, &birthday, &email, &genreID, &genre, &branchID, &branch, &levelID, &level,
&gradeID, &grade, &classroomAttendanceID); err != nil { &gradeID, &grade, &classroomAttendanceID); err != nil {
return nil, fmt.Errorf("scan student row: %w", err) return nil, nil, fmt.Errorf("scan student row: %w", err)
} }
doc := &Student{ students[studentID] = &Student{
StudentID: studentID, StudentID: studentID,
StudentCode: studentCode, StudentCode: studentCode,
StudentInternalCode: studentInternalCode, StudentInternalCode: studentInternalCode,
...@@ -76,10 +101,85 @@ func (f *QueryReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) { ...@@ -76,10 +101,85 @@ func (f *QueryReader) FetchAll(ctx context.Context) ([]appsync.SyncRow, error) {
Grade: grade, Grade: grade,
ClassroomAttendanceID: classroomAttendanceID, ClassroomAttendanceID: classroomAttendanceID,
} }
result = append(result, appsync.SyncRow{ID: studentID, Doc: doc}) order = append(order, studentID)
} }
if err := rows.Err(); err != nil { if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate student rows: %w", err) return nil, nil, fmt.Errorf("iterate student rows: %w", err)
} }
return result, nil return students, order, nil
}
// paymentRow es una fila de paymentplans.Query antes de agruparse por año.
type paymentRow struct {
Payment
year int
}
// fetchPaymentPlans lee todos los planes de pago y los agrupa por student_id.
func (f *QueryReader) fetchPaymentPlans(ctx context.Context) (map[int][]paymentRow, error) {
rows, err := f.pool.Pool.Query(ctx, paymentplans.Query)
if err != nil {
return nil, fmt.Errorf("query payment_plans: %w", err)
}
defer rows.Close()
plansByStudent := make(map[int][]paymentRow)
for rows.Next() {
var (
paymentPlanID int
studentID int
enrollmentID int
order int
concept *string
year int
date *string
additionalConcepts *string
issueDate *string
total *float64
debt *int
)
if err := rows.Scan(&paymentPlanID, &studentID, &enrollmentID, &order, &concept, &year,
&date, &additionalConcepts, &issueDate, &total, &debt); err != nil {
return nil, fmt.Errorf("scan payment_plan row: %w", err)
}
plansByStudent[studentID] = append(plansByStudent[studentID], paymentRow{
Payment: Payment{
PaymentPlanID: paymentPlanID,
EnrollmentID: enrollmentID,
Order: order,
Concept: concept,
Date: date,
AdditionalConcepts: additionalConcepts,
IssueDate: issueDate,
Total: total,
Debt: debt,
},
year: year,
})
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate payment_plan rows: %w", err)
}
return plansByStudent, nil
}
// buildPaymentPlanYears agrupa los pagos de un estudiante por año, orden descendente.
func buildPaymentPlanYears(payments []paymentRow) []PaymentPlanYear {
byYear := make(map[int][]Payment)
for _, p := range payments {
byYear[p.year] = append(byYear[p.year], p.Payment)
}
years := make([]int, 0, len(byYear))
for y := range byYear {
years = append(years, y)
}
sort.Sort(sort.Reverse(sort.IntSlice(years)))
result := make([]PaymentPlanYear, 0, len(years))
for _, y := range years {
result = append(result, PaymentPlanYear{Year: y, Payments: byYear[y]})
}
return result
} }
...@@ -5,8 +5,8 @@ import ( ...@@ -5,8 +5,8 @@ import (
"fmt" "fmt"
"time" "time"
"intranet-sycronizacion/internal/db" "intranet-synchronizer/internal/db"
appsync "intranet-sycronizacion/internal/sync" appsync "intranet-synchronizer/internal/sync"
) )
// QueryReader ejecuta Query directamente y arma una SyncRow por fila del resultado. // QueryReader ejecuta Query directamente y arma una SyncRow por fila del resultado.
......
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