diff --git a/CLAUDE.md b/CLAUDE.md index 928681c..b6b4339 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -18,7 +18,7 @@ go test -v ./ccMessage -run TestJSONEncode # Run single test go test -cover ./... # Tests with coverage ``` -No Makefile or special build tools. The `ccTopology` package requires hwloc C library (`sudo apt install hwloc`). +`make fmt` formats all source files with `gofumpt`. The `ccTopology` package requires hwloc C library (`sudo apt install hwloc`). ## Code Style Requirements @@ -37,7 +37,7 @@ lp "github.com/ClusterCockpit/cc-lib/v2/ccMessage" mp "github.com/ClusterCockpit/cc-lib/v2/messageProcessor" ``` -**Formatting:** Code must be formatted with `gofumpt` (a stricter `gofmt`). +**Formatting:** Code must be formatted with `gofumpt` (a stricter `gofmt`). Run `make fmt` to format all source files. **Import grouping:** stdlib first, then third-party separated by blank line. diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..13b9a6a --- /dev/null +++ b/Makefile @@ -0,0 +1,11 @@ +.PHONY: fmt test upgrade + +fmt: + gofumpt -l -w . + +test: + go test ./... + +upgrade: + go get -u ./... + go mod tidy diff --git a/ccTopology/ccTopology.go b/ccTopology/ccTopology.go index 3a44983..f7dffbc 100644 --- a/ccTopology/ccTopology.go +++ b/ccTopology/ccTopology.go @@ -364,7 +364,7 @@ func convertObject(hwloc_obj C.hwloc_obj_t) (Object, bool, error) { value := C._hwloc_get_info_value_by_idx(hwloc_obj, C.int(i)) o.Infos[C.GoString(name)] = C.GoString(value) } - if (hwloc_obj.attr) != nil { + if hwloc_obj.attr != nil { switch hwloc_obj._type { case C.HWLOC_OBJ_NUMANODE: o.Infos["local_memory"] = fmt.Sprintf("%d", uint64(C._hwloc_read_numanode_memory(hwloc_obj))) diff --git a/ccUnits/ccUnits_test.go b/ccUnits/ccUnits_test.go index f3c1b21..5f597d2 100644 --- a/ccUnits/ccUnits_test.go +++ b/ccUnits/ccUnits_test.go @@ -99,7 +99,7 @@ func TestUnitUnitConversion(t *testing.T) { {"Flops/s", NewUnit("GFlops/s"), 1e-9}, {"MHz", NewUnit("Hertz"), 1e6}, {"kb", NewUnit("Kib"), 1000.0 / 1024}, - {"Mib", NewUnit("MBytes"), (1024 * 1024.0) / (1e6)}, + {"Mib", NewUnit("MBytes"), (1024 * 1024.0) / 1e6}, {"mb", NewUnit("MBytes"), 1.0}, } compareUnitWithPrefix := func(in, out Unit, factor float64) bool { diff --git a/go.mod b/go.mod index a92d175..c715ef4 100644 --- a/go.mod +++ b/go.mod @@ -22,23 +22,22 @@ require ( github.com/apapsch/go-jsonmerge/v2 v2.0.0 // indirect github.com/beorn7/perks v1.0.1 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect - github.com/coder/websocket v1.8.14 // indirect + github.com/coder/websocket v1.8.15 // indirect github.com/google/go-tpm v0.9.8 // indirect github.com/google/uuid v1.6.0 // indirect github.com/influxdata/line-protocol v0.0.0-20210922203350-b1ad95c89adf // indirect - github.com/klauspost/compress v1.18.6 // indirect + github.com/klauspost/compress v1.19.0 // indirect github.com/minio/highwayhash v1.0.4 // indirect github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/nats-io/jwt/v2 v2.8.2 // indirect github.com/nats-io/nkeys v0.4.16 // indirect github.com/nats-io/nuid v1.0.1 // indirect - github.com/oapi-codegen/runtime v1.4.0 // indirect + github.com/oapi-codegen/runtime v1.4.2 // indirect github.com/prometheus/client_model v0.6.2 // indirect - github.com/prometheus/common v0.67.5 // indirect - github.com/prometheus/procfs v0.20.1 // indirect - go.yaml.in/yaml/v2 v2.4.4 // indirect + github.com/prometheus/common v0.69.0 // indirect + github.com/prometheus/procfs v0.21.1 // indirect golang.org/x/crypto v0.53.0 // indirect - golang.org/x/net v0.55.0 // indirect + golang.org/x/net v0.56.0 // indirect golang.org/x/sys v0.46.0 // indirect golang.org/x/time v0.15.0 // indirect google.golang.org/protobuf v1.36.11 // indirect diff --git a/go.sum b/go.sum index b7d1aa2..ee3fbcf 100644 --- a/go.sum +++ b/go.sum @@ -22,8 +22,8 @@ github.com/cenkalti/backoff/v4 v4.2.1 h1:y4OZtCnogmCPw98Zjyt5a6+QwPLGkiQsYW5oUqy github.com/cenkalti/backoff/v4 v4.2.1/go.mod h1:Y3VNntkOUPxTVeUxJ/G5vcM//AlwfmyYozVcomhLiZE= github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= -github.com/coder/websocket v1.8.14 h1:9L0p0iKiNOibykf283eHkKUHHrpG7f65OE3BhhO7v9g= -github.com/coder/websocket v1.8.14/go.mod h1:NX3SzP+inril6yawo5CQXx8+fk145lPDC6pumgx0mVg= +github.com/coder/websocket v1.8.15 h1:6B2JPeOGlpff2Uz6vOEH1Vzpi0iUz20A+lPVhPHtNUA= +github.com/coder/websocket v1.8.15/go.mod h1:NX3SzP+inril6yawo5CQXx8+fk145lPDC6pumgx0mVg= github.com/containerd/containerd v1.7.12 h1:+KQsnv4VnzyxWcfO9mlxxELaoztsDEjOuCMPAuPqgU0= github.com/containerd/containerd v1.7.12/go.mod h1:/5OMpE1p0ylxtEUGY8kuCYkDRzJm9NO1TFMWjUpdevk= github.com/containerd/log v0.1.0 h1:TCJt7ioM2cr/tfR8GPbGf9/VRAX8D2B4PjzCpfX540I= @@ -68,8 +68,8 @@ github.com/influxdata/line-protocol v0.0.0-20210922203350-b1ad95c89adf/go.mod h1 github.com/influxdata/line-protocol-corpus v0.0.0-20210922080147-aa28ccfb8937 h1:MHJNQ+p99hFATQm6ORoLmpUCF7ovjwEFshs/NHzAbig= github.com/influxdata/line-protocol-corpus v0.0.0-20210922080147-aa28ccfb8937/go.mod h1:BKR9c0uHSmRgM/se9JhFHtTT7JTO67X23MtKMHtZcpo= github.com/juju/gnuflag v0.0.0-20171113085948-2ce1bb71843d/go.mod h1:2PavIy+JPciBPrBUjwbNvtwB6RQlve+hkpll6QSNmOE= -github.com/klauspost/compress v1.18.6 h1:2jupLlAwFm95+YDR+NwD2MEfFO9d4z4Prjl1XXDjuao= -github.com/klauspost/compress v1.18.6/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2oVvCQ= +github.com/klauspost/compress v1.19.0/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= @@ -102,8 +102,10 @@ github.com/nats-io/nkeys v0.4.16 h1:rd5oAuLOb8mnAycB0xleuEBNS1pVVnN0fv/FF34Eypg= github.com/nats-io/nkeys v0.4.16/go.mod h1:llLgWoI0o4z/Q57q2R1kHfmocyhGV6VG/U18Glg1Afs= github.com/nats-io/nuid v1.0.1 h1:5iA8DT8V7q8WK2EScv2padNa/rTESc1KdnPw4TC2paw= github.com/nats-io/nuid v1.0.1/go.mod h1:19wcPz3Ph3q0Jbyiqsd0kePYG7A95tJPxeL+1OSON2c= -github.com/oapi-codegen/runtime v1.4.0 h1:KLOSFOp7UzkbS7Cs1ms6NBEKYr0WmH2wZG0KKbd2er4= -github.com/oapi-codegen/runtime v1.4.0/go.mod h1:5sw5fxCDmnOzKNYmkVNF8d34kyUeejJEY8HNT2WaPec= +github.com/oapi-codegen/nullable v1.1.0 h1:eAh8JVc5430VtYVnq00Hrbpag9PFRGWLjxR1/3KntMs= +github.com/oapi-codegen/nullable v1.1.0/go.mod h1:KUZ3vUzkmEKY90ksAmit2+5juDIhIZhfDl+0PwOQlFY= +github.com/oapi-codegen/runtime v1.4.2 h1:GMxFVYLzoYLua+/KvzgSphkyK1lLTReQI9Vf4hvATKE= +github.com/oapi-codegen/runtime v1.4.2/go.mod h1:GwV7hC2hviaMzj+ITfHVRESK5J2W/GefVwIND/bMGvU= github.com/opencontainers/go-digest v1.0.0 h1:apOUWs51W5PlhuyGyz9FCeeBIOUDA/6nW8Oi/yOhh5U= github.com/opencontainers/go-digest v1.0.0/go.mod h1:0JzlMkj0TRzQZfJkVvzbP0HBR3IKzErnv2BNG4W4MAM= github.com/opencontainers/image-spec v1.1.0-rc5 h1:Ygwkfw9bpDvs+c9E34SdgGOj41dX/cbdlwvlWt0pnFI= @@ -120,10 +122,10 @@ github.com/prometheus/client_golang v1.23.2 h1:Je96obch5RDVy3FDMndoUsjAhG5Edi49h github.com/prometheus/client_golang v1.23.2/go.mod h1:Tb1a6LWHB3/SPIzCoaDXI4I8UHKeFTEQ1YCr+0Gyqmg= github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= -github.com/prometheus/common v0.67.5 h1:pIgK94WWlQt1WLwAC5j2ynLaBRDiinoAb86HZHTUGI4= -github.com/prometheus/common v0.67.5/go.mod h1:SjE/0MzDEEAyrdr5Gqc6G+sXI67maCxzaT3A2+HqjUw= -github.com/prometheus/procfs v0.20.1 h1:XwbrGOIplXW/AU3YhIhLODXMJYyC1isLFfYCsTEycfc= -github.com/prometheus/procfs v0.20.1/go.mod h1:o9EMBZGRyvDrSPH1RqdxhojkuXstoe4UlK79eF5TGGo= +github.com/prometheus/common v0.69.0 h1:OA85nJQS/T/MaYh/Q2CcgDKSGWqNIgrBDvDH85CuiNk= +github.com/prometheus/common v0.69.0/go.mod h1:ZzL3f6u94qUxh9p+tJTrF+FvBS1XXbbRAZCQkytAL0Y= +github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI= +github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY= github.com/questdb/go-questdb-client/v4 v4.2.0 h1:+d0HJwCjUWMj7zmY6qmhoqTJzTyoYKl+LSTYGN0T8T8= github.com/questdb/go-questdb-client/v4 v4.2.0/go.mod h1:/2x93LK1wjM4JX/b5c6q7Yqk22htjWY1lE6p1X8iLbE= github.com/rogpeppe/go-internal v1.11.0 h1:cWPaGQEPrBb5/AsnsZesgZZ9yb1OQ+GOISoDNXVBh4M= @@ -157,14 +159,14 @@ go.yaml.in/yaml/v2 v2.4.4 h1:tuyd0P+2Ont/d6e2rl3be67goVK4R6deVxCUX5vyPaQ= go.yaml.in/yaml/v2 v2.4.4/go.mod h1:gMZqIpDtDqOfM0uNfy0SkpRhvUryYH0Z6wdMYcacYXQ= golang.org/x/crypto v0.53.0 h1:QZ4Muo8THX6CizN2vPPd5fBGHyogrdK9fG4wLPFUsto= golang.org/x/crypto v0.53.0/go.mod h1:DNLU434OwVakk9PzuwV8w62mAJpRJL3vsgcfp4Qnsio= -golang.org/x/exp v0.0.0-20231005195138-3e424a577f31 h1:9k5exFQKQglLo+RoP+4zMjOFE14P6+vyR0baDAi0Rcs= -golang.org/x/exp v0.0.0-20231005195138-3e424a577f31/go.mod h1:S2oDrQGGwySpoQPVqRShND87VCbxmc6bL1Yd2oYrm6k= +golang.org/x/exp v0.0.0-20240404231335-c0f41cb1a7a0 h1:985EYyeCOxTpcgOTJpflJUwOeEz0CQOdPt73OzpE9F8= +golang.org/x/exp v0.0.0-20240404231335-c0f41cb1a7a0/go.mod h1:/lliqkxwWAhPjf5oSOIJup2XcqJaw8RGS6k3TGEc7GI= golang.org/x/mod v0.33.0 h1:tHFzIWbBifEmbwtGz65eaWyGiGZatSrT9prnU8DbVL8= golang.org/x/mod v0.33.0/go.mod h1:swjeQEj+6r7fODbD2cqrnje9PnziFuw4bmLbBZFrQ5w= -golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8= -golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww= -golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4= -golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/net v0.56.0 h1:Rw8j/hFzGvJUZwNBXnAtf5sVDVt+65SK2C7IxCxZt5o= +golang.org/x/net v0.56.0/go.mod h1:D3Ku6r+V6JROoZK144D2XfMHFcMq/0zSfLelVTCFKec= +golang.org/x/sync v0.21.0 h1:HLII4xRRTtCRkxYp4HNFF0Js/Og6q2i++KXbg0gHCwM= +golang.org/x/sync v0.21.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.21.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.46.0 h1:noSf2Fq6F8DBgS+LysIkx7rIExoNHJsxOAtPp4rthXw= golang.org/x/sys v0.46.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= diff --git a/receivers/ipmiReceiver.go b/receivers/ipmiReceiver.go index bd0f01e..abf16a3 100644 --- a/receivers/ipmiReceiver.go +++ b/receivers/ipmiReceiver.go @@ -65,7 +65,8 @@ func (r *IPMIReceiver) doReadMetric() { clientConfig := &r.config.ClientConfigs[i] var cmd_options []string if clientConfig.Protocol == "ipmi-sensors" { - cmd_options = append(cmd_options, + cmd_options = append( + cmd_options, "--always-prefix", "--sdr-cache-recreate", // Attempt to interpret OEM data, such as event data, sensor readings, or general extra info @@ -139,7 +140,10 @@ func (r *IPMIReceiver) doReadMetric() { name := strings.ToLower( strings.ReplaceAll( strings.TrimSpace( - v2[idxName]), " ", "_")) + v2[idxName], + ), " ", "_", + ), + ) // remove prefix enumeration like 01-... if v := numPrefixRegex.FindStringSubmatch(name); v != nil { name = v[1] @@ -207,7 +211,8 @@ func (r *IPMIReceiver) doReadMetric() { // Debug output for unprocessed metrics fmt.Printf( "host: '%s', metric: '%s', name: '%s', unit: '%s'\n", - host, metric, name, unit) + host, metric, name, unit, + ) } continue } @@ -238,7 +243,8 @@ func (r *IPMIReceiver) doReadMetric() { map[string]any{ "value": value, }, - time.Now()) + time.Now(), + ) if err == nil { r.sink <- y } @@ -527,7 +533,8 @@ func NewIPMIReceiver(name string, config json.RawMessage) (Receiver, error) { Password: password, CLIOptions: cliOptions, isExcluded: isExcluded, - }) + }, + ) } if totalNumHosts == 0 { diff --git a/receivers/redfishReceiver.go b/receivers/redfishReceiver.go index cd586d3..5671574 100644 --- a/receivers/redfishReceiver.go +++ b/receivers/redfishReceiver.go @@ -37,7 +37,7 @@ type RedfishReceiverClientConfig struct { Hostname string // is metric excluded globally or per client - isExcluded map[string](bool) + isExcluded map[string]bool doPowerMetric bool doProcessorMetrics bool @@ -294,7 +294,8 @@ func (r *RedfishReceiver) readSensors( return } }, - clientConfig.readSensorURLs[chassis.ID]) + clientConfig.readSensorURLs[chassis.ID], + ) } return nil } @@ -622,7 +623,8 @@ func (r *RedfishReceiver) readMetrics(clientConfig *RedfishReceiverClientConfig) clientConfig.gofish.BasicAuth, clientConfig.gofish.HTTPClient.Timeout, clientConfig.gofish.HTTPClient.Transport.(*http.Transport).TLSClientConfig.InsecureSkipVerify, - err) + err, + ) } defer c.Logout() @@ -1021,7 +1023,8 @@ func NewRedfishReceiver(name string, config json.RawMessage) (Receiver, error) { HTTPClient: httpClient, }, mp: p, - }) + }, + ) } } diff --git a/schema/job.go b/schema/job.go index e2c09ce..cc5362c 100644 --- a/schema/job.go +++ b/schema/job.go @@ -6,6 +6,7 @@ package schema import ( + "encoding/json" "errors" "fmt" "io" @@ -28,38 +29,38 @@ import ( // Job model // @Description Information of a HPC job. type Job struct { - Cluster string `json:"cluster" db:"cluster" example:"fritz"` - SubCluster string `json:"subCluster" db:"subcluster" example:"main"` - Partition string `json:"partition,omitempty" db:"cluster_partition" example:"main"` - Project string `json:"project" db:"project" example:"abcd200"` - User string `json:"user" db:"hpc_user" example:"abcd100h"` - Shared string `json:"shared" db:"shared" enums:"none,single_user,multi_user"` - State JobState `json:"jobState" db:"job_state" example:"completed" enums:"boot_fail,cancelled,completed,deadline,failed,node_fail,out-of-memory,pending,preempted,running,suspended,timeout"` - Tags []*Tag `json:"tags,omitempty"` - RawEnergyFootprint []byte `json:"-" db:"energy_footprint"` - RawFootprint []byte `json:"-" db:"footprint"` - RawMetaData []byte `json:"-" db:"meta_data"` - RawResources []byte `json:"-" db:"resources"` - Resources []*Resource `json:"resources"` - EnergyFootprint map[string]float64 `json:"energyFootprint"` - Footprint map[string]float64 `json:"footprint"` - MetaData map[string]string `json:"metaData"` - ConcurrentJobs JobLinkResultList `json:"concurrentJobs"` - Energy float64 `json:"energy" db:"energy"` - ArrayJobID int64 `json:"arrayJobId,omitempty" db:"array_job_id" example:"123000"` - Walltime int64 `json:"walltime,omitempty" db:"walltime" example:"86400" minimum:"1"` - RequestedMemory int64 `json:"requestedMemory,omitempty" db:"requested_memory" example:"128000" minimum:"1"` // in MB - JobID int64 `json:"jobId" db:"job_id" example:"123000"` - Duration int32 `json:"duration" db:"duration" example:"43200" minimum:"1"` - SMT int32 `json:"smt,omitempty" db:"smt" example:"4"` - MonitoringStatus int32 `json:"monitoringStatus,omitempty" db:"monitoring_status" example:"1" minimum:"0" maximum:"3"` - NumAcc int32 `json:"numAcc,omitempty" db:"num_acc" example:"2" minimum:"1"` - NumHWThreads int32 `json:"numHwthreads,omitempty" db:"num_hwthreads" example:"20" minimum:"1"` - NumNodes int32 `json:"numNodes" db:"num_nodes" example:"2" minimum:"1"` - Statistics map[string]JobStatistics `json:"statistics"` - ID *int64 `json:"id,omitempty" db:"id"` - SubmitTime int64 `json:"submitTime,omitempty" db:"submit_time" example:"1649723812"` - StartTime int64 `json:"startTime" db:"start_time" example:"1649723812"` + Cluster string `json:"cluster" db:"cluster" example:"fritz"` + SubCluster string `json:"subCluster" db:"subcluster" example:"main"` + Partition string `json:"partition,omitempty" db:"cluster_partition" example:"main"` + Project string `json:"project" db:"project" example:"abcd200"` + User string `json:"user" db:"hpc_user" example:"abcd100h"` + Shared string `json:"shared" db:"shared" enums:"none,single_user,multi_user"` + State JobState `json:"jobState" db:"job_state" example:"completed" enums:"boot_fail,cancelled,completed,deadline,failed,node_fail,out-of-memory,pending,preempted,running,suspended,timeout"` + Tags []*Tag `json:"tags,omitempty"` + RawEnergyFootprint []byte `json:"-" db:"energy_footprint"` + RawFootprint []byte `json:"-" db:"footprint"` + RawMetaData []byte `json:"-" db:"meta_data"` + RawResources []byte `json:"-" db:"resources"` + Resources []*Resource `json:"resources"` + EnergyFootprint map[string]float64 `json:"energyFootprint"` + Footprint map[string]float64 `json:"footprint"` + MetaData map[string]string `json:"metaData"` + ConcurrentJobs JobLinkResultList `json:"concurrentJobs"` + Energy float64 `json:"energy" db:"energy"` + ArrayJobID int64 `json:"arrayJobId,omitempty" db:"array_job_id" example:"123000"` + Walltime int64 `json:"walltime,omitempty" db:"walltime" example:"86400" minimum:"1"` + RequestedMemory int64 `json:"requestedMemory,omitempty" db:"requested_memory" example:"128000" minimum:"1"` // in MB + JobID int64 `json:"jobId" db:"job_id" example:"123000"` + Duration int32 `json:"duration" db:"duration" example:"43200" minimum:"1"` + SMT int32 `json:"smt,omitempty" db:"smt" example:"4"` + MonitoringStatus int32 `json:"monitoringStatus,omitempty" db:"monitoring_status" example:"1" minimum:"0" maximum:"3"` + NumAcc int32 `json:"numAcc,omitempty" db:"num_acc" example:"2" minimum:"1"` + NumHWThreads int32 `json:"numHwthreads,omitempty" db:"num_hwthreads" example:"20" minimum:"1"` + NumNodes int32 `json:"numNodes" db:"num_nodes" example:"2" minimum:"1"` + Statistics JobStatisticsSet `json:"statistics"` + ID *int64 `json:"id,omitempty" db:"id"` + SubmitTime int64 `json:"submitTime,omitempty" db:"submit_time" example:"1649723812"` + StartTime int64 `json:"startTime" db:"start_time" example:"1649723812"` } // JobLink represents a lightweight reference to a job, typically used for linking related jobs. @@ -108,6 +109,151 @@ type JobStatistics struct { Max float64 `json:"max" example:"3000" minimum:"0"` // Job metric maximum } +// StatsGroupInstance is one named element of an array-valued statistics group +// (a single filesystem, later a single interconnect). Its per-instance metrics +// are flat JobStatistics (the job-meta schema does not scope filesystem stats). +type StatsGroupInstance struct { + Name string `json:"name"` + Type string `json:"type"` + Metrics map[string]JobStatistics `json:"-"` // flattened onto the instance object by MarshalJSON +} + +// StatsGroup is an array-valued statistics group identified by its top-level +// JSON key (e.g. "filesystems"). Instance order is preserved. +type StatsGroup struct { + Key string + Instances []StatsGroupInstance +} + +// JobStatisticsSet is the value of Job.Statistics. Metrics holds the flat +// per-metric aggregates; Groups holds array-valued groups (filesystems, and +// later interconnects) whose members each carry their own statistics. Custom +// (Un)MarshalJSON keep the job-meta "statistics" layout: flat metrics as +// objects plus each group as a nested array. +type JobStatisticsSet struct { + Metrics map[string]JobStatistics + Groups []StatsGroup +} + +// MarshalJSON renders the job-meta "statistics" layout: each flat metric as +// {unit,avg,min,max}, and each group as an array of +// {name, type, : {unit,avg,min,max}, ...} instances. +func (s JobStatisticsSet) MarshalJSON() ([]byte, error) { + out := make(map[string]json.RawMessage, len(s.Metrics)+len(s.Groups)) + + for name, stat := range s.Metrics { + b, err := json.Marshal(stat) + if err != nil { + return nil, err + } + out[name] = b + } + + for _, group := range s.Groups { + arr := make([]json.RawMessage, 0, len(group.Instances)) + for _, inst := range group.Instances { + obj := make(map[string]json.RawMessage, len(inst.Metrics)+2) + name, err := json.Marshal(inst.Name) + if err != nil { + return nil, err + } + obj["name"] = name + typ, err := json.Marshal(inst.Type) + if err != nil { + return nil, err + } + obj["type"] = typ + for metric, stat := range inst.Metrics { + b, err := json.Marshal(stat) + if err != nil { + return nil, err + } + obj[metric] = b + } + b, err := json.Marshal(obj) + if err != nil { + return nil, err + } + arr = append(arr, b) + } + b, err := json.Marshal(arr) + if err != nil { + return nil, err + } + out[group.Key] = b + } + + return json.Marshal(out) +} + +// UnmarshalJSON parses the job-meta "statistics" layout into a JobStatisticsSet. +// Array-valued group keys (registered, or fallback: value starts with '[') are +// decoded into Groups; all other keys into flat Metrics. +func (s *JobStatisticsSet) UnmarshalJSON(b []byte) error { + var raw map[string]json.RawMessage + if err := json.Unmarshal(b, &raw); err != nil { + return err + } + + s.Metrics = make(map[string]JobStatistics, len(raw)) + s.Groups = nil + + for key, msg := range raw { + if IsMetricGroupKey(key) || firstToken(msg) == '[' { + group := StatsGroup{Key: key} + var items []map[string]json.RawMessage + if err := json.Unmarshal(msg, &items); err != nil { + return err + } + for _, item := range items { + inst := StatsGroupInstance{Metrics: make(map[string]JobStatistics, len(item))} + for field, fmsg := range item { + switch field { + case "name": + if err := json.Unmarshal(fmsg, &inst.Name); err != nil { + return err + } + case "type": + if err := json.Unmarshal(fmsg, &inst.Type); err != nil { + return err + } + default: + var stat JobStatistics + if err := json.Unmarshal(fmsg, &stat); err != nil { + return err + } + inst.Metrics[field] = stat + } + } + group.Instances = append(group.Instances, inst) + } + s.Groups = append(s.Groups, group) + continue + } + + var stat JobStatistics + if err := json.Unmarshal(msg, &stat); err != nil { + return err + } + s.Metrics[key] = stat + } + + return nil +} + +// AddGroupInstance appends a named statistics instance to the group identified +// by groupKey, creating the group if necessary. +func (s *JobStatisticsSet) AddGroupInstance(groupKey, name, typ string, metrics map[string]JobStatistics) { + inst := StatsGroupInstance{Name: name, Type: typ, Metrics: metrics} + for i := range s.Groups { + if s.Groups[i].Key == groupKey { + s.Groups[i].Instances = append(s.Groups[i].Instances, inst) + return + } + } + s.Groups = append(s.Groups, StatsGroup{Key: groupKey, Instances: []StatsGroupInstance{inst}}) +} + // Tag model // @Description Defines a tag using name and type. type Tag struct { diff --git a/schema/metrics.go b/schema/metrics.go index cc78e8c..85328be 100644 --- a/schema/metrics.go +++ b/schema/metrics.go @@ -6,6 +6,7 @@ package schema import ( + "encoding/json" "fmt" "io" "math" @@ -15,19 +16,352 @@ import ( "github.com/ClusterCockpit/cc-lib/v2/util" ) -// JobData maps metric names to their data organized by scope. -// Structure: map[metricName]map[scope]*JobMetric +// ScopedMetrics maps a hierarchical scope (node/socket/core/...) to its metric +// data. It is the value stored for a single metric name. +type ScopedMetrics map[MetricScope]*JobMetric + +// MetricGroupInstance is one named element of an array-valued metric group, +// such as a single filesystem (and, in the future, a single network +// interconnect). Alongside its identity (Name/Type) it carries its own set of +// scoped metrics keyed by metric name (e.g. "read_bw", "write_bw"). // -// For example: jobData["cpu_load"][MetricScopeNode] contains node-level CPU load data. -// This structure allows efficient lookup of metrics at different hierarchical levels. -type JobData map[string]map[MetricScope]*JobMetric +// In the job-archive JSON, per-instance metrics are node-scoped; the scope map +// is kept general so finer scopes remain representable without type changes. +type MetricGroupInstance struct { + Name string `json:"name"` + Type string `json:"type"` + Metrics map[string]ScopedMetrics `json:"-"` // flattened onto the instance object by MarshalJSON +} + +// MetricGroup is an array-valued group of named instances, identified by its +// top-level JSON key (e.g. "filesystems"). Instance order is preserved because +// the schema array order is meaningful for rendering. +type MetricGroup struct { + Key string + Instances []MetricGroupInstance +} -// ScopedJobStats maps metric names to statistical summaries organized by scope. -// Structure: map[metricName]map[scope][]*ScopedStats +// JobData holds all metric data of a HPC job. // -// Used to store pre-computed statistics without the full time series data, -// reducing memory footprint when only aggregated values are needed. -type ScopedJobStats map[string]map[MetricScope][]*ScopedStats +// Metrics maps a metric name to its scope-organized data, for example +// jobData.Metrics["cpu_load"][MetricScopeNode]. Groups holds array-valued +// metric groups (filesystems, and later interconnects) whose members each +// carry their own scoped metrics; these cannot be expressed by the flat +// Metrics map because a plain map has nowhere to store per-instance identity. +// +// Custom (Un)MarshalJSON keep the on-disk/on-wire layout identical to the JSON +// schema: flat metrics are top-level scope objects and each group is a +// top-level array (see job-data.schema.json). +type JobData struct { + Metrics map[string]ScopedMetrics + Groups []MetricGroup +} + +// ScopedMetricStats maps a scope to its per-source statistical summaries. +type ScopedMetricStats map[MetricScope][]*ScopedStats + +// ScopedStatsGroupInstance is the ScopedJobStats analogue of +// MetricGroupInstance: one named instance carrying scope-organized statistics. +type ScopedStatsGroupInstance struct { + Name string `json:"name"` + Type string `json:"type"` + Metrics map[string]ScopedMetricStats `json:"-"` +} + +// ScopedStatsGroup is an array-valued group of ScopedStatsGroupInstance. +type ScopedStatsGroup struct { + Key string + Instances []ScopedStatsGroupInstance +} + +// ScopedJobStats stores pre-computed statistics without the full time series +// data, reducing memory footprint when only aggregated values are needed. It +// mirrors JobData: flat metrics in Metrics, array-valued groups in Groups. +type ScopedJobStats struct { + Metrics map[string]ScopedMetricStats + Groups []ScopedStatsGroup +} + +// metricGroupKeys is the set of top-level JSON keys that are array-valued +// metric groups rather than scope objects. Adding a new group (e.g. +// "interconnects") is a single registration; no new types or codec are needed. +var metricGroupKeys = map[string]bool{ + "filesystems": true, +} + +// RegisterMetricGroup registers an additional top-level key as an array-valued +// metric group so that it is (de)serialized into JobData.Groups. +func RegisterMetricGroup(key string) { + metricGroupKeys[key] = true +} + +// IsMetricGroupKey reports whether key denotes an array-valued metric group. +func IsMetricGroupKey(key string) bool { + return metricGroupKeys[key] +} + +// firstToken returns the first non-whitespace byte of a JSON value, or 0 if +// the value is empty. Used to distinguish an array group ('[') from a scope +// object ('{') when a key is not (yet) registered. +func firstToken(b []byte) byte { + for _, c := range b { + switch c { + case ' ', '\t', '\n', '\r': + continue + default: + return c + } + } + return 0 +} + +// MarshalJSON renders JobData in the job-data schema layout: every flat metric +// as a top-level scope object, and every group as a top-level array of +// {name, type, : , ...} instances. Per-metric values reuse +// the existing *JobMetric/Series marshalling (NaN -> null preserved). +func (jd JobData) MarshalJSON() ([]byte, error) { + out := make(map[string]json.RawMessage, len(jd.Metrics)+len(jd.Groups)) + + for name, scopes := range jd.Metrics { + b, err := json.Marshal(scopes) + if err != nil { + return nil, err + } + out[name] = b + } + + for _, group := range jd.Groups { + arr := make([]json.RawMessage, 0, len(group.Instances)) + for _, inst := range group.Instances { + obj := make(map[string]json.RawMessage, len(inst.Metrics)+2) + name, err := json.Marshal(inst.Name) + if err != nil { + return nil, err + } + obj["name"] = name + typ, err := json.Marshal(inst.Type) + if err != nil { + return nil, err + } + obj["type"] = typ + for metric, scopes := range inst.Metrics { + b, err := json.Marshal(scopes) + if err != nil { + return nil, err + } + obj[metric] = b + } + b, err := json.Marshal(obj) + if err != nil { + return nil, err + } + arr = append(arr, b) + } + b, err := json.Marshal(arr) + if err != nil { + return nil, err + } + out[group.Key] = b + } + + return json.Marshal(out) +} + +// UnmarshalJSON parses the job-data schema layout into JobData. Top-level keys +// are dispatched by group-key registration; as a robustness fallback an +// unregistered key whose value is a JSON array is also treated as a group, so a +// reader can ingest a new array group before it is registered. +func (jd *JobData) UnmarshalJSON(b []byte) error { + var raw map[string]json.RawMessage + if err := json.Unmarshal(b, &raw); err != nil { + return err + } + + jd.Metrics = make(map[string]ScopedMetrics, len(raw)) + jd.Groups = nil + + for key, msg := range raw { + if IsMetricGroupKey(key) || firstToken(msg) == '[' { + group, err := unmarshalMetricGroup(key, msg) + if err != nil { + return err + } + jd.Groups = append(jd.Groups, group) + continue + } + + var scopes ScopedMetrics + if err := json.Unmarshal(msg, &scopes); err != nil { + return err + } + jd.Metrics[key] = scopes + } + + return nil +} + +// unmarshalMetricGroup decodes a top-level array value into a MetricGroup, +// splitting each instance's name/type from its per-metric scope objects. +func unmarshalMetricGroup(key string, msg json.RawMessage) (MetricGroup, error) { + group := MetricGroup{Key: key} + + var items []map[string]json.RawMessage + if err := json.Unmarshal(msg, &items); err != nil { + return group, err + } + + for _, item := range items { + inst := MetricGroupInstance{Metrics: make(map[string]ScopedMetrics, len(item))} + for field, fmsg := range item { + switch field { + case "name": + if err := json.Unmarshal(fmsg, &inst.Name); err != nil { + return group, err + } + case "type": + if err := json.Unmarshal(fmsg, &inst.Type); err != nil { + return group, err + } + default: + var scopes ScopedMetrics + if err := json.Unmarshal(fmsg, &scopes); err != nil { + return group, err + } + inst.Metrics[field] = scopes + } + } + group.Instances = append(group.Instances, inst) + } + + return group, nil +} + +// FlatMap returns the flat scoped-metrics map, providing the pre-struct +// access shape (metric -> scope -> *JobMetric) for consumers that do not need +// the grouped metrics. +func (jd JobData) FlatMap() map[string]ScopedMetrics { + return jd.Metrics +} + +// AddGroupInstance appends a named instance (with its own scoped metrics) to the +// group identified by groupKey, creating the group if necessary. This is the +// converter seam used to assemble grouped data from flat, selector-style query +// results without the caller knowing the JobData layout. +func (jd *JobData) AddGroupInstance(groupKey, name, typ string, metrics map[string]ScopedMetrics) { + inst := MetricGroupInstance{Name: name, Type: typ, Metrics: metrics} + for i := range jd.Groups { + if jd.Groups[i].Key == groupKey { + jd.Groups[i].Instances = append(jd.Groups[i].Instances, inst) + return + } + } + jd.Groups = append(jd.Groups, MetricGroup{Key: groupKey, Instances: []MetricGroupInstance{inst}}) +} + +// MarshalJSON renders ScopedJobStats in the same top-level layout as JobData: +// flat metrics as scope objects, groups as top-level arrays of +// {name, type, : } instances. +func (sjs ScopedJobStats) MarshalJSON() ([]byte, error) { + out := make(map[string]json.RawMessage, len(sjs.Metrics)+len(sjs.Groups)) + + for name, scopes := range sjs.Metrics { + b, err := json.Marshal(scopes) + if err != nil { + return nil, err + } + out[name] = b + } + + for _, group := range sjs.Groups { + arr := make([]json.RawMessage, 0, len(group.Instances)) + for _, inst := range group.Instances { + obj := make(map[string]json.RawMessage, len(inst.Metrics)+2) + name, err := json.Marshal(inst.Name) + if err != nil { + return nil, err + } + obj["name"] = name + typ, err := json.Marshal(inst.Type) + if err != nil { + return nil, err + } + obj["type"] = typ + for metric, scopes := range inst.Metrics { + b, err := json.Marshal(scopes) + if err != nil { + return nil, err + } + obj[metric] = b + } + b, err := json.Marshal(obj) + if err != nil { + return nil, err + } + arr = append(arr, b) + } + b, err := json.Marshal(arr) + if err != nil { + return nil, err + } + out[group.Key] = b + } + + return json.Marshal(out) +} + +// UnmarshalJSON parses the top-level layout into ScopedJobStats, dispatching +// array-valued group keys into Groups (see JobData.UnmarshalJSON). +func (sjs *ScopedJobStats) UnmarshalJSON(b []byte) error { + var raw map[string]json.RawMessage + if err := json.Unmarshal(b, &raw); err != nil { + return err + } + + sjs.Metrics = make(map[string]ScopedMetricStats, len(raw)) + sjs.Groups = nil + + for key, msg := range raw { + if IsMetricGroupKey(key) || firstToken(msg) == '[' { + group := ScopedStatsGroup{Key: key} + var items []map[string]json.RawMessage + if err := json.Unmarshal(msg, &items); err != nil { + return err + } + for _, item := range items { + inst := ScopedStatsGroupInstance{Metrics: make(map[string]ScopedMetricStats, len(item))} + for field, fmsg := range item { + switch field { + case "name": + if err := json.Unmarshal(fmsg, &inst.Name); err != nil { + return err + } + case "type": + if err := json.Unmarshal(fmsg, &inst.Type); err != nil { + return err + } + default: + var scopes ScopedMetricStats + if err := json.Unmarshal(fmsg, &scopes); err != nil { + return err + } + inst.Metrics[field] = scopes + } + } + group.Instances = append(group.Instances, inst) + } + sjs.Groups = append(sjs.Groups, group) + continue + } + + var scopes ScopedMetricStats + if err := json.Unmarshal(msg, &scopes); err != nil { + return err + } + sjs.Metrics[key] = scopes + } + + return nil +} // JobMetric contains time series data and statistics for a single metric. // @@ -166,7 +500,7 @@ func (e MetricScope) Valid() bool { func (jd *JobData) Size() int { n := 128 - for _, scopes := range *jd { + sizeScopes := func(scopes ScopedMetrics) { for _, metric := range scopes { if metric.StatisticsSeries != nil { n += len(metric.StatisticsSeries.Max) @@ -180,6 +514,17 @@ func (jd *JobData) Size() int { } } } + + for _, scopes := range jd.Metrics { + sizeScopes(scopes) + } + for _, group := range jd.Groups { + for _, inst := range group.Instances { + for _, scopes := range inst.Metrics { + sizeScopes(scopes) + } + } + } return n * int(unsafe.Sizeof(Float(0))) } @@ -271,7 +616,7 @@ func (jm *JobMetric) AddStatisticsSeries() { } func (jd *JobData) AddNodeScope(metric string) bool { - scopes, ok := (*jd)[metric] + scopes, ok := jd.Metrics[metric] if !ok { return false } @@ -340,7 +685,7 @@ func (jd *JobData) AddNodeScope(metric string) bool { func (jd *JobData) RoundMetricStats() { // TODO: Make Digit-Precision Configurable? (Currently: Fixed to 2 Digits) - for _, scopes := range *jd { + roundScopes := func(scopes ScopedMetrics) { for _, jm := range scopes { for index := range jm.Series { jm.Series[index].Statistics = MetricStatistics{ @@ -351,11 +696,22 @@ func (jd *JobData) RoundMetricStats() { } } } + + for _, scopes := range jd.Metrics { + roundScopes(scopes) + } + for _, group := range jd.Groups { + for _, inst := range group.Instances { + for _, scopes := range inst.Metrics { + roundScopes(scopes) + } + } + } } func (sjs *ScopedJobStats) RoundScopedMetricStats() { // TODO: Make Digit-Precision Configurable? (Currently: Fixed to 2 Digits) - for _, scopes := range *sjs { + roundScopes := func(scopes ScopedMetricStats) { for _, stats := range scopes { for index := range stats { roundedStats := MetricStatistics{ @@ -367,6 +723,17 @@ func (sjs *ScopedJobStats) RoundScopedMetricStats() { } } } + + for _, scopes := range sjs.Metrics { + roundScopes(scopes) + } + for _, group := range sjs.Groups { + for _, inst := range group.Instances { + for _, scopes := range inst.Metrics { + roundScopes(scopes) + } + } + } } func (jm *JobMetric) AddPercentiles(ps []int) bool { diff --git a/schema/metrics_groups_test.go b/schema/metrics_groups_test.go new file mode 100644 index 0000000..39f6b5d --- /dev/null +++ b/schema/metrics_groups_test.go @@ -0,0 +1,179 @@ +// Copyright (C) NHR@FAU, University Erlangen-Nuremberg. +// All rights reserved. This file is part of cc-lib. +// Use of this source code is governed by a MIT-style +// license that can be found in the LICENSE file. + +package schema + +import ( + "bytes" + "encoding/json" + "reflect" + "strings" + "testing" +) + +// one scope object of job-metric-data, reused across the fixtures below. +const nodeMetricJSON = `{"node":{"unit":{"base":"B/s"},"timestep":60,"series":[{"hostname":"h1","statistics":{"avg":2,"min":1,"max":3},"data":[1,2,3]}]}}` + +// jobDataFixture is a schema-valid job-data document. It contains every required +// top-level metric, a top-level metric named "open" that collides with a +// filesystem sub-metric of the same name, and a two-instance filesystems array. +func jobDataFixture() string { + m := nodeMetricJSON + return `{` + + `"cpu_user":` + m + `,` + + `"cpu_load":` + m + `,` + + `"mem_used":` + m + `,` + + `"flops_any":` + m + `,` + + `"mem_bw":` + m + `,` + + `"net_bw":` + m + `,` + + `"open":` + m + `,` + + `"filesystems":[` + + `{"name":"home","type":"nfs","read_bw":` + m + `,"write_bw":` + m + `,"open":` + m + `},` + + `{"name":"scratch","type":"lustre","read_bw":` + m + `,"write_bw":` + m + `}` + + `]}` +} + +func TestJobData_GroupsRoundTrip(t *testing.T) { + fixture := jobDataFixture() + + if err := Validate(Data, strings.NewReader(fixture)); err != nil { + t.Fatalf("fixture does not validate against job-data schema: %v", err) + } + + var jd JobData + if err := json.Unmarshal([]byte(fixture), &jd); err != nil { + t.Fatalf("unmarshal fixture: %v", err) + } + + // Flat metric and the colliding top-level "open" survive. + if _, ok := jd.Metrics["flops_any"]; !ok { + t.Error("flat metric flops_any missing") + } + if _, ok := jd.Metrics["open"]; !ok { + t.Error("top-level metric open missing (collision)") + } + if _, ok := jd.Metrics["filesystems"]; ok { + t.Error("filesystems must not land in flat Metrics") + } + + // Group parsed with instance identity and its own metrics (incl. its own "open"). + if len(jd.Groups) != 1 || jd.Groups[0].Key != "filesystems" { + t.Fatalf("expected one filesystems group, got %+v", jd.Groups) + } + insts := jd.Groups[0].Instances + if len(insts) != 2 { + t.Fatalf("expected 2 filesystem instances, got %d", len(insts)) + } + if insts[0].Name != "home" || insts[0].Type != "nfs" { + t.Errorf("instance[0] identity wrong: %+v", insts[0]) + } + if insts[1].Name != "scratch" || insts[1].Type != "lustre" { + t.Errorf("instance[1] identity wrong: %+v", insts[1]) + } + if _, ok := insts[0].Metrics["open"]; !ok { + t.Error("filesystem sub-metric open missing (collision with top-level open)") + } + if _, ok := insts[0].Metrics["read_bw"]; !ok { + t.Error("filesystem sub-metric read_bw missing") + } + + // Marshal back -> still schema-valid, and filesystems is an array. + out, err := json.Marshal(jd) + if err != nil { + t.Fatalf("marshal: %v", err) + } + if err := Validate(Data, bytes.NewReader(out)); err != nil { + t.Fatalf("re-marshalled job data does not validate: %v", err) + } + if !bytes.Contains(out, []byte(`"filesystems":[`)) { + t.Errorf("marshalled output does not contain filesystems array: %s", out) + } + + // Re-unmarshal and compare structurally (instance order preserved). + var jd2 JobData + if err := json.Unmarshal(out, &jd2); err != nil { + t.Fatalf("re-unmarshal: %v", err) + } + if !reflect.DeepEqual(jd, jd2) { + t.Errorf("round-trip mismatch:\n first: %+v\nsecond: %+v", jd, jd2) + } +} + +func TestJobData_NaNMarshalsAsNull(t *testing.T) { + id := "0" + jd := JobData{ + Metrics: map[string]ScopedMetrics{ + "flops_any": { + MetricScopeNode: &JobMetric{ + Unit: Unit{Base: "F/s"}, + Timestep: 60, + Series: []Series{ + {Hostname: "h1", ID: &id, Data: []Float{1, NaN, 3}, Statistics: MetricStatistics{Min: 1, Avg: 2, Max: 3}}, + }, + }, + }, + }, + } + + out, err := json.Marshal(jd) + if err != nil { + t.Fatalf("marshal: %v", err) + } + if !bytes.Contains(out, []byte("null")) { + t.Errorf("NaN data point not rendered as null: %s", out) + } +} + +func TestJobData_SizeIncludesGroups(t *testing.T) { + mk := func() ScopedMetrics { + return ScopedMetrics{ + MetricScopeNode: &JobMetric{ + Unit: Unit{Base: "B/s"}, + Timestep: 60, + Series: []Series{{Hostname: "h1", Data: []Float{1, 2, 3, 4, 5}}}, + }, + } + } + + base := JobData{Metrics: map[string]ScopedMetrics{"read_bw": mk()}} + withGroup := JobData{Metrics: map[string]ScopedMetrics{"read_bw": mk()}} + withGroup.AddGroupInstance("filesystems", "home", "nfs", map[string]ScopedMetrics{"read_bw": mk()}) + + if withGroup.Size() <= base.Size() { + t.Errorf("Size() must account for group metrics: base=%d withGroup=%d", base.Size(), withGroup.Size()) + } +} + +func TestJobStatisticsSet_RoundTrip(t *testing.T) { + set := JobStatisticsSet{ + Metrics: map[string]JobStatistics{ + "cpu_load": {Unit: Unit{Base: ""}, Avg: 2, Min: 1, Max: 3}, + "mem_bw": {Unit: Unit{Base: "B/s", Prefix: "G"}, Avg: 20, Min: 10, Max: 30}, + }, + } + set.AddGroupInstance("filesystems", "home", "nfs", map[string]JobStatistics{ + "read_bw": {Unit: Unit{Base: "B/s"}, Avg: 5, Min: 1, Max: 9}, + "write_bw": {Unit: Unit{Base: "B/s"}, Avg: 6, Min: 2, Max: 8}, + }) + + out, err := json.Marshal(set) + if err != nil { + t.Fatalf("marshal: %v", err) + } + if !bytes.Contains(out, []byte(`"filesystems":[`)) { + t.Errorf("statistics output missing filesystems array: %s", out) + } + if !bytes.Contains(out, []byte(`"cpu_load"`)) { + t.Errorf("statistics output missing flat metric: %s", out) + } + + var got JobStatisticsSet + if err := json.Unmarshal(out, &got); err != nil { + t.Fatalf("unmarshal: %v", err) + } + if !reflect.DeepEqual(set, got) { + t.Errorf("round-trip mismatch:\n first: %+v\nsecond: %+v", set, got) + } +} diff --git a/schema/metrics_test.go b/schema/metrics_test.go index 27d3fae..9997864 100644 --- a/schema/metrics_test.go +++ b/schema/metrics_test.go @@ -14,15 +14,17 @@ func TestAddNodeScope(t *testing.T) { // host1 has 2 cores with 3 data points each. // host2 has 2 cores with 5 data points each. jd := JobData{ - "flops_any": { - MetricScopeCore: &JobMetric{ - Unit: Unit{Base: "F/s"}, - Timestep: 10, - Series: []Series{ - {Hostname: "host1", Data: []Float{1, 2, 3}, Statistics: MetricStatistics{Min: 1, Avg: 2, Max: 3}}, - {Hostname: "host1", Data: []Float{4, 5, 6}, Statistics: MetricStatistics{Min: 4, Avg: 5, Max: 6}}, - {Hostname: "host2", Data: []Float{10, 20, 30, 40, 50}, Statistics: MetricStatistics{Min: 10, Avg: 30, Max: 50}}, - {Hostname: "host2", Data: []Float{100, 200, 300, 400, 500}, Statistics: MetricStatistics{Min: 100, Avg: 300, Max: 500}}, + Metrics: map[string]ScopedMetrics{ + "flops_any": { + MetricScopeCore: &JobMetric{ + Unit: Unit{Base: "F/s"}, + Timestep: 10, + Series: []Series{ + {Hostname: "host1", Data: []Float{1, 2, 3}, Statistics: MetricStatistics{Min: 1, Avg: 2, Max: 3}}, + {Hostname: "host1", Data: []Float{4, 5, 6}, Statistics: MetricStatistics{Min: 4, Avg: 5, Max: 6}}, + {Hostname: "host2", Data: []Float{10, 20, 30, 40, 50}, Statistics: MetricStatistics{Min: 10, Avg: 30, Max: 50}}, + {Hostname: "host2", Data: []Float{100, 200, 300, 400, 500}, Statistics: MetricStatistics{Min: 100, Avg: 300, Max: 500}}, + }, }, }, }, @@ -33,7 +35,7 @@ func TestAddNodeScope(t *testing.T) { t.Fatal("AddNodeScope returned false") } - nodeMetric, exists := jd["flops_any"][MetricScopeNode] + nodeMetric, exists := jd.Metrics["flops_any"][MetricScopeNode] if !exists { t.Fatal("node scope not created") } @@ -76,13 +78,15 @@ func TestAddNodeScope(t *testing.T) { func TestAddNodeScopeUnevenCores(t *testing.T) { // Same host, cores with different data lengths. jd := JobData{ - "mem_bw": { - MetricScopeCore: &JobMetric{ - Unit: Unit{Base: "B/s"}, - Timestep: 10, - Series: []Series{ - {Hostname: "node1", Data: []Float{1, 2, 3}, Statistics: MetricStatistics{Min: 1, Avg: 2, Max: 3}}, - {Hostname: "node1", Data: []Float{10, 20, 30, 40, 50}, Statistics: MetricStatistics{Min: 10, Avg: 30, Max: 50}}, + Metrics: map[string]ScopedMetrics{ + "mem_bw": { + MetricScopeCore: &JobMetric{ + Unit: Unit{Base: "B/s"}, + Timestep: 10, + Series: []Series{ + {Hostname: "node1", Data: []Float{1, 2, 3}, Statistics: MetricStatistics{Min: 1, Avg: 2, Max: 3}}, + {Hostname: "node1", Data: []Float{10, 20, 30, 40, 50}, Statistics: MetricStatistics{Min: 10, Avg: 30, Max: 50}}, + }, }, }, }, @@ -93,7 +97,7 @@ func TestAddNodeScopeUnevenCores(t *testing.T) { t.Fatal("AddNodeScope returned false") } - nodeMetric := jd["mem_bw"][MetricScopeNode] + nodeMetric := jd.Metrics["mem_bw"][MetricScopeNode] if len(nodeMetric.Series) != 1 { t.Fatalf("expected 1 node series, got %d", len(nodeMetric.Series)) } diff --git a/sinks/httpSink.go b/sinks/httpSink.go index dc6ca6a..e4b21f3 100644 --- a/sinks/httpSink.go +++ b/sinks/httpSink.go @@ -117,7 +117,8 @@ func (s *HttpSink) Write(msg lp.CCMessage) error { if err := s.Flush(); err != nil { cclog.ComponentError(s.name, "Flush triggered by flush timer: flush failed:", err) } - }) + }, + ) } } diff --git a/sinks/influxSink.go b/sinks/influxSink.go index faf3f97..04d9c8f 100644 --- a/sinks/influxSink.go +++ b/sinks/influxSink.go @@ -130,7 +130,8 @@ func (s *InfluxSink) connect() error { s.name, "connect():", "Influx client options HTTPRequestTimeout:", - time.Second*time.Duration(clientOptions.HTTPRequestTimeout())) + time.Second*time.Duration(clientOptions.HTTPRequestTimeout()), + ) // Set retry interval if len(s.config.InfluxRetryInterval) > 0 { @@ -145,7 +146,8 @@ func (s *InfluxSink) connect() error { s.name, "connect():", "Influx client options RetryInterval:", - time.Millisecond*time.Duration(clientOptions.RetryInterval())) + time.Millisecond*time.Duration(clientOptions.RetryInterval()), + ) // Set the maximum delay between each retry attempt if len(s.config.InfluxMaxRetryInterval) > 0 { @@ -160,7 +162,8 @@ func (s *InfluxSink) connect() error { s.name, "connect():", "Influx client options MaxRetryInterval:", - time.Millisecond*time.Duration(clientOptions.MaxRetryInterval())) + time.Millisecond*time.Duration(clientOptions.MaxRetryInterval()), + ) // Set the base for the exponential retry delay if s.config.InfluxExponentialBase != 0 { @@ -170,7 +173,8 @@ func (s *InfluxSink) connect() error { s.name, "connect():", "Influx client options ExponentialBase:", - clientOptions.ExponentialBase()) + clientOptions.ExponentialBase(), + ) // Set maximum count of retry attempts of failed writes if s.config.InfluxMaxRetries != 0 { @@ -180,7 +184,8 @@ func (s *InfluxSink) connect() error { s.name, "connect():", "Influx client options MaxRetries:", - clientOptions.MaxRetries()) + clientOptions.MaxRetries(), + ) // Set the maximum total retry timeout if len(s.config.InfluxMaxRetryTime) > 0 { @@ -196,7 +201,8 @@ func (s *InfluxSink) connect() error { s.name, "connect():", "Influx client options MaxRetryTime:", - time.Millisecond*time.Duration(clientOptions.MaxRetryTime())) + time.Millisecond*time.Duration(clientOptions.MaxRetryTime()), + ) // Specify whether to use GZip compression in write requests clientOptions.SetUseGZip(s.config.InfluxUseGzip) @@ -204,7 +210,8 @@ func (s *InfluxSink) connect() error { s.name, "connect():", "Influx client options UseGZip:", - clientOptions.UseGZip()) + clientOptions.UseGZip(), + ) // Do not check InfluxDB certificate clientOptions.SetTLSConfig( @@ -356,7 +363,8 @@ func (s *InfluxSink) Write(msg lp.CCMessage) error { if err := s.Flush(); err != nil { cclog.ComponentError(s.name, "Flush triggered by flush timer: flush failed:", err) } - }) + }, + ) } } diff --git a/sinks/natsSink.go b/sinks/natsSink.go index e4a6a79..b1dfdc4 100644 --- a/sinks/natsSink.go +++ b/sinks/natsSink.go @@ -117,7 +117,8 @@ func (s *NatsSink) Write(m lp.CCMessage) error { if err := s.Flush(); err != nil { cclog.ComponentError(s.name, "Flush triggered by flush timer: flush failed:", err) } - }) + }, + ) } }