dsci-runner.git | job/ | job.go
package job
import (
"database/sql"
"dsci_runner/types"
"dsci_runner/utils"
"encoding/json"
"errors"
"fmt"
_ "github.com/mattn/go-sqlite3"
"log"
"os"
"path/filepath"
"strings"
"io"
)
func JobQueueFs(r types.JobRequest, cr string) {
project := r.Config.Project
job_id := r.Config.JobId
fmt.Printf("start job for project: %s, job_id: %s\n", project, job_id)
jsonData, err := json.MarshalIndent(r.Config, "", " ") // Use MarshalIndent for pretty printing
//log.Printf("JobQueueFs queue job: r.Config: %s\n",jsonData)
sparky_project_dir := utils.CreateSparkyProjectDir(project)
cache_dir := utils.CreateSparkyCacheDir(job_id)
jsonData, err = json.MarshalIndent(r.SparrowdoConfig, "", " ") // Use MarshalIndent for pretty printing
if err != nil {
log.Fatal("Error marshaling to JSON:", err)
}
err = os.WriteFile(fmt.Sprintf("%s/config.json", cache_dir), jsonData, 0644)
if err != nil {
log.Fatal(err)
}
r.Trigger.Cwd = utils.SparkyCacheDirDocker(r.Config.JobId, cr)
if r.Config.Sparrowdo.Localhost {
r.Trigger.Sparrowdo.Docker = ""
r.Trigger.Sparrowdo.Host = ""
r.Trigger.Sparrowdo.Localhost = true
} else if r.Config.Sparrowdo.Host != "" {
r.Trigger.Sparrowdo.Docker = ""
r.Trigger.Sparrowdo.Localhost = false
r.Trigger.Sparrowdo.Host = r.Config.Sparrowdo.Host
} else if r.Config.Sparrowdo.Docker != "" {
r.Trigger.Sparrowdo.Localhost = false
r.Trigger.Sparrowdo.Host = ""
r.Trigger.Sparrowdo.Localhost = false
r.Trigger.Sparrowdo.Docker = r.Config.Sparrowdo.Docker
}
if r.Config.Sparrowdo.Sudo {
r.Trigger.Sparrowdo.NoSudo = false
r.Trigger.Sparrowdo.Sudo = true
}
if r.Config.Sparrowdo.NoSudo {
r.Trigger.Sparrowdo.Sudo = false
r.Trigger.Sparrowdo.NoSudo = true
}
if r.Config.Description != "" {
r.Trigger.Description = r.Config.Description
} else {
r.Trigger.Description = "spawned job"
}
if r.Config.Sparrowdo.SshUser != "" {
r.Trigger.Sparrowdo.SshUser = r.Config.Sparrowdo.SshUser
}
if r.Config.Sparrowdo.Image != "" {
r.Trigger.Sparrowdo.Image = r.Config.Sparrowdo.Image
}
if r.Config.Sparrowdo.Conf != "" {
r.Trigger.Sparrowdo.Conf = r.Config.Sparrowdo.Conf
}
if r.Config.Sparrowdo.Bootstrap == true {
r.Trigger.Sparrowdo.Bootstrap = r.Config.Sparrowdo.Bootstrap
}
var tags []string
for k, vv := range r.Config.Tags {
switch v := vv.(type) {
case string:
v_safe := strings.ReplaceAll(v, ",", "___comma___")
v_safe = strings.ReplaceAll(v_safe, "=", "___eq___")
tags = append(tags, fmt.Sprintf("%s=%s", k, v_safe))
case int:
tags = append(tags, fmt.Sprintf("%s=%d", k, v))
case bool:
tags = append(tags, fmt.Sprintf("%s", k))
}
}
if len(tags) != 0 {
fmt.Printf(
"JobQueueFs: append tags from Config.Tags to Trigger.Sparrowdo.Tags string: %s\n",
strings.Join(tags, ","),
)
r.Trigger.Sparrowdo.Tags = r.Trigger.Sparrowdo.Tags + "," + strings.Join(tags, ",")
}
if r.SparrowdoConfig != nil {
r.Trigger.Sparrowdo.Conf = "config.json"
}
fmt.Printf(
"job-queue-fs: create trigger file: %s/.triggers/%s\n",
sparky_project_dir,
job_id,
)
utils.CreateSparkyTriggersDir(r.Config.Project)
jsonData, err = json.MarshalIndent(r.Trigger, "", " ") // Use MarshalIndent for pretty printing
err = os.WriteFile(fmt.Sprintf("%s/.triggers/%s", sparky_project_dir, job_id), jsonData, 0644)
if err != nil {
log.Fatalf("JobQueueFs: can't create trigger file: %s", err)
}
path := fmt.Sprintf("%s/sparrowfile", sparky_project_dir)
if _, err := os.Stat(path); errors.Is(err, os.ErrNotExist) {
err = os.WriteFile(path, []byte("# dummy file, generated by dsci"), 0644)
if err != nil {
log.Fatalf("JobQueueFs: can't create sparrowfile: %s", err)
}
}
err = os.WriteFile(fmt.Sprintf("%s/sparrowfile", cache_dir), []byte(r.Sparrowfile), 0644)
if err != nil {
log.Fatal(err)
}
log.Printf("JobQueueFs queue job: %s\ncache dir: %s\n", jsonData, cache_dir)
}
func PutJobStash(p string, job_id string, data interface{}) {
utils.CreateSparkyStashDir(p)
path := filepath.Join(utils.SparkyStashDir(p), job_id)
file, err := os.Create(path)
if err != nil {
log.Fatalf("PutJobStash: Error creating file:", err)
}
defer file.Close()
encoder := json.NewEncoder(file)
encoder.SetIndent("", " ") // Optional: prettify
err = encoder.Encode(data)
if err != nil {
log.Fatalf("Error encoding JSON:", err)
}
}
func GetJobStash(p string, job_id string) string {
path := filepath.Join(utils.SparkyStashDir(p), job_id)
if _, err := os.Stat(path); errors.Is(err, os.ErrNotExist) {
return "{}"
}
dat, err := os.ReadFile(path)
if err != nil {
log.Fatalf("Error reading file:", err)
}
return string(dat)
}
func PutJobFile(p string, job_id string, filename string, data io.Reader) (int64, error) {
utils.CreateSparkyFilesDir(p)
path := filepath.Join(utils.SparkyFilesDir(p), job_id, filename)
log.Printf("PutJobFile: saving file: %s", path)
dir := filepath.Join(utils.SparkyFilesDir(p), job_id)
err := os.MkdirAll(dir, 0755)
if err != nil {
log.Printf("PutJobFile: error creating directory: %s", err)
return 0, err
}
file, err := os.Create(path)
if err != nil {
log.Printf("PutJobFile: Error creating file:", err)
return 0, err
}
defer file.Close()
bytesWritten, err := io.Copy(file, data)
if err != nil {
log.Printf("PutJobFile: error occured during write to blob file: %s")
return 0, err
}
return bytesWritten, nil
}
func GetJobFile(p string, job_id string, filename string) ([]byte, error) {
path := filepath.Join(utils.SparkyFilesDir(p), job_id, filename)
data, err := os.ReadFile(path)
if err != nil {
log.Printf("GetJobFile error reading file: %s", err)
return []byte{}, err
}
return data, nil
}
func GetSparkyScenarioFile(p string, filename string) string {
path := filepath.Join(utils.SparkyProjectDir(p), filename)
dat, err := os.ReadFile(path)
if err != nil {
log.Fatalf("GetSparkyScenarioFile: Error reading file:", err)
}
return string(dat)
}
func GetSparrowdoConfig(p string, filename string) interface{} {
path := filepath.Join(utils.SparkyProjectDir(p), filename)
dat, err := os.ReadFile(path)
if err != nil {
log.Fatalf("GetSparrowdoConfig: Error reading file:", err)
}
var c interface{}
err = json.Unmarshal(dat, &c)
if err != nil {
log.Fatalf("GetSparrowdoConfig: error unmarshaling JSON: %v", err)
}
return c
}
func JobState(p string, job_id string) string {
// default state - unknown
state := "-2"
active_state := ""
path := filepath.Join(utils.ProjectStateDir(p), job_id)
_, err := os.Stat(path)
if err == nil {
// job either finished or running
// set active_state
dat, err := os.ReadFile(path)
if err != nil {
log.Fatalf("Error reading file:", err)
}
switch st := string(dat); st {
case "0":
log.Printf("JobState: project: %s job_id: %s, state: running ", p, job_id)
case "1":
log.Printf("JobState: project: %s job_id: %s, state: success ", p, job_id)
case "-1":
log.Printf("JobState: project: %s job_id: %s, state: failed ", p, job_id)
}
// set active state
active_state = string(dat)
}
// try to see if there is job in queue
path = filepath.Join(utils.SparkyTriggersDir(p), job_id)
_, err = os.Stat(path)
if errors.Is(err, os.ErrNotExist) {
// trigger file does not exist
log.Printf("JobState: project: %s job_id: %s, trigger file does not exist", p, job_id)
if active_state != "" {
return active_state
} else {
// return default state (unknown) if trigger file
// does not exist and there is no active_state
return state
}
} else if err == nil {
log.Printf("JobState: project: %s job_id: %s, state: in queue", p, job_id)
if active_state == "0" {
return active_state
} else {
// return state "in queue"
// if trigger file exists
// and there active state is finshed
return "-3"
}
} else {
log.Printf("JobState: error accessing trigger file: %s", err)
if active_state != "" {
return active_state
} else {
// return default state (unknown) if
// there is error accessing trigger file
// and there is no active_state
return state
}
}
}
func ReportByBuildId(p string, build_id string) string {
log.Printf("Report request project: %s, build id: %s", p, build_id)
path := fmt.Sprintf(
"%s/%s/build-%s.txt",
utils.SparkyReportsDir(),
p,
build_id,
)
dat, err := os.ReadFile(path)
if err != nil {
log.Fatalf("JobReport: error reading file %s %s:", path, err)
}
log.Printf("Report return: path: %s", path)
return string(dat)
}
func Report(p string, job_id string) string {
log.Printf("Report request project: %s, job_id: %s", p, job_id)
db, err := sql.Open("sqlite3", utils.SparkyDbFile())
defer db.Close()
if err != nil {
log.Fatalf("JobReport: error opening db file: %s", err)
}
sqlQuery := `SELECT id FROM builds WHERE job_id = ? order by id desc LIMIT 1`
// QueryRow returns a *sql.Row
row := db.QueryRow(sqlQuery, job_id)
r := struct {
ID int
}{}
err = row.Scan(&r.ID)
if err != nil {
log.Printf("Report return emtpy, database error: %s", err)
return ""
}
path := fmt.Sprintf(
"%s/%s/build-%d.txt",
utils.SparkyReportsDir(),
p,
r.ID,
)
dat, err := os.ReadFile(path)
if err != nil {
log.Fatalf("JobReport: error reading file %s %s:", path, err)
}
log.Printf("Report return: path: %s", path)
return string(dat)
}
func JobTriggerFile(p string, job_id string) string {
log.Printf("JobTriggerFile. look up trigger: %s %s\n", p, job_id)
hdir, _ := os.UserHomeDir()
path := fmt.Sprintf("%s/.dsci/.sparky/projects/%s/.triggers/%s", hdir, p, job_id)
_, err := os.Stat(path)
if errors.Is(err, os.ErrNotExist) {
path = fmt.Sprintf("%s/.dsci/.sparky/work/%s/.triggers/%s", hdir, p, job_id)
_, err = os.Stat(path)
if errors.Is(err, os.ErrNotExist) {
} else if err != nil {
log.Fatalf("JobTriggerFile: can't read trigger: %s | %s", path, err)
} else {
dat, err := os.ReadFile(path)
if err != nil {
log.Fatalf("JobTriggerFile. Error reading file: %s | %s", path, err)
}
return string(dat)
}
} else if err != nil {
log.Fatalf("JobTriggerFile: can't read trigger: %s | %s", path, err)
} else {
dat, err := os.ReadFile(path)
if err != nil {
log.Fatalf("JobTriggerFile. Error reading file: %s | %s", path, err)
}
return string(dat)
}
log.Printf("JobTriggerFile. trigger file not found\n")
return ""
}
func Builds(db *sql.DB) []types.JobBuild {
q := `SELECT id, project, job_id, description, dt, state FROM builds order by id desc LIMIT 30`
rows, err := db.Query(q)
if err != nil {
log.Fatalf("Builds: database select error: %s", err)
}
var builds []types.JobBuild
for rows.Next() {
var b types.JobBuild
err = rows.Scan(&b.ID, &b.Project, &b.JobId, &b.Description, &b.Dt, &b.State)
if err != nil {
log.Fatalf("Builds, rows.Scan error: ", err)
}
builds = append(builds, b)
}
err = rows.Err()
if err != nil {
log.Fatalf("Builds rows.Err: %s", err)
}
return builds
}