Files
control/control_case/services.go
T
2026-07-12 20:42:21 +05:00

136 lines
2.9 KiB
Go

package control_case
import (
"context"
"errors"
"fmt"
"io/fs"
"os"
"path/filepath"
"sync"
"time"
"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
}
const BASE_DIR = "./data"
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(BASE_DIR, 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
}
}
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, &params); 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").
Preload("Params.InitialCondition").
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)
// }
}