140 lines
3.1 KiB
Go
140 lines
3.1 KiB
Go
package control_case
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"fmt"
|
|
"io/fs"
|
|
"os"
|
|
"path/filepath"
|
|
"sync"
|
|
"time"
|
|
|
|
"control/analize"
|
|
"github.com/gocarina/gocsv"
|
|
"gorm.io/gorm"
|
|
)
|
|
|
|
func exists(path string) (bool, error) {
|
|
_, err := os.Stat(path)
|
|
if err == nil {
|
|
return true, nil
|
|
}
|
|
if errors.Is(err, fs.ErrNotExist) {
|
|
return false, nil
|
|
}
|
|
return false, err
|
|
}
|
|
|
|
func worker(id int, jobs <-chan ControlCase, ctx context.Context, db *gorm.DB) {
|
|
for j := range jobs {
|
|
fmt.Printf("Worker %d started job %d \n", id, j.ID)
|
|
p := j.Params
|
|
p.FolderPath = filepath.Join(".", fmt.Sprintf("%d", p.ID))
|
|
|
|
exists, err := p.Exist()
|
|
if err != nil {
|
|
fmt.Printf("Worker %d failed to check folder for job %d: %v\n", id, j.ID, err)
|
|
continue
|
|
}
|
|
|
|
if !exists {
|
|
if err := p.CreateFolder(); err != nil {
|
|
fmt.Printf("Worker %d failed to prepare job %d: %v\n", id, j.ID, err)
|
|
continue
|
|
}
|
|
if err := p.Run(ctx); err != nil {
|
|
fmt.Printf("Worker %d failed to run job %d: %v\n", id, j.ID, err)
|
|
continue
|
|
}
|
|
}
|
|
|
|
csvPath := filepath.Join(p.FolderPath, "foo.csv")
|
|
if _, err := analize.AnalyzeCSV(db, j.ID, j.Name, csvPath); err != nil {
|
|
fmt.Printf("Worker %d failed to analyze job %d: %v\n", id, j.ID, err)
|
|
continue
|
|
}
|
|
|
|
j.Status = "R"
|
|
if err := db.Model(&ControlCase{}).Where("id = ?", j.ID).Update("status", j.Status).Error; err != nil {
|
|
fmt.Printf("Worker %d failed to update status for job %d: %v\n", id, j.ID, err)
|
|
continue
|
|
}
|
|
time.Sleep(time.Second) // Имитация длительной задачи
|
|
fmt.Printf("Worker %d finished job %d \n", id, j.ID)
|
|
// results <- j * 2
|
|
}
|
|
}
|
|
|
|
func readCSVWithGocsv() []Params {
|
|
file, err := os.Open("db.csv")
|
|
if err != nil {
|
|
fmt.Printf("failed to open csv: %v\n", err)
|
|
return nil
|
|
}
|
|
defer file.Close()
|
|
|
|
var params []Params
|
|
if err := gocsv.UnmarshalFile(file, ¶ms); err != nil {
|
|
fmt.Printf("failed to parse csv: %v\n", err)
|
|
return nil
|
|
}
|
|
|
|
for _, person := range params {
|
|
fmt.Printf("%+v\n", person)
|
|
}
|
|
return params
|
|
}
|
|
|
|
func readFromDb(db *gorm.DB) []ControlCase {
|
|
var cases []ControlCase
|
|
statuses := []string{"N"}
|
|
_ = db.
|
|
Where("status IN ?", statuses).
|
|
Preload("Params").
|
|
Find(&cases).Error
|
|
return cases
|
|
}
|
|
|
|
func RunControllWorker(db *gorm.DB, results chan int) {
|
|
const numJobs = 1
|
|
const numWorkers = 10
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
defer func() {
|
|
fmt.Printf("cencel")
|
|
cancel()
|
|
}() // Cancel if main exits early
|
|
|
|
jobs := make(chan ControlCase, numJobs)
|
|
//results := make(chan int, numJobs)
|
|
|
|
// Запуск воркеров
|
|
var wg sync.WaitGroup
|
|
for w := 1; w <= numWorkers; w++ {
|
|
wg.Add(1)
|
|
go func(w int) {
|
|
defer wg.Done()
|
|
worker(w, jobs, ctx, db)
|
|
}(w)
|
|
}
|
|
|
|
//params := readCSVWithGocsv()
|
|
for {
|
|
cases := readFromDb(db)
|
|
// Отправка задач
|
|
for _, c := range cases {
|
|
jobs <- c
|
|
}
|
|
time.Sleep(10 * time.Second)
|
|
}
|
|
// Ожидание завершения всех воркеров
|
|
//fmt.Printf("wait")
|
|
// wg.Wait()
|
|
// close(jobs)
|
|
//
|
|
// // Получение результатов
|
|
// for r := range results {
|
|
// fmt.Println("Result:", r)
|
|
// }
|
|
}
|