All checks were successful
ci/woodpecker/push/woodpecker Pipeline was successful
407 lines
10 KiB
Go
407 lines
10 KiB
Go
package importer
|
|
|
|
import (
|
|
"bytes"
|
|
csvlib "cargo-erp-backend/pkg/csv"
|
|
"cargo-erp-backend/pkg/helpers"
|
|
"cargo-erp-backend/pkg/jobqueue"
|
|
"context"
|
|
"encoding/csv"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"net/http"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
type RowData interface {
|
|
TableName() string
|
|
}
|
|
|
|
type ImportHandler interface {
|
|
ValidateRow(row []string, rowNumber int) (RowData, []error)
|
|
IsDuplicate(db *gorm.DB, data RowData) bool
|
|
GetDataString(data RowData) string
|
|
GetTemplateHeaders() []string
|
|
GetTemplateRows() [][]string
|
|
}
|
|
|
|
type ImportLog struct {
|
|
Row int `json:"row"`
|
|
Status string `json:"status"`
|
|
Data string `json:"data"`
|
|
Message string `json:"message"`
|
|
}
|
|
|
|
type ImportResult struct {
|
|
JobID string `json:"job_id,omitempty"`
|
|
TotalRows int `json:"total_rows"`
|
|
Imported int `json:"imported"`
|
|
Skipped int `json:"skipped"`
|
|
Logs []ImportLog `json:"logs"`
|
|
Errors []string `json:"errors"`
|
|
IsAsync bool `json:"is_async"`
|
|
Message string `json:"message,omitempty"`
|
|
}
|
|
|
|
type Importer struct {
|
|
Helper helpers.HelperInterface
|
|
Handler ImportHandler
|
|
JobQ *jobqueue.JobQueue
|
|
}
|
|
|
|
func NewImporter(helper helpers.HelperInterface, handler ImportHandler) *Importer {
|
|
return &Importer{
|
|
Helper: helper,
|
|
Handler: handler,
|
|
JobQ: jobqueue.NewJobQueue(helper.GetRedis("master")),
|
|
}
|
|
}
|
|
|
|
func (i *Importer) Process(file io.Reader, userID string, filename string) (*ImportResult, error) {
|
|
if !strings.HasSuffix(strings.ToLower(filename), ".csv") {
|
|
return nil, fmt.Errorf("file must be .csv format")
|
|
}
|
|
|
|
fileBytes, err := io.ReadAll(file)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to read file: %w", err)
|
|
}
|
|
|
|
if err := i.validateHeaders(bytes.NewReader(fileBytes)); err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
rowCount, err := csvlib.CountRows(bytes.NewReader(fileBytes))
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to count rows: %w", err)
|
|
}
|
|
|
|
if rowCount <= 1 {
|
|
return &ImportResult{
|
|
TotalRows: 0,
|
|
Imported: 0,
|
|
Skipped: 0,
|
|
Logs: []ImportLog{},
|
|
Errors: []string{"CSV file is empty or has only headers"},
|
|
}, nil
|
|
}
|
|
|
|
totalDataRows := rowCount - 1
|
|
|
|
if totalDataRows >= 10000 {
|
|
return i.processAsync(bytes.NewReader(fileBytes), userID, totalDataRows)
|
|
}
|
|
|
|
return i.processSync(bytes.NewReader(fileBytes), userID, totalDataRows)
|
|
}
|
|
|
|
func (i *Importer) validateHeaders(file io.Reader) error {
|
|
reader := csv.NewReader(file)
|
|
header, err := reader.Read()
|
|
if err != nil {
|
|
return fmt.Errorf("failed to read CSV headers: %w", err)
|
|
}
|
|
|
|
expectedHeaders := i.Handler.GetTemplateHeaders()
|
|
|
|
if len(header) != len(expectedHeaders) {
|
|
return fmt.Errorf("invalid CSV headers. Expected: %s. Got: %s",
|
|
strings.Join(expectedHeaders, ", "),
|
|
strings.Join(header, ", "))
|
|
}
|
|
|
|
for idx, h := range header {
|
|
if strings.TrimSpace(strings.ToLower(h)) != strings.TrimSpace(strings.ToLower(expectedHeaders[idx])) {
|
|
return fmt.Errorf("invalid CSV headers. Expected: %s. Got: %s",
|
|
strings.Join(expectedHeaders, ", "),
|
|
strings.Join(header, ", "))
|
|
}
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
func (i *Importer) processSync(file io.Reader, userID string, totalRows int) (*ImportResult, error) {
|
|
db := i.Helper.GetDB("master")
|
|
|
|
var imported, skipped int
|
|
var logs []ImportLog
|
|
var errors []string
|
|
var batch []RowData
|
|
var batchRows []int
|
|
batchSize := 500
|
|
|
|
csvlib.ReadCSV(file, func(row []string, rowNumber int) error {
|
|
data, rowErrors := i.Handler.ValidateRow(row, rowNumber)
|
|
if len(rowErrors) > 0 {
|
|
for _, e := range rowErrors {
|
|
errorMsg := fmt.Sprintf("Row %d: %s", rowNumber, e.Error())
|
|
errors = append(errors, errorMsg)
|
|
logs = append(logs, ImportLog{
|
|
Row: rowNumber,
|
|
Status: "error",
|
|
Data: strings.Join(row, ","),
|
|
Message: errorMsg,
|
|
})
|
|
}
|
|
skipped++
|
|
return nil
|
|
}
|
|
|
|
if i.Handler.IsDuplicate(db, data) {
|
|
dataStr := i.Handler.GetDataString(data)
|
|
logs = append(logs, ImportLog{
|
|
Row: rowNumber,
|
|
Status: "skipped",
|
|
Data: dataStr,
|
|
Message: fmt.Sprintf("Row %d: '%s' already exists, skipped", rowNumber, dataStr),
|
|
})
|
|
skipped++
|
|
return nil
|
|
}
|
|
|
|
batch = append(batch, data)
|
|
batchRows = append(batchRows, rowNumber)
|
|
if len(batch) >= batchSize {
|
|
if err := i.insertBatch(db, batch); err != nil {
|
|
errorMsg := fmt.Sprintf("Batch insert error: %v", err)
|
|
errors = append(errors, errorMsg)
|
|
logs = append(logs, ImportLog{
|
|
Row: rowNumber,
|
|
Status: "error",
|
|
Data: strings.Join(row, ","),
|
|
Message: errorMsg,
|
|
})
|
|
} else {
|
|
for idx, b := range batch {
|
|
dataStr := i.Handler.GetDataString(b)
|
|
logs = append(logs, ImportLog{
|
|
Row: batchRows[idx],
|
|
Status: "inserted",
|
|
Data: dataStr,
|
|
Message: fmt.Sprintf("Row %d: '%s' inserted successfully", batchRows[idx], dataStr),
|
|
})
|
|
}
|
|
imported += len(batch)
|
|
}
|
|
batch = nil
|
|
batchRows = nil
|
|
}
|
|
return nil
|
|
})
|
|
|
|
if len(batch) > 0 {
|
|
if err := i.insertBatch(db, batch); err != nil {
|
|
errorMsg := fmt.Sprintf("Batch insert error: %v", err)
|
|
errors = append(errors, errorMsg)
|
|
logs = append(logs, ImportLog{
|
|
Row: batchRows[0],
|
|
Status: "error",
|
|
Data: "",
|
|
Message: errorMsg,
|
|
})
|
|
} else {
|
|
for idx, b := range batch {
|
|
dataStr := i.Handler.GetDataString(b)
|
|
logs = append(logs, ImportLog{
|
|
Row: batchRows[idx],
|
|
Status: "inserted",
|
|
Data: dataStr,
|
|
Message: fmt.Sprintf("Row %d: '%s' inserted successfully", batchRows[idx], dataStr),
|
|
})
|
|
}
|
|
imported += len(batch)
|
|
}
|
|
}
|
|
|
|
return &ImportResult{
|
|
TotalRows: totalRows,
|
|
Imported: imported,
|
|
Skipped: skipped,
|
|
Logs: logs,
|
|
Errors: errors,
|
|
IsAsync: false,
|
|
}, nil
|
|
}
|
|
|
|
func (i *Importer) processAsync(file io.Reader, userID string, totalRows int) (*ImportResult, error) {
|
|
jobID := uuid.New().String()
|
|
ctx := context.Background()
|
|
|
|
metadata := map[string]string{
|
|
"user_id": userID,
|
|
"table": i.Handler.GetTemplateHeaders()[0],
|
|
}
|
|
|
|
if err := i.JobQ.CreateJob(ctx, jobID, totalRows, metadata); err != nil {
|
|
return nil, fmt.Errorf("failed to create job: %w", err)
|
|
}
|
|
|
|
go i.processAsyncWorker(file, jobID, userID, totalRows)
|
|
|
|
return &ImportResult{
|
|
JobID: jobID,
|
|
TotalRows: totalRows,
|
|
IsAsync: true,
|
|
Message: fmt.Sprintf("Import started. Poll /import/status/%s for progress.", jobID),
|
|
}, nil
|
|
}
|
|
|
|
func (i *Importer) processAsyncWorker(file io.Reader, jobID string, userID string, totalRows int) {
|
|
db := i.Helper.GetDB("master")
|
|
ctx := context.Background()
|
|
|
|
var imported, skipped int
|
|
var logs []ImportLog
|
|
var errors []string
|
|
var batch []RowData
|
|
var batchRows []int
|
|
batchSize := 500
|
|
|
|
i.JobQ.UpdateProgress(ctx, jobID, 0, 0, nil)
|
|
|
|
csvlib.ReadCSV(file, func(row []string, rowNumber int) error {
|
|
data, rowErrors := i.Handler.ValidateRow(row, rowNumber)
|
|
if len(rowErrors) > 0 {
|
|
for _, e := range rowErrors {
|
|
errorMsg := fmt.Sprintf("Row %d: %s", rowNumber, e.Error())
|
|
errors = append(errors, errorMsg)
|
|
logs = append(logs, ImportLog{
|
|
Row: rowNumber,
|
|
Status: "error",
|
|
Data: strings.Join(row, ","),
|
|
Message: errorMsg,
|
|
})
|
|
}
|
|
skipped++
|
|
return nil
|
|
}
|
|
|
|
if i.Handler.IsDuplicate(db, data) {
|
|
dataStr := i.Handler.GetDataString(data)
|
|
logs = append(logs, ImportLog{
|
|
Row: rowNumber,
|
|
Status: "skipped",
|
|
Data: dataStr,
|
|
Message: fmt.Sprintf("Row %d: '%s' already exists, skipped", rowNumber, dataStr),
|
|
})
|
|
skipped++
|
|
return nil
|
|
}
|
|
|
|
batch = append(batch, data)
|
|
batchRows = append(batchRows, rowNumber)
|
|
if len(batch) >= batchSize {
|
|
if err := i.insertBatch(db, batch); err != nil {
|
|
errorMsg := fmt.Sprintf("Batch insert error: %v", err)
|
|
errors = append(errors, errorMsg)
|
|
logs = append(logs, ImportLog{
|
|
Row: rowNumber,
|
|
Status: "error",
|
|
Data: strings.Join(row, ","),
|
|
Message: errorMsg,
|
|
})
|
|
} else {
|
|
for idx, b := range batch {
|
|
dataStr := i.Handler.GetDataString(b)
|
|
logs = append(logs, ImportLog{
|
|
Row: batchRows[idx],
|
|
Status: "inserted",
|
|
Data: dataStr,
|
|
Message: fmt.Sprintf("Row %d: '%s' inserted successfully", batchRows[idx], dataStr),
|
|
})
|
|
}
|
|
imported += len(batch)
|
|
}
|
|
batch = nil
|
|
batchRows = nil
|
|
i.JobQ.UpdateProgress(ctx, jobID, imported, skipped, errors)
|
|
i.saveJobLogs(ctx, jobID, logs, imported, skipped, errors)
|
|
}
|
|
return nil
|
|
})
|
|
|
|
if len(batch) > 0 {
|
|
if err := i.insertBatch(db, batch); err != nil {
|
|
errorMsg := fmt.Sprintf("Batch insert error: %v", err)
|
|
errors = append(errors, errorMsg)
|
|
logs = append(logs, ImportLog{
|
|
Row: batchRows[0],
|
|
Status: "error",
|
|
Data: "",
|
|
Message: errorMsg,
|
|
})
|
|
} else {
|
|
for idx, b := range batch {
|
|
dataStr := i.Handler.GetDataString(b)
|
|
logs = append(logs, ImportLog{
|
|
Row: batchRows[idx],
|
|
Status: "inserted",
|
|
Data: dataStr,
|
|
Message: fmt.Sprintf("Row %d: '%s' inserted successfully", batchRows[idx], dataStr),
|
|
})
|
|
}
|
|
imported += len(batch)
|
|
}
|
|
}
|
|
|
|
i.JobQ.CompleteJob(ctx, jobID, imported, skipped, errors)
|
|
i.saveJobLogs(ctx, jobID, logs, imported, skipped, errors)
|
|
}
|
|
|
|
func (i *Importer) saveJobLogs(ctx context.Context, jobID string, logs []ImportLog, imported int, skipped int, errors []string) {
|
|
logsJSON, _ := json.Marshal(logs)
|
|
key := fmt.Sprintf("import_job:%s:logs", jobID)
|
|
i.Helper.GetRedis("master").Set(ctx, key, logsJSON, 24*time.Hour)
|
|
}
|
|
|
|
func (i *Importer) insertBatch(db *gorm.DB, batch []RowData) error {
|
|
if len(batch) == 0 {
|
|
return nil
|
|
}
|
|
tableName := batch[0].TableName()
|
|
for _, item := range batch {
|
|
if err := db.Table(tableName).Create(item).Error; err != nil {
|
|
return err
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func (i *Importer) GetTemplate() string {
|
|
headers := i.Handler.GetTemplateHeaders()
|
|
rows := i.Handler.GetTemplateRows()
|
|
return csvlib.WriteCSV(headers, rows)
|
|
}
|
|
|
|
func (i *Importer) ServeTemplate(w http.ResponseWriter) {
|
|
template := i.GetTemplate()
|
|
|
|
w.Header().Set("Content-Type", "text/csv")
|
|
w.Header().Set("Content-Disposition", fmt.Sprintf("attachment; filename=%s_template.csv", time.Now().Format("20060102")))
|
|
w.Header().Set("Content-Length", fmt.Sprintf("%d", len(template)))
|
|
w.Write([]byte(template))
|
|
}
|
|
|
|
func (i *Importer) GetJobStatus(jobID string) (*jobqueue.Job, []ImportLog, error) {
|
|
ctx := context.Background()
|
|
job, err := i.JobQ.GetJob(ctx, jobID)
|
|
if err != nil {
|
|
return nil, nil, err
|
|
}
|
|
|
|
var logs []ImportLog
|
|
key := fmt.Sprintf("import_job:%s:logs", jobID)
|
|
logsJSON, err := i.Helper.GetRedis("slave").Get(ctx, key).Result()
|
|
if err == nil {
|
|
json.Unmarshal([]byte(logsJSON), &logs)
|
|
}
|
|
|
|
return job, logs, nil
|
|
}
|