cc-backend/pkg/archive/archive.go

152 lines
3.2 KiB
Go
Raw Permalink Normal View History

// Copyright (C) 2022 NHR@FAU, University Erlangen-Nuremberg.
// All rights reserved.
// Use of this source code is governed by a MIT-style
// license that can be found in the LICENSE file.
package archive
import (
"encoding/json"
"fmt"
2023-03-27 13:24:06 +02:00
"github.com/ClusterCockpit/cc-backend/pkg/log"
"github.com/ClusterCockpit/cc-backend/pkg/lrucache"
"github.com/ClusterCockpit/cc-backend/pkg/schema"
)
const Version uint64 = 1
2023-03-27 13:24:06 +02:00
type ArchiveBackend interface {
Init(rawConfig json.RawMessage) (uint64, error)
2023-05-15 14:32:23 +02:00
Info()
Exists(job *schema.Job) bool
LoadJobMeta(job *schema.Job) (*schema.JobMeta, error)
LoadJobData(job *schema.Job) (schema.JobData, error)
LoadClusterCfg(name string) (*schema.Cluster, error)
StoreJobMeta(jobMeta *schema.JobMeta) error
ImportJob(jobMeta *schema.JobMeta, jobData *schema.JobData) error
2022-09-06 08:57:38 +02:00
GetClusters() []string
2023-05-09 16:33:26 +02:00
CleanUp(jobs []*schema.Job)
2023-04-18 07:43:21 +02:00
Move(jobs []*schema.Job, path string)
2023-05-15 14:32:23 +02:00
Clean(before int64, after int64)
2023-05-09 16:33:26 +02:00
Compress(jobs []*schema.Job)
2023-05-09 09:34:03 +02:00
CompressLast(starttime int64) int64
Iter(loadMetricData bool) <-chan JobContainer
}
type JobContainer struct {
Meta *schema.JobMeta
Data *schema.JobData
}
var cache *lrucache.Cache = lrucache.New(128 * 1024 * 1024)
var ar ArchiveBackend
2022-11-08 16:49:45 +01:00
var useArchive bool
2022-11-08 16:49:45 +01:00
func Init(rawConfig json.RawMessage, disableArchive bool) error {
useArchive = !disableArchive
2023-05-09 09:34:03 +02:00
var cfg struct {
2022-09-06 08:57:38 +02:00
Kind string `json:"kind"`
}
2023-05-09 09:34:03 +02:00
if err := json.Unmarshal(rawConfig, &cfg); err != nil {
log.Warn("Error while unmarshaling raw config json")
2022-09-06 08:57:38 +02:00
return err
}
2023-05-09 09:34:03 +02:00
switch cfg.Kind {
2022-09-06 08:57:38 +02:00
case "file":
ar = &FsArchive{}
// case "s3":
// ar = &S3Archive{}
2022-09-06 08:57:38 +02:00
default:
2023-05-09 09:34:03 +02:00
return fmt.Errorf("ARCHIVE/ARCHIVE > unkown archive backend '%s''", cfg.Kind)
2022-09-06 08:57:38 +02:00
}
2023-03-27 13:24:06 +02:00
version, err := ar.Init(rawConfig)
if err != nil {
log.Error("Error while initializing archiveBackend")
2022-09-06 08:57:38 +02:00
return err
}
2023-03-27 13:24:06 +02:00
log.Infof("Load archive version %d", version)
2023-05-09 09:34:03 +02:00
return initClusterConfig()
}
func GetHandle() ArchiveBackend {
return ar
}
// Helper to metricdata.LoadAverages().
func LoadAveragesFromArchive(
job *schema.Job,
metrics []string,
data [][]schema.Float) error {
metaFile, err := ar.LoadJobMeta(job)
if err != nil {
log.Warn("Error while loading job metadata from archiveBackend")
return err
}
for i, m := range metrics {
if stat, ok := metaFile.Statistics[m]; ok {
data[i] = append(data[i], schema.Float(stat.Avg))
} else {
data[i] = append(data[i], schema.NaN)
}
}
return nil
}
func GetStatistics(job *schema.Job) (map[string]schema.JobStatistics, error) {
metaFile, err := ar.LoadJobMeta(job)
if err != nil {
log.Warn("Error while loading job metadata from archiveBackend")
return nil, err
}
return metaFile.Statistics, nil
}
// If the job is archived, find its `meta.json` file and override the tags list
// in that JSON file. If the job is not archived, nothing is done.
func UpdateTags(job *schema.Job, tags []*schema.Tag) error {
2022-11-08 16:49:45 +01:00
if job.State == schema.JobStateRunning || !useArchive {
return nil
}
jobMeta, err := ar.LoadJobMeta(job)
if err != nil {
log.Warn("Error while loading job metadata from archiveBackend")
return err
}
jobMeta.Tags = make([]*schema.Tag, 0)
for _, tag := range tags {
jobMeta.Tags = append(jobMeta.Tags, &schema.Tag{
Name: tag.Name,
Type: tag.Type,
})
}
return ar.StoreJobMeta(jobMeta)
}