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

}