101 lines
2.9 KiB
Go
101 lines
2.9 KiB
Go
package cif
|
|
|
|
import (
|
|
"time"
|
|
|
|
"git.fjla.uk/owlboard/go-types/pkg/database"
|
|
"git.fjla.uk/owlboard/go-types/pkg/upstreamApi"
|
|
"git.fjla.uk/owlboard/timetable-mgr/dbAccess"
|
|
"git.fjla.uk/owlboard/timetable-mgr/log"
|
|
"go.uber.org/zap"
|
|
)
|
|
|
|
// Processes parsed CIF data and applies the data to the database
|
|
func processParsedCif(data *parsedData) error {
|
|
createTasks := make([]*upstreamApi.JsonScheduleV1, 0)
|
|
deleteTasks := make([]*upstreamApi.JsonScheduleV1, 0)
|
|
|
|
for _, item := range data.sched {
|
|
switch item.TransactionType {
|
|
case "Delete":
|
|
deleteTasks = append(deleteTasks, &item)
|
|
case "Create":
|
|
createTasks = append(createTasks, &item)
|
|
default:
|
|
log.Msg.Error("Unknown transaction type in CIF Schedule", zap.String("TransactionType", item.TransactionType))
|
|
}
|
|
}
|
|
|
|
err := doDeletions(deleteTasks)
|
|
if err != nil {
|
|
log.Msg.Error("Error deleting CIF Entries", zap.Error(err))
|
|
return err
|
|
}
|
|
err = doCreations(createTasks)
|
|
if err != nil {
|
|
log.Msg.Error("Error creating CIF Entries", zap.Error(err))
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// Create delete query types and pass to the function which batches the deletions
|
|
func doDeletions(deletions []*upstreamApi.JsonScheduleV1) error {
|
|
log.Msg.Info("Preparing CIF update Delete tasks", zap.Int("Delete task count", len(deletions)))
|
|
|
|
deleteQueries := make([]database.DeleteQuery, 0)
|
|
for _, item := range deletions {
|
|
query := database.DeleteQuery{
|
|
ScheduleStartDate: ParseCifDate(&item.ScheduleStartDate, "start"),
|
|
StpIndicator: item.CifStpIndicator,
|
|
TrainUid: item.CifTrainUid,
|
|
}
|
|
|
|
deleteQueries = append(deleteQueries, query)
|
|
}
|
|
|
|
err := dbAccess.DeleteCifEntries(deleteQueries)
|
|
if err != nil {
|
|
log.Msg.Error("Error deleting documents", zap.Error(err))
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// Convert to the correct struct for the database and pass to the function which batches insertions
|
|
func doCreations(creations []*upstreamApi.JsonScheduleV1) error {
|
|
log.Msg.Info("Preparing CIF update Create tasks", zap.Int("Create task count", len(creations)))
|
|
|
|
createDocuments := make([]database.Service, 0)
|
|
for _, item := range creations {
|
|
document, err := ConvertServiceType(item, false)
|
|
if err != nil {
|
|
log.Msg.Error("Error converting JsonSchedule to Service type", zap.Error(err))
|
|
}
|
|
|
|
createDocuments = append(createDocuments, *document)
|
|
}
|
|
|
|
err := dbAccess.CreateCifEntries(createDocuments)
|
|
if err != nil {
|
|
log.Msg.Error("Error creating documents", zap.Error(err))
|
|
return err
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// Accepts the JsonTimetableV1 struct which contains CIF File metadata and produces a DB Ready struct.
|
|
func generateMetadata(header *upstreamApi.JsonTimetableV1) *dbAccess.CifMetadata {
|
|
newMetadata := dbAccess.CifMetadata{
|
|
Doctype: dbAccess.Doctype,
|
|
LastTimestamp: header.Timestamp,
|
|
LastUpdate: time.Now().In(londonTimezone),
|
|
LastSequence: header.Metadata.Sequence,
|
|
}
|
|
|
|
return &newMetadata
|
|
}
|