[ADD] cambios

parent d76900b4
Subproject commit b98ee98b582510a212e30ce7e230f019ec9e22c0
Subproject commit 3647943bdf49414e5868b6638ba6912d0387cd75
// 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
}
var out []T
if err := json.Unmarshal(raw, &out); err != nil {
return nil, err
}
return out, nil
}
// Package parents contiene la entidad, query SQL y reader de la colección "parents".
package parents
import "time"
// StudentRef es un estudiante asociado a un padre/apoderado, embebido dentro de Record.
type StudentRef struct {
StudentID int `bson:"student_id" json:"student_id"`
}
// Record es la entidad de la colección "parents".
type Record struct {
ParentID int `bson:"parent_id"`
DNI *string `bson:"parent_dni"`
PaternalLastName *string `bson:"parent_paternal_last_name"`
MaternalLastName *string `bson:"parent_maternal_last_name"`
Name *string `bson:"parent_name"`
Email *string `bson:"parent_email"`
Phone *string `bson:"parent_phone"`
Students []StudentRef `bson:"students"`
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 parents
// Query es el SQL crudo de la tubería de padres/apoderados, con estudiantes a cargo embebidos.
const Query = `
SELECT pp.persona_id as parent_id,
pp.persona_numero_documento_identidad as parent_dni,
public.to_camel_case(pp.persona_apellido_paterno) as parent_paternal_last_name,
public.to_camel_case(pp.persona_apellido_materno) as parent_maternal_last_name,
public.to_camel_case(pp.persona_nombre) as parent_name,
pp.persona_correo as parent_email,
pp.persona_telefono as parent_phone,
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
FROM matricula.ma_estudiante_apoderado apoderado
WHERE apoderado.estudiante_apoderado_estado = '1'
GROUP BY apoderado.persona_id) est ON est.persona_id = pp.persona_id;
`
package parents
import (
"context"
"fmt"
"intranet-sycronizacion/internal/db"
"intranet-sycronizacion/internal/jsonutil"
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 parents 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 parents: %w", err)
}
defer rows.Close()
var result []appsync.SyncRow
for rows.Next() {
var (
parentID int
dni *string
paternalLastName *string
maternalLastName *string
name *string
email *string
phone *string
students []byte
)
if err := rows.Scan(&parentID, &dni, &paternalLastName, &maternalLastName, &name, &email, &phone, &students); err != nil {
return nil, fmt.Errorf("scan parent row: %w", err)
}
studentRefs, err := jsonutil.DecodeList[StudentRef](students)
if err != nil {
return nil, fmt.Errorf("decode students for parent %d: %w", parentID, err)
}
doc := &Record{
ParentID: parentID,
DNI: dni,
PaternalLastName: paternalLastName,
MaternalLastName: maternalLastName,
Name: name,
Email: email,
Phone: phone,
Students: studentRefs,
}
result = append(result, appsync.SyncRow{ID: parentID, Doc: doc})
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate parent rows: %w", err)
}
return result, nil
}
// 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
// Query es el SQL crudo de la tubería de planes de pago: junta caja.ca_plan_de_pago con matricula.ma_matricula.
const Query = `
WITH tb_lista_plan_pago AS (select c.plan_de_pago_id as payment_plan_id,
m.estudiante_id as student_id,
c.matricula_id as enrollment_id,
COALESCE(c.plan_de_pago_num_orden,0) as payment_plan_order,
UPPER(c.plan_de_pago_concepto) as payment_plan_concept,
c.plan_de_pago_anio as payment_plan_year,
to_char(c.plan_de_pago_fecha_pago, 'dd-mm-YYYY') as payment_plan_date,
UPPER(COALESCE(c.plan_de_pago_conceptos_adicionales, '') ||
COALESCE(c.plan_de_pago_conceptos_adicionales3, '')) as payment_plan_additional_concepts,
to_char(plan_de_pago_fecha_ven, 'dd-mm-YYYY') as payment_plan_issue_date,
c.plan_de_pago_subtotal as payment_plan_total,
CASE WHEN c.plan_de_pago_deuda THEN 1 ELSE 0 END as payment_plan_debt
FROM caja.ca_plan_de_pago c
INNER JOIN matricula.ma_matricula m ON c.matricula_id = m.matricula_id
WHERE plan_de_pago_anulado = false
and m.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;
`
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
}
// Package professors contiene la entidad, query SQL y reader de la colección "professors".
package professors
import "time"
// Record es la entidad de la colección "professors".
type Record 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"`
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 professors
// Query es el SQL crudo de la tubería de profesores.
const Query = `
SELECT profesor_id as professor_id,
public.to_camel_case(p.persona_apellido_paterno) as professor_paternal_last_name,
public.to_camel_case(p.persona_apellido_materno) as professor_maternal_last_name,
public.to_camel_case(p.persona_nombre) as professor_name,
p.persona_numero_documento_identidad as professor_dni
FROM horario.ho_profesor pro
INNER JOIN persona.pe_persona p on pro.persona_id = p.persona_id;
`
package professors
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 professors 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 professors: %w", err)
}
defer rows.Close()
var result []appsync.SyncRow
for rows.Next() {
var (
professorID int
paternalLastName *string
maternalLastName *string
name *string
dni *string
)
if err := rows.Scan(&professorID, &paternalLastName, &maternalLastName, &name, &dni); err != nil {
return nil, fmt.Errorf("scan professor row: %w", err)
}
doc := &Record{
ProfessorID: professorID,
ProfessorPaternalLastName: paternalLastName,
ProfessorMaternalLastName: maternalLastName,
ProfessorName: name,
ProfessorDNI: dni,
}
result = append(result, appsync.SyncRow{ID: professorID, Doc: doc})
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate professor rows: %w", err)
}
return result, nil
}
// Package students contiene la entidad, query SQL y reader de la colección "students".
package students
import "time"
// Student es la entidad de la colección "students".
type Student struct {
StudentID int `bson:"student_id"`
StudentCode *string `bson:"student_code"`
StudentInternalCode *string `bson:"student_internal_code"`
PaternalLastName *string `bson:"student_paternal_last_name"`
MaternalLastName *string `bson:"student_maternal_last_name"`
Name *string `bson:"student_name"`
DNI *string `bson:"student_dni"`
Birthday *time.Time `bson:"student_birthday"`
Email *string `bson:"student_email"`
GenreID *int `bson:"genre_id"`
Genre *string `bson:"genre"`
BranchID *int `bson:"branch_id"`
Branch *string `bson:"branch"`
LevelID *int `bson:"level_id"`
Level *string `bson:"level"`
GradeID *int `bson:"grade_id"`
Grade *string `bson:"grade"`
ClassroomAttendanceID *int `bson:"classroom_attendance_id"`
UpdatedAt time.Time `bson:"updated_at"`
RowHash string `bson:"row_hash,omitempty" json:"-"`
}
// 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
}
package students
// Query es el SQL crudo de la tubería de estudiantes: junta estudiante + persona + matrícula + apertura (sede/nivel/grado).
const Query = `
SELECT e.estudiante_id as student_id,
e.estudiante_codigo as student_code,
e.estudiante_codigo_interno as student_internal_code,
public.to_camel_case(pp.persona_apellido_paterno) as student_paternal_last_name,
public.to_camel_case(pp.persona_apellido_materno) as student_maternal_last_name,
public.to_camel_case(pp.persona_nombre) as student_name,
pp.persona_numero_documento_identidad as student_dni,
pp.persona_fecha_nacimiento as student_birthday,
pp.persona_correo as student_email,
sexo.catalogo_siiaa_id as genre_id,
sexo.catalogo_siiaa_nombre as genre,
ac.sede_id as branch_id,
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,
ac.grado_id as grade_id,
public.to_camel_case(grado.grado_sigla) as grade,
ma.aula_id_asiste as classroom_attendance_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
INNER JOIN matricula.ma_matricula ma on e.estudiante_id = ma.estudiante_id
INNER JOIN academico.ac_apertura ac on ma.apertura_id = ac.apertura_id
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 ma.periodo_academico_id >= 14;
`
package students
import (
"context"
"fmt"
"time"
"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 estudiantes 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 students: %w", err)
}
defer rows.Close()
var result []appsync.SyncRow
for rows.Next() {
var (
studentID int
studentCode *string
studentInternalCode *string
paternalLastName *string
maternalLastName *string
name *string
dni *string
birthday *time.Time
email *string
genreID *int
genre *string
branchID *int
branch *string
levelID *int
level *string
gradeID *int
grade *string
classroomAttendanceID *int
)
if err := rows.Scan(&studentID, &studentCode, &studentInternalCode, &paternalLastName, &maternalLastName,
&name, &dni, &birthday, &email, &genreID, &genre, &branchID, &branch, &levelID, &level,
&gradeID, &grade, &classroomAttendanceID); err != nil {
return nil, fmt.Errorf("scan student row: %w", err)
}
doc := &Student{
StudentID: studentID,
StudentCode: studentCode,
StudentInternalCode: studentInternalCode,
PaternalLastName: paternalLastName,
MaternalLastName: maternalLastName,
Name: name,
DNI: dni,
Birthday: birthday,
Email: email,
GenreID: genreID,
Genre: genre,
BranchID: branchID,
Branch: branch,
LevelID: levelID,
Level: level,
GradeID: gradeID,
Grade: grade,
ClassroomAttendanceID: classroomAttendanceID,
}
result = append(result, appsync.SyncRow{ID: studentID, Doc: doc})
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate student rows: %w", err)
}
return result, nil
}
package sync
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
}
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
}
func (s DailySchedule) Next(now time.Time) time.Duration {
next := time.Date(now.Year(), now.Month(), now.Day(), s.Hour, s.Minute, 0, 0, now.Location())
if !next.After(now) {
next = next.AddDate(0, 0, 1)
}
return next.Sub(now)
}
// Package users contiene la entidad, query SQL y reader de la colección "users".
package users
import "time"
// Record es la entidad de la colección "users".
type Record struct {
UserID int `bson:"user_id"`
Login *string `bson:"user_login"`
Password *string `bson:"user_password"`
CreationDate *time.Time `bson:"user_creation_date"`
ParentID *int `bson:"parent_id"`
Status *int `bson:"user_status"`
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 users
// Query es el SQL crudo de la tubería de usuarios.
const Query = `
SELECT user_id, user_login, user_password, user_creation_date, parent_id, user_status
FROM intranet.users;
`
package users
import (
"context"
"fmt"
"time"
"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 users 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 users: %w", err)
}
defer rows.Close()
var result []appsync.SyncRow
for rows.Next() {
var (
userID int
login *string
password *string
creationDate *time.Time
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)
}
doc := &Record{
UserID: userID,
Login: login,
Password: password,
CreationDate: creationDate,
ParentID: parentID,
Status: status,
}
result = append(result, appsync.SyncRow{ID: userID, Doc: doc})
}
if err := rows.Err(); err != nil {
return nil, fmt.Errorf("iterate user rows: %w", err)
}
return result, 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